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/crates/rds-desktop/src/delivery_rate.rs b/crates/rds-desktop/src/delivery_rate.rs new file mode 100644 index 0000000..e21fd44 --- /dev/null +++ b/crates/rds-desktop/src/delivery_rate.rs @@ -0,0 +1,97 @@ +//! Conservative recent goodput from successful, timely media receipts only. +//! Low traffic is not a link-capacity measurement; silence never lowers a rate. +#[derive(Default)] +pub(crate) struct DeliveryRate { + path: Option, + start_ms: u64, + bytes: u64, + receipts: u64, + previous_rate: Option, + floor: Option<(u64, u64)>, +} + +impl DeliveryRate { + pub(crate) fn sample( + &mut self, + path: Option, + now_ms: u64, + bytes: u64, + receipts: u64, + impaired: bool, + ) -> Option { + if path.is_none() || path != self.path || impaired || now_ms < self.start_ms { + *self = Self { + path, + start_ms: now_ms, + bytes, + receipts, + ..Self::default() + }; + return None; + } + let elapsed = now_ms.saturating_sub(self.start_ms); + if elapsed >= 1000 { + let delivered = bytes.saturating_sub(self.bytes); + let count = receipts.saturating_sub(self.receipts); + let rate = (elapsed <= 2500 && count >= 3 && delivered >= 4096) + .then(|| delivered.saturating_mul(8000) / elapsed); + if let Some((previous, current)) = self.previous_rate.zip(rate) { + // Two adjacent sufficiently populated windows and 20% headroom. + // This is a floor justified by delivery, never a capacity ceiling. + let minimum = previous.min(current); + self.floor = Some((minimum / 5 * 4, now_ms)); + } else { + self.floor = None; + } + self.previous_rate = rate; + self.start_ms = now_ms; + self.bytes = bytes; + self.receipts = receipts; + } + self.floor + .filter(|(_, sampled)| now_ms.saturating_sub(*sampled) <= 2000) + .map(|(rate, _)| rate) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn requires_two_windows_and_uses_the_lower_rate_with_headroom() { + let mut rate = DeliveryRate::default(); + assert_eq!(rate.sample(Some(1), 0, 500, 4, false), None); + assert_eq!(rate.sample(Some(1), 1000, 250_500, 14, false), None); + assert_eq!( + rate.sample(Some(1), 2000, 450_500, 24, false), + Some(1_280_000) + ); + assert_eq!( + rate.sample(Some(1), 2250, 450_500, 24, false), + Some(1_280_000) + ); + assert_eq!( + rate.sample(Some(1), 3000, 450_500, 24, false), + None, + "idle traffic is not capacity proof" + ); + } + + #[test] + fn failure_path_change_and_silence_revoke_the_observation() { + for reason in 0..4 { + let mut rate = DeliveryRate::default(); + rate.sample(Some(1), 0, 0, 0, false); + rate.sample(Some(1), 1000, 250_000, 10, false); + assert!(rate.sample(Some(1), 2000, 500_000, 20, false).is_some()); + let (path, time, impaired) = match reason { + 0 => (Some(1), 2250, true), + 1 => (Some(2), 2250, false), + 2 => (None, 2250, false), + _ => (Some(1), 5000, false), + }; + assert_eq!(rate.sample(path, time, 500_000, 20, impaired), None); + } + } +} diff --git a/crates/rds-desktop/src/lib.rs b/crates/rds-desktop/src/lib.rs index b697c76..b83ec86 100644 --- a/crates/rds-desktop/src/lib.rs +++ b/crates/rds-desktop/src/lib.rs @@ -25,6 +25,7 @@ pub mod client; pub mod clipboard; pub mod codec; mod decode_work; +mod delivery_rate; pub mod input; pub mod mailbox; mod order; diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index e98af6d..24eb6a7 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -60,6 +60,8 @@ struct DeliveryFeedback { late_pending: AtomicU64, failed: AtomicU64, acknowledged: AtomicU64, + timely_bytes: AtomicU64, + timely_receipts: AtomicU64, obsolete: AtomicU64, last_ack_ms: AtomicU64, producing: AtomicBool, @@ -70,6 +72,16 @@ struct DeliveryFeedback { max_input_inject_ms: AtomicU64, } +impl DeliveryFeedback { + fn acknowledged(&self, bytes: usize, elapsed: Duration, budget: Duration) { + self.acknowledged.fetch_add(1, Ordering::Relaxed); + if elapsed <= budget { + self.timely_bytes.fetch_add(bytes as u64, Ordering::Relaxed); + self.timely_receipts.fetch_add(1, Ordering::Relaxed); + } + } +} + struct LateReceipt(Arc); impl LateReceipt { @@ -320,6 +332,7 @@ pub struct BitrateController { last_path: Option, delivery_hold_ticks: u8, delivery_cut_cooldown_ticks: u8, + rtt_reduction: bool, } impl BitrateController { @@ -341,6 +354,7 @@ impl BitrateController { last_path: None, delivery_hold_ticks: 0, delivery_cut_cooldown_ticks: 0, + rtt_reduction: false, } } @@ -361,6 +375,7 @@ impl BitrateController { pub fn step(&mut self, path: Option, deadline_misses: u64) -> u64 { let mut next = self.current; self.reduction_reason = None; + self.rtt_reduction = false; if let Some(p) = path { if self.last_path != Some(p.path_id) { self.last_path = Some(p.path_id); @@ -422,6 +437,7 @@ impl BitrateController { // burst for a second, as with correlated media-receipt pressure. if self.primed && (loss_high || rtt_high) { if self.path_cut_cooldown_ticks == 0 { + self.rtt_reduction = rtt_high; next = (next / 10 * 7 + next % 10 * 7 / 10).max(self.floor); self.path_cut_cooldown_ticks = PATH_CUT_COOLDOWN_TICKS; self.reduction_reason = Some(if loss_high { @@ -451,14 +467,27 @@ impl BitrateController { // one second, and increase only on fresh successful delivery // and at 1% per sample so clean relay packet counters cannot immediately // drive the encoder back into the same backlog. + #[cfg(test)] fn step_with_delivery( &mut self, path: Option, deadline_misses: u64, impaired: bool, delivered: bool, + ) -> u64 { + self.step_with_delivery_floor(path, deadline_misses, impaired, delivered, None) + } + + fn step_with_delivery_floor( + &mut self, + path: Option, + deadline_misses: u64, + impaired: bool, + delivered: bool, + delivery_floor: Option, ) -> u64 { let previous = self.current; + let held = self.delivery_hold_ticks > 0; let proposed = self.step(path, deadline_misses); self.delivery_cut_cooldown_ticks = self.delivery_cut_cooldown_ticks.saturating_sub(1); self.current = if impaired { @@ -495,6 +524,21 @@ impl BitrateController { }; proposed.min(previous.saturating_add((previous / 100).max(1))) }; + // Independent QUIC loss declarations must not lower a healthy media + // stream beneath conservatively observed successful timely goodput. + // RTT/producer pressure and real delivery problems retain their cuts. + if !impaired + && !held + && delivered + && self.delivery_hold_ticks == 0 + && !self.rtt_reduction + && self.rtt_rise_baseline.is_none() + && deadline_misses == 0 + && self.reduction_reason == Some("sampled_packet_loss") + && let Some(floor) = delivery_floor + { + self.current = self.current.max(floor.min(self.ceiling)); + } self.current } } @@ -681,6 +725,7 @@ pub async fn serve_desktop_with( let pending_key = keyframe_pending.clone(); let progress_clock = clock.clone(); let mut controller = BitrateController::new(4_000_000, ceiling); + let mut delivery_rate = crate::delivery_rate::DeliveryRate::default(); workers.spawn(async move { let mut last_delayed = 0; let mut last_failed = 0; @@ -708,17 +753,25 @@ pub async fn serve_desktop_with( let delivered = acknowledged > last_acknowledged; let late_pending = feedback.late_pending.load(Ordering::Relaxed); let impaired = pressure.sample(late_pending, delivered, failed_frames > 0); - let bps = controller.step_with_delivery( + let delivery_floor = delivery_rate.sample( + path.map(|p| p.path_id), progress_clock.now_ms(), + feedback.timely_bytes.load(Ordering::Relaxed), + feedback.timely_receipts.load(Ordering::Relaxed), + impaired || failed_frames > 0 || late_pending > 0, + ); + let bps = controller.step_with_delivery_floor( path, missed, impaired, delivered, + delivery_floor, ); (last_delayed, last_failed, last_acknowledged) = (delayed, failed, acknowledged); if bps != previous { tracing::debug!( previous_bps = previous, bitrate_bps = bps, + timely_delivery_floor_bps = delivery_floor, deadline_misses = missed, delayed_frames, failed_frames, @@ -766,6 +819,9 @@ pub async fn serve_desktop_with( keyframe_pending = pending_key.load(Ordering::Acquire), bitrate_bps = bps, delayed_delivery = delayed, + timely_delivery_floor_bps = delivery_floor, + timely_delivery_bytes = feedback.timely_bytes.load(Ordering::Relaxed), + timely_deliveries = feedback.timely_receipts.load(Ordering::Relaxed), late_pending, delivery_stalled_ticks = pressure.stalled_ticks, failed_delivery = failed, @@ -812,21 +868,14 @@ pub async fn serve_desktop_with( let mut failed_delivery = 0u64; let mut health = Instant::now(); 'writer: loop { - while let Some(result) = acknowledgements.try_join_next() { - match result { - Ok((_, FrameReceipt::Delivered)) => acknowledged += 1, - Ok((_, FrameReceipt::Obsolete)) => obsolete += 1, - Ok((seq, FrameReceipt::Failed)) => { - failed_delivery += 1; - chain.failed_receipt(seq); - } - _ => { - failed_delivery += 1; - chain.next = None; - writer_idr.store(true, Ordering::Relaxed); - } - } - } + drain_frame_receipts( + &mut acknowledgements, + &mut chain, + &writer_idr, + &mut acknowledged, + &mut obsolete, + &mut failed_delivery, + ); if acknowledgements.len() >= MAX_PENDING_FRAME_ACKS { match acknowledgements.join_next().await { Some(Ok((_, FrameReceipt::Delivered))) => acknowledged += 1, @@ -855,17 +904,43 @@ pub async fn serve_desktop_with( request_frame_repair(&latest_key_seq, &writer_idr, produced.header.seq); continue; } + // Receipts can finish while the writer waits for the next frame. + // Retire them before testing whether a recovered key has an empty + // media queue; completed tasks are not outstanding delivery. + drain_frame_receipts( + &mut acknowledgements, + &mut chain, + &writer_idr, + &mut acknowledged, + &mut obsolete, + &mut failed_delivery, + ); let bps = writer_bitrate.load(Ordering::Relaxed).max(50_000) as f64 / 8.0; let now = Instant::now(); budget = (budget + now.duration_since(last).as_secs_f64() * bps).min(bps * 0.25); last = now; - let cost = produced.payload.len() as f64 + 64.0; - if cost > budget { - let wait = ((cost - budget) / bps).min(0.5); - tokio::time::sleep(Duration::from_secs_f64(wait)).await; - budget = (budget - cost).max(-bps * 0.5); - } else { - budget -= cost; + // With no unconfirmed media, a bounded independent key can start + // immediately. QUIC still paces its packets and the existing key + // receipt barrier prevents dependent capture from accumulating. + let wait = frame_pacing_wait( + &mut budget, + bps, + produced.header.keyframe, + produced.payload.len(), + acknowledgements.len(), + ); + if wait >= Duration::from_millis(100) { + tracing::info!( + frame_seq = produced.header.seq, + keyframe = produced.header.keyframe, + payload_bytes = produced.payload.len(), + wait_ms = wait.as_millis(), + bitrate_bps = writer_bitrate.load(Ordering::Relaxed), + "desktop frame pacing delayed" + ); + } + if !wait.is_zero() { + tokio::time::sleep(wait).await; } // Admission follows the final selection: advancing this before // the pacing wait would lose track of frames collapsed afterward. @@ -1197,13 +1272,35 @@ fn advance_cadence( false } +fn resume_cadence(next_due: &mut Instant, interval: Duration, now: Instant) { + *next_due = (*next_due).max(now.checked_sub(interval).unwrap_or(now)); +} + +fn frame_pacing_wait( + budget: &mut f64, + bytes_per_second: f64, + keyframe: bool, + payload_bytes: usize, + pending_receipts: usize, +) -> Duration { + let cost = payload_bytes as f64 + 64.0; + let cold_key = keyframe && payload_bytes <= 64 * 1024 && pending_receipts == 0; + if cost <= *budget { + *budget -= cost; + return Duration::ZERO; + } + let wait = if cold_key { + Duration::ZERO + } else { + Duration::from_secs_f64(((cost - *budget) / bytes_per_second).min(0.5)) + }; + *budget = (*budget - cost).max(-bytes_per_second * 0.5); + wait +} + impl FrameProducer for SyntheticProducer { fn resume_after_backpressure(&mut self) { - self.next_due = self.next_due.max( - Instant::now() - .checked_sub(self.interval) - .unwrap_or_else(Instant::now), - ); + resume_cadence(&mut self.next_due, self.interval, Instant::now()); } fn produce( &mut self, @@ -1341,11 +1438,7 @@ mod x11 { .interval .max(self.last_work) .min(Duration::from_millis(500)); - self.next_due = self.next_due.max( - Instant::now() - .checked_sub(interval) - .unwrap_or_else(Instant::now), - ); + resume_cadence(&mut self.next_due, interval, Instant::now()); } fn preserves_reference(&self) -> bool { self.skipped @@ -1548,6 +1641,31 @@ async fn send_frame( } } +fn drain_frame_receipts( + receipts: &mut JoinSet<(u64, FrameReceipt)>, + chain: &mut FrameChain, + idr: &AtomicBool, + acknowledged: &mut u64, + obsolete: &mut u64, + failed: &mut u64, +) { + while let Some(result) = receipts.try_join_next() { + match result { + Ok((_, FrameReceipt::Delivered)) => *acknowledged += 1, + Ok((_, FrameReceipt::Obsolete)) => *obsolete += 1, + Ok((seq, FrameReceipt::Failed)) => { + *failed += 1; + chain.failed_receipt(seq); + } + _ => { + *failed += 1; + chain.next = None; + idr.store(true, Ordering::Relaxed); + } + } + } +} + async fn send_frame_inner( conn: &Connection, route: rds_core::UniHello, @@ -1556,6 +1674,7 @@ async fn send_frame_inner( acknowledgements: &mut JoinSet<(u64, FrameReceipt)>, delivery: FrameDelivery, ) -> SendOutcome { + let transfer_started = Instant::now(); let AdmittedFrame { produced, permit } = produced; let mut sending = match conn.open_uni().await { Ok(stream) => FrameSend { @@ -1650,7 +1769,7 @@ async fn send_frame_inner( let acknowledged = match result { Ok(Ok(None)) => { sending.finished = true; - feedback.acknowledged.fetch_add(1, Ordering::Relaxed); + feedback.acknowledged(payload_bytes, transfer_started.elapsed(), delay_budget); tracing::trace!(frame_seq=seq,payload_bytes,ack_ms=started.elapsed().as_millis(),"desktop frame transport acknowledged"); FrameReceipt::Delivered } @@ -1859,20 +1978,22 @@ mod tests { #[test] fn admission_pause_does_not_report_encoder_starvation() { - let controls = ProducerControls::new(4_000_000); - let mut source = SyntheticProducer::new(60, 64, 64, 256); - source.next_due = Instant::now() - Duration::from_secs(2); - source.resume_after_backpressure(); - assert!( - source - .produce(0, &controls, &SessionClock::default()) - .is_some() - ); - assert_eq!(controls.deadline_misses.load(Ordering::Relaxed), 0); + let now = Instant::now(); + let interval = Duration::from_millis(16); + let mut due = now - Duration::from_secs(2); + // Exercise the production resume/advance functions on one explicit + // clock. OS preemption between two calls is actual scheduling delay, + // not proof that admission itself counted as encoder starvation. + resume_cadence(&mut due, interval, now); + assert!(!advance_cadence(&mut due, interval, now, false)); + assert_eq!(due, now); // Real lateness without a deliberate admission pause still reports. - source.next_due = Instant::now() - Duration::from_secs(2); - source.produce(1, &controls, &SessionClock::default()); - assert_eq!(controls.deadline_misses.load(Ordering::Relaxed), 1); + assert!(advance_cadence( + &mut due, + interval, + now + interval * 3, + false + )); } fn produced(seq: u64, keyframe: bool) -> Produced { @@ -2227,6 +2348,8 @@ mod tests { "obsolete disposal cannot justify bitrate growth" ); assert_eq!(feedback.failed.load(Ordering::Acquire), 0); + assert_eq!(feedback.timely_bytes.load(Ordering::Acquire), 0); + assert_eq!(feedback.timely_receipts.load(Ordering::Acquire), 0); assert_eq!( feedback.obsolete.load(Ordering::Acquire), u64::from(obsolete) @@ -2600,6 +2723,160 @@ mod tests { assert_eq!(c.current(), 8_000_000); } + #[test] + fn loss_only_cuts_respect_recent_timely_confirmed_goodput() { + let mut c = BitrateController::new(4_000_000, 8_000_000); + c.step_with_delivery_floor(Some(path(1000, 0, 120, 0)), 0, false, true, None); + for i in 1..=40 { + let bps = c.step_with_delivery_floor( + Some(path(1000 + i * 100, i * 15, 120, i)), + 0, + false, + true, + Some(1_600_000), + ); + assert!( + bps >= 1_600_000, + "confirmed timely media must not collapse: {bps}" + ); + assert!(bps <= 8_000_000); + } + // No receipt evidence keeps the conservative original loss response. + let mut old = BitrateController::new(4_000_000, 8_000_000); + old.step_with_delivery(Some(path(1000, 0, 120, 0)), 0, false, true); + for i in 1..=40 { + old.step_with_delivery(Some(path(1000 + i * 100, i * 15, 120, i)), 0, false, true); + } + assert!( + old.current() < 1_600_000, + "fixture must distinguish the old response" + ); + } + + #[test] + fn delivery_rate_credits_only_complete_timely_receipts() { + let feedback = DeliveryFeedback::default(); + feedback.acknowledged( + 10_000, + Duration::from_millis(100), + Duration::from_millis(250), + ); + feedback.acknowledged(90_000, Duration::from_secs(2), Duration::from_millis(250)); + assert_eq!(feedback.acknowledged.load(Ordering::Relaxed), 2); + assert_eq!(feedback.timely_bytes.load(Ordering::Relaxed), 10_000); + assert_eq!(feedback.timely_receipts.load(Ordering::Relaxed), 1); + } + + #[test] + fn empty_media_key_starts_without_artificial_wait_but_retains_bounded_debt() { + let mut budget = 0.0; + let rate = 100_000.0 / 8.0; + let old_wait = Duration::from_secs_f64((50_000.0f64 / rate).min(0.5)); + assert_eq!( + old_wait, + Duration::from_millis(500), + "fixture must distinguish prior delay" + ); + assert_eq!( + frame_pacing_wait(&mut budget, rate, true, 50_000, 0), + Duration::ZERO + ); + assert_eq!(budget, -rate * 0.5); + assert_eq!( + frame_pacing_wait(&mut budget, rate, false, 1000, 0), + Duration::from_millis(500) + ); + assert_eq!(budget, -rate * 0.5, "a later frame cannot grow debt"); + let mut full = 50_000.0; + assert_eq!( + frame_pacing_wait(&mut full, rate, false, 1000, 0), + Duration::ZERO + ); + assert_eq!(full, 48_936.0); + for (keyframe, bytes, pending) in [(false, 50_000, 0), (true, 50_000, 1), (true, 65_537, 0)] + { + let mut budget = 0.0; + assert_eq!( + frame_pacing_wait(&mut budget, rate, keyframe, bytes, pending), + Duration::from_millis(500), + "nonempty/large/dependent frames retain pacing" + ); + } + } + + #[tokio::test] + async fn receipt_completed_during_idle_wait_does_not_delay_the_next_independent_key() { + let mut receipts = JoinSet::new(); + let done = receipts.spawn(async { (1, FrameReceipt::Delivered) }); + tokio::time::timeout(Duration::from_secs(2), async { + while !done.is_finished() { + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + assert_eq!( + receipts.len(), + 1, + "completed task still needs to be retired" + ); + let (mut acknowledged, mut obsolete, mut failed) = (0, 0, 0); + let mut chain = FrameChain::default(); + let idr = AtomicBool::new(false); + drain_frame_receipts( + &mut receipts, + &mut chain, + &idr, + &mut acknowledged, + &mut obsolete, + &mut failed, + ); + assert_eq!((acknowledged, obsolete, failed), (1, 0, 0)); + assert!(!idr.load(Ordering::Relaxed)); + let mut budget = 0.0; + assert_eq!( + frame_pacing_wait(&mut budget, 12_500.0, true, 50_000, receipts.len()), + Duration::ZERO + ); + } + + #[test] + fn confirmed_goodput_cannot_override_delivery_rtt_or_producer_pressure() { + for reason in 0..4 { + let mut c = BitrateController::new(4_000_000, 8_000_000); + c.step_with_delivery_floor(Some(path(1000, 0, 120, 0)), 0, false, true, None); + let (impaired, delivered, misses, rtt) = match reason { + 0 => (true, false, 0, 120), + 1 => (false, false, 0, 120), + 2 => (false, true, 1, 120), + _ => (false, true, 0, 240), + }; + let bps = c.step_with_delivery_floor( + Some(path(1100, 15, rtt, 1)), + misses, + impaired, + delivered, + Some(8_000_000), + ); + assert!( + bps < 4_000_000, + "real pressure was overridden: reason={reason} bps={bps}" + ); + } + let mut capped = BitrateController::new(4_000_000, 4_000_000); + capped.step_with_delivery_floor(Some(path(1000, 0, 120, 0)), 0, false, true, None); + assert_eq!( + capped.step_with_delivery_floor( + Some(path(1100, 15, 120, 1)), + 0, + false, + true, + Some(u64::MAX) + ), + 4_000_000 + ); + } + #[test] fn steer_survives_adaptation_ticks() { // The SetBitrate semantics the viewer sees: a filed target is not diff --git a/crates/rds-sync/tests/session_v2.rs b/crates/rds-sync/tests/session_v2.rs index 9111b2b..ee17db5 100644 --- a/crates/rds-sync/tests/session_v2.rs +++ b/crates/rds-sync/tests/session_v2.rs @@ -100,8 +100,13 @@ async fn pair() -> ( .with_max_level(tracing::level_filters::LevelFilter::DEBUG) .with_test_writer() .try_init(); - let server_ep = bind_endpoint(EndpointConfig::default()).await.unwrap(); - let client_ep = bind_endpoint(EndpointConfig::default()).await.unwrap(); + let config = EndpointConfig { + discovery: false, + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + ..Default::default() + }; + let server_ep = bind_endpoint(config.clone()).await.unwrap(); + let client_ep = bind_endpoint(config).await.unwrap(); let dir = scratch("server"); let (outcomes, rx) = mpsc::channel(); let target = server_ep.addr(); @@ -251,10 +256,16 @@ async fn v1_route_still_serves_unchanged_for_compat() { .unwrap(); let ack: HelloAck = read_frame(&mut recv).await.unwrap(); assert!(matches!(ack, HelloAck::Ok)); - let stats = Transfer::new(id) + let stats = match Transfer::new(id) .send_file(&conn, &src, (send, recv), Duration::from_secs(60)) .await - .unwrap(); + { + Ok(stats) => stats, + Err(error) => { + let server = rx.recv_timeout(Duration::from_secs(10)); + panic!("v1 client failed: {error:#}; server outcome: {server:?}"); + } + }; assert_eq!(stats.bytes, data.len() as u64); assert_eq!(std::fs::read(server_dir.join("v1.bin")).unwrap(), data); assert_eq!(rx.recv_timeout(Duration::from_secs(10)).unwrap(), "ok"); diff --git a/crates/rds-sync/tests/sync_e2e.rs b/crates/rds-sync/tests/sync_e2e.rs index df05293..02c9dd0 100644 --- a/crates/rds-sync/tests/sync_e2e.rs +++ b/crates/rds-sync/tests/sync_e2e.rs @@ -98,7 +98,16 @@ fn count_parts(state_dir: &Path) -> usize { .map(|e| e.path().join("parts")) .map(|p| { std::fs::read_dir(&p) - .map(|d| d.flatten().filter(|f| f.path().is_file()).count()) + .map(|d| { + d.flatten() + .filter(|f| { + f.path().is_file() + && f.file_name().to_str().is_some_and(|name| { + name.len() == 64 && name.bytes().all(|b| b.is_ascii_hexdigit()) + }) + }) + .count() + }) .unwrap_or(0) }) .sum() @@ -123,6 +132,35 @@ async fn push( send_file(&conn, path, send, recv).await } +/// Client-task cancellation does not join a receiver's running filesystem +/// operation. Wait for its actual exclusive journal admission before the one +/// resumed transfer; only a typed lock-contention error is transient here. +async fn wait_receive_journal( + dir: &Path, + rel: &str, + manifest: &Manifest, +) -> rds_sync::journal::Journal { + tokio::time::timeout(Duration::from_secs(3), async { + loop { + let (dir, rel, manifest) = (dir.to_owned(), rel.to_owned(), manifest.clone()); + let opened = tokio::task::spawn_blocking(move || { + rds_sync::journal::Journal::open(&dir, &rel, &manifest) + }) + .await + .unwrap(); + match opened { + Ok(journal) => break journal, + Err(rds_sync::SyncError::Io(error)) + if error.kind() == std::io::ErrorKind::WouldBlock => {} + Err(error) => panic!("unexpected receive admission error: {error}"), + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("canceled receive retained its journal lock beyond the cleanup bound") +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn transfer_completes_byte_identical() { let (_s, c_ep, target, _task, server_dir) = pair().await; @@ -354,8 +392,16 @@ async fn kill_mid_transfer_resumes_identical() { tokio::time::sleep(Duration::from_millis(10)).await; } attempt.abort(); - // Let the server observe the drop before reconnecting. - tokio::time::sleep(Duration::from_millis(150)).await; + let _ = attempt.await; + assert!(landed >= 3, "the kill must follow verified part progress"); + // A fixed sleep is not a receiver cleanup fence. Observe the exact root's + // lock and retain its verified progress before making one resumed push. + let journal = wait_receive_journal(&server_dir, "killme.bin", &manifest_of(&data)).await; + assert!( + journal.have_set().len() >= 3, + "canceled parts did not survive" + ); + drop(journal); let stats = push(&c_ep, target.clone(), &src).await.unwrap(); assert_eq!(std::fs::read(server_dir.join("killme.bin")).unwrap(), data); @@ -367,6 +413,35 @@ async fn kill_mid_transfer_resumes_identical() { ); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn receive_admission_waits_for_the_actual_lock_release() { + use rds_sync::journal::Journal; + let dir = scratch("release-fence"); + let manifest = manifest_of(b"retained receiver work"); + let held = Journal::open(&dir, "data.bin", &manifest).unwrap(); + // Reproduce the old assumption: 150ms elapsed, but an executing receiver + // still owns the real lock. This must remain a refusal, not a forced release. + tokio::time::sleep(Duration::from_millis(150)).await; + assert!(matches!(Journal::open(&dir, "data.bin", &manifest), + Err(rds_sync::SyncError::Io(error)) if error.kind() == std::io::ErrorKind::WouldBlock)); + let path = dir.clone(); + let waiting = + tokio::spawn(async move { wait_receive_journal(&path, "data.bin", &manifest).await }); + tokio::time::sleep(Duration::from_millis(50)).await; + assert!( + !waiting.is_finished(), + "admission bypassed the receiver lock" + ); + drop(held); + let journal = tokio::time::timeout(Duration::from_secs(1), waiting) + .await + .unwrap() + .unwrap(); + assert_eq!(journal.need(), [0]); + drop(journal); + std::fs::remove_dir_all(dir).unwrap(); +} + /// G6: kill at randomized progress points until done — byte-identical /// every time. 20 iterations of a small file keep it fast. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] 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..1d59b36 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. @@ -56,6 +67,22 @@ sample with fresh receipts prevent immediate return to a sustained backlog. Media reductions coalesce over one second; simultaneous path/media observations apply the stronger response once. Existing bitrate bounds, frame deadlines, three-frame admission and reference-preserving live encoder updates still apply. + +Loss-only reductions now also respect recent timely transport goodput. Two +adjacent one-second windows must each contain at least three complete timely +frame receipts and 4096 payload bytes; the lower window's rate supplies a +conservative floor with 20% headroom. Late, obsolete and failed frames supply +no such credit. Unknown/changed paths, outstanding late receipts, actual failure +or stale/idle evidence clear the observation. RTT/producer pressure, delivery +holds and negotiated ceilings retain their existing limits. This estimates +confirmed transport delivery, not decoded/displayed quality or total capacity. +See [the qualification boundary](reports/rds-confirmed-goodput-20261002.md). + +With no unconfirmed media, an independent key of at most 64 KiB starts under +QUIC pacing without an extra application wait. Larger keys, dependent frames +and nonempty media retain pacing and the same half-second debt bound. The +existing key-receipt capture barrier remains; waits of at least 100 ms are +logged with frame metadata, never pixels. Sender health includes delayed-delivery counts and latest receipt duration. Independent production health continues during a stopped video writer and records production activity, intentional codec skips, latest produced-frame @@ -73,6 +100,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-confirmed-goodput-20261002.md b/docs/reports/rds-confirmed-goodput-20261002.md new file mode 100644 index 0000000..1bd50c4 --- /dev/null +++ b/docs/reports/rds-confirmed-goodput-20261002.md @@ -0,0 +1,60 @@ +# Preserve confirmed timely media goodput + +Scope: W6.4/W6.7 adaptive media and diagnostic remediation. Native physical +latency, quality and sustained stability acceptance remains open. + +Repeated packet-loss samples could multiply the offered bitrate down to its +floor even while recent media receipts established higher timely delivery. +Reliable QUIC packet loss and actual media goodput measure different things. +The increment credits payload bytes only after successful normal transport +receipt and only when full stream open/write/receipt duration stays inside the +existing delivery-delay budget. Peer disposal, unknown stops, timeouts and late +successes do not contribute. Pacing before opening the stream is still included +in the aggregate wall-time goodput window; the rate does not subtract that wait. + +Two adjacent one-second windows each need at least three timely receipts and +4096 bytes. The smaller measured rate, reduced by 20%, can bound a loss-only +cut. It never exceeds the negotiated ceiling. No ACK, media pressure, delivery +holds, producer misses or suspected/sustained RTT growth retain their original +responses. Unknown/changed path identity, outstanding late receipts, failure, +silence or insufficient traffic clears the observation. This is not total +network capacity, codec correctness or physical display timing. + +INFO health adds timely byte/receipt counts and the eligible delivery floor. +These are metadata only. No protocol/identity/authorization format changes. + +An independent key up to 64 KiB can start without an extra application pacing +wait when no media receipt is unconfirmed. QUIC still paces packets and the +key-receipt barrier bounds following capture. Half-second debt bounds remain; +larger keys, dependent frames and nonempty media keep ordinary pacing. This +removes an artificial initial wait (up to 500 ms at the bitrate floor), not +transmission time. Waits of at least 100 ms gain frame-metadata diagnostics. + +## Regression evidence + +- A seeded 15% packet-loss replay with timely confirmed goodput keeps at least + that conservative observed floor. The equivalent original response without + the evidence falls below it; the fixture distinguishes both behaviors. +- Failure, missing delivery, producer pressure and RTT growth cannot use an + asserted goodput floor to suppress their cuts. Negotiated ceilings hold. +- Explicit-clock window fixtures cover two-window admission, unequal rates, + idle/sparse traffic, path changes, missing selection and expired evidence. +- Receipt classification gives no timely-byte credit to a late completion; + real Iroh/Noq obsolete/unknown-stop regressions additionally check zero credit. +- The admission-cadence fixture exposed by the prior macOS CI is repaired with + the production resume/advance functions on an explicit clock. It no longer + assumes the OS schedules two real-time calls within one 16 ms frame slot. + +The combined macOS run initially retained an independent v1 sync failure: +`control stream ended`. Its isolated investigation passed. The negotiated-sync +pair had default public discovery/port mapping; it is now explicitly local UDP +with discovery disabled, and its failure reports both client and server reasons. +No accepted operation is retried and no failure assertion is removed. The +original negative log remains retained. Final platform and public CI results +are recorded with the candidate PR; incomplete lanes are never called green. + +```sh +cargo fmt --check +cargo test --locked -p rds-desktop -p rds-client -p rds-cli -p rds-sync --all-features +cargo clippy --locked --workspace --all-targets --all-features -- -D warnings +``` 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/reports/rds-sync-release-fence-20261002.md b/docs/reports/rds-sync-release-fence-20261002.md new file mode 100644 index 0000000..97ae5ce --- /dev/null +++ b/docs/reports/rds-sync-release-fence-20261002.md @@ -0,0 +1,34 @@ +# Sync cancellation fixture: observe receiver admission + +Scope: repair an independent W2.5/W8 fixture assumption exposed by the native +viewer PR's full Ubuntu CI; no runtime or wire change. + +The original [Ubuntu failure](https://github.com/NDDev-OpenNetwork/remote-device-sync/actions/runs/36969086539/job/110719084734) +is retained in [issue79](https://github.com/NDDev-OpenNetwork/remote-device-sync/issues/79): +`kill_mid_transfer_resumes_identical` reopened a transfer after a fixed 150 ms +sleep and failed with `offer refused: cannot open transfer journal`. +Canceling the sender does not join an executing receiver filesystem operation; +the exclusive journal lock correctly remains held until that operation ends. + +The fixture now joins the aborted sender, counts only hash-named part files, +requires actual part progress, and probes the receiver's journal admission under +one three-second bound. Only `SyncError::Io(WouldBlock)` is retried; all other +errors fail immediately. The admitted journal verifies retained chunks before +the single resumed network transfer. Byte identity and fewer-than-total fetched +chunks remain required. Accepted transfers are not retried by this helper. + +A real-lock fixture retains a Journal past 150 ms and confirms typed refusal. +The bounded helper remains pending while the lock is held and admits after the +owner releases it. This reproduces the invalid timing assumption with an actual +filesystem lock, without pretending to cancel a running syscall. + +The corrected macOS sync E2E lane passed 15 tests with zero failures or ignores; +strict all-target/all-feature sync lint and formatting passed. Linux and full +CI results accompany the repairing PR. These fixtures do not close physical +power-loss, native desktop latency/quality or sustained stability acceptance. + +```sh +cargo test --locked -p rds-sync --test sync_e2e +cargo clippy --locked -p rds-sync --all-targets --all-features -- -D warnings +cargo fmt --check +``` diff --git a/docs/research.md b/docs/research.md index 7230823..4fa750a 100644 --- a/docs/research.md +++ b/docs/research.md @@ -8,6 +8,30 @@ decisions and a build order. ## 0. Executive summary — what changed vs v0.1 +The October 2 [confirmed-goodput increment](reports/rds-confirmed-goodput-20261002.md) +uses timely frame-receipt bytes to constrain loss-only bitrate reductions. +The maintained [WebRTC LossBasedBweV2](https://webrtc.googlesource.com/src/+/refs/heads/main/modules/congestion_controller/goog_cc/loss_based_bwe_v2.cc) +separates inherent loss, acknowledged rate and delay estimates; its +`CalculateInstantLowerBound` can retain a rate backed by acknowledged traffic. +RDS uses its own smaller conservative observation rule, not a port of GCC: +two populated recent windows, 20% headroom, selected-path scope and immediate +revocation under actual media/RTT pressure. Static application-limited traffic +is not capacity proof. [QUIC loss recovery](https://www.rfc-editor.org/rfc/rfc9002.html) +still controls the underlying transport independently. Increasing encoder +bitrate without delivery evidence can increase queues; [FQ-CoDel](https://www.rfc-editor.org/rfc/rfc8290.html) +documents why short per-flow queues matter to interactive traffic. Physical +route bandwidth, NAT/migration and native long-session acceptance remain open. + +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. |