diff --git a/crates/rds-cli/src/desktop.rs b/crates/rds-cli/src/desktop.rs index 1214bca..2a5d9a3 100644 --- a/crates/rds-cli/src/desktop.rs +++ b/crates/rds-cli/src/desktop.rs @@ -475,9 +475,8 @@ mod native { Some(rds_core::DesktopEvent::Heartbeat { ts_ms, .. }) => view .control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), Some(rds_core::DesktopEvent::InputAck { seq, .. }) => view.input_ack(seq), - Some(rds_core::DesktopEvent::ClipboardReady { bytes, .. }) => { - view.clipboard_ready(bytes); - tracing::info!(bytes, "remote clipboard ready"); + Some(rds_core::DesktopEvent::ClipboardReady { id, bytes }) => { + view.clipboard_ack(id, bytes); } None => { tracing::warn!("managed desktop event channel ended"); @@ -566,7 +565,7 @@ mod native { event = session.events.recv() => match event { Some(rds_core::DesktopEvent::Heartbeat { ts_ms,.. }) => view.control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), Some(rds_core::DesktopEvent::InputAck { seq,.. }) => view.input_ack(seq), - Some(rds_core::DesktopEvent::ClipboardReady { bytes,.. }) => {view.clipboard_ready(bytes);tracing::info!(bytes,"remote clipboard ready");}, + Some(rds_core::DesktopEvent::ClipboardReady { id, bytes }) => {view.clipboard_ack(id, bytes);}, None => break Ok(false), }, frame = session.frames.recv() => match frame { diff --git a/crates/rds-cli/src/desktop/diagnostics.rs b/crates/rds-cli/src/desktop/diagnostics.rs index 484036f..596cea2 100644 --- a/crates/rds-cli/src/desktop/diagnostics.rs +++ b/crates/rds-cli/src/desktop/diagnostics.rs @@ -18,6 +18,7 @@ pub(super) struct Recorder { pending: Option, previous_reconnects: Option, previous_slow_acks: Option, + previous_slow_clipboards: Option, cooldown_until_ms: u64, } @@ -36,6 +37,10 @@ impl Recorder { .previous_slow_acks .is_some_and(|previous| snapshot.report.slow_input_acks > previous); self.previous_slow_acks = Some(snapshot.report.slow_input_acks); + let slow_clipboard = self + .previous_slow_clipboards + .is_some_and(|previous| snapshot.report.slow_clipboard_transfers > previous); + self.previous_slow_clipboards = Some(snapshot.report.slow_clipboard_transfers); if let Some(incident) = &mut self.pending { // Keep memory bounded even if the diagnostic timer runs rapidly. if incident.snapshots.len() < HISTORY + 6 { @@ -55,6 +60,9 @@ impl Recorder { // Completed stalls can fall entirely between periodic snapshots. reasons.push("input_ack_delayed"); } + if slow_clipboard { + reasons.push("clipboard_ready_delayed"); + } if snapshot .report .oldest_input_ack_age_ms @@ -62,6 +70,13 @@ impl Recorder { { reasons.push("input_ack_pending"); } + if snapshot + .report + .clipboard_oldest_pending_age_ms + .is_some_and(|ms| ms >= 250) + { + reasons.push("clipboard_ready_pending"); + } if snapshot.decoded_frame_age_ms.is_some_and(|ms| ms >= 3000) { reasons.push("decoded_video_stalled"); } @@ -195,6 +210,28 @@ mod tests { assert_eq!(data["reasons"][0], "input_ack_delayed"); } + #[test] + fn pending_and_completed_slow_clipboard_transfers_preserve_an_incident() { + for completed in [false, true] { + let mut recorder = Recorder::default(); + recorder.observe(&snapshot(0)); + let mut delayed = snapshot(2000); + let reason = if completed { + delayed.report.slow_clipboard_transfers = 1; + delayed.report.last_clipboard_transfer_ms = Some(600.); + "clipboard_ready_delayed" + } else { + delayed.report.clipboard_pending_transfers = 1; + delayed.report.clipboard_oldest_pending_age_ms = Some(600); + "clipboard_ready_pending" + }; + recorder.observe(&delayed); + let data: serde_json::Value = + serde_json::from_slice(&recorder.finish(true).unwrap()).unwrap(); + assert_eq!(data["reasons"], serde_json::json!([reason])); + } + } + #[test] fn idle_and_a_fresh_update_do_not_report_a_visible_renderer_stall() { let mut recorder = Recorder::default(); diff --git a/crates/rds-desktop/src/render/viewer.rs b/crates/rds-desktop/src/render/viewer.rs index 482e30f..a6fc69d 100644 --- a/crates/rds-desktop/src/render/viewer.rs +++ b/crates/rds-desktop/src/render/viewer.rs @@ -139,6 +139,12 @@ pub struct ViewerReport { pub managed_events_separated: bool, pub input_ack_p50_ms: Option, pub input_ack_p95_ms: Option, + pub keyboard_input_ack_p50_ms: Option, + pub keyboard_input_ack_p95_ms: Option, + pub keyboard_input_ack_max_ms: Option, + pub button_input_ack_p50_ms: Option, + pub button_input_ack_p95_ms: Option, + pub button_input_ack_max_ms: Option, pub input_queue_p95_ms: Option, pub input_pointer_coalesced: u64, pub input_events_dropped: u64, @@ -147,6 +153,16 @@ pub struct ViewerReport { pub oldest_input_ack_age_ms: Option, pub clipboard_transfers: u64, pub last_clipboard_bytes: u32, + pub clipboard_transfer_p50_ms: Option, + pub clipboard_transfer_p95_ms: Option, + pub clipboard_transfer_max_ms: Option, + pub last_clipboard_transfer_ms: Option, + pub slow_clipboard_transfers: u64, + pub clipboard_pending_transfers: usize, + pub clipboard_oldest_pending_age_ms: Option, + pub clipboard_unmatched_replies: u64, + pub clipboard_tracking_evicted: u64, + pub clipboard_transfers_canceled: u64, pub reconnects: u64, pub video_repair_requests: u64, pub last_recovery_ms: Option, @@ -203,29 +219,81 @@ impl SubmissionDebt { /// Local event-to-ack measurements never compare clocks on different hosts. /// Coalesced pointer events are tracked only after dequeue; the event's own /// local creation time still includes time spent waiting in the UI queue. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum InputClass { + Keyboard, + Button, + Other, +} +impl InputClass { + fn of(kind: &InputKind) -> Self { + match kind { + InputKind::KeyDown { .. } | InputKind::KeyUp { .. } => Self::Keyboard, + InputKind::PointerButton { .. } => Self::Button, + _ => Self::Other, + } + } +} + #[derive(Default)] struct InputLatency { - pending: VecDeque<(u64, u64)>, + pending: VecDeque<(u64, u64, InputClass)>, queued: Vec, acknowledged: Vec, + keyboard: Vec, + buttons: Vec, evicted: u64, last_slow_log_ms: Option, } + +#[derive(Default)] +struct ClipboardLatency { + pending: VecDeque<(u64, u32, u64)>, + acknowledged: Vec, + evicted: u64, +} +impl ClipboardLatency { + fn sent(&mut self, id: u64, bytes: u32, now_ms: u64) { + if self.pending.len() == 8 { + self.pending.pop_front(); + self.evicted += 1; + } + self.pending.push_back((id, bytes, now_ms)); + } + fn ack(&mut self, id: u64, bytes: u32, now_ms: u64) -> Option { + let index = self + .pending + .iter() + .position(|(pending, size, _)| *pending == id && *size == bytes)?; + let (_, _, sent_ms) = self.pending.remove(index)?; + let latency = now_ms.saturating_sub(sent_ms) as f64; + sample(&mut self.acknowledged, latency); + Some(latency) + } +} impl InputLatency { - fn sent(&mut self, seq: u64, created_ms: u64, now_ms: u64) { + fn sent(&mut self, seq: u64, created_ms: u64, now_ms: u64, class: InputClass) { if self.pending.len() == INPUT_QUEUE_CAPACITY { self.pending.pop_front(); self.evicted += 1; } - self.pending.push_back((seq, created_ms)); + self.pending.push_back((seq, created_ms, class)); sample(&mut self.queued, now_ms.saturating_sub(created_ms) as f64); } fn ack(&mut self, seq: u64, now_ms: u64) -> Option { - if let Some(index) = self.pending.iter().position(|(pending, _)| *pending == seq) - && let Some((_, created_ms)) = self.pending.remove(index) + if let Some(index) = self + .pending + .iter() + .position(|(pending, _, _)| *pending == seq) + && let Some((_, created_ms, class)) = self.pending.remove(index) { let latency = now_ms.saturating_sub(created_ms) as f64; sample(&mut self.acknowledged, latency); + match class { + InputClass::Keyboard => sample(&mut self.keyboard, latency), + InputClass::Button => sample(&mut self.buttons, latency), + InputClass::Other => {} + } return Some(latency); } None @@ -243,6 +311,7 @@ struct State { encoding: Vec, sending: Vec, input_latency: InputLatency, + clipboard_latency: ClipboardLatency, close: bool, interrupted: Option, network_stage: String, @@ -364,6 +433,9 @@ impl ViewerHandle { if state.status == "Reconnecting" { state.report.input_acks_canceled += state.input_latency.pending.len() as u64; state.input_latency.pending.clear(); + state.report.clipboard_transfers_canceled += + state.clipboard_latency.pending.len() as u64; + state.clipboard_latency.pending.clear(); if let Some(probe) = &mut state.visual_probe { probe.reset(); } @@ -407,9 +479,12 @@ impl ViewerHandle { { let mut state = lock(&self.state); state.report.inputs_dispatched += 1; - state - .input_latency - .sent(event.seq, event.event_ts_ms, sent_ms); + state.input_latency.sent( + event.seq, + event.event_ts_ms, + sent_ms, + InputClass::of(&event.kind), + ); } let event_class = match event.kind { InputKind::KeyDown { .. } => "key_down", @@ -465,6 +540,24 @@ impl ViewerHandle { state.report.clipboard_transfers += 1; state.report.last_clipboard_bytes = bytes; } + /// Correlate explicit paste publication using only metadata and local clocks. + pub fn clipboard_ack(&self, id: u64, bytes: u32) { + let now_ms = self.started.elapsed().as_millis() as u64; + let mut state = lock(&self.state); + let latency = state.clipboard_latency.ack(id, bytes, now_ms); + if latency.is_some() { + state.report.clipboard_transfers += 1; + state.report.last_clipboard_bytes = bytes; + state.report.last_clipboard_transfer_ms = latency; + if latency.is_some_and(|ms| ms >= 250.) { + state.report.slow_clipboard_transfers += 1; + } + } else { + state.report.clipboard_unmatched_replies += 1; + } + drop(state); + tracing::info!(transfer_id=id, bytes, transfer_ms=?latency, "native clipboard ready observed"); + } /// Sender stage durations use only that sender's monotonic clock. They /// are separate from network transit and local receive-to-submit timing. pub fn media_timing(&self, header: &rds_core::FrameHeader) { @@ -497,12 +590,49 @@ impl ViewerHandle { (report.encode_to_send_p50_ms, report.encode_to_send_p95_ms) = quantiles(&state.sending); (report.input_ack_p50_ms, report.input_ack_p95_ms) = quantiles(&state.input_latency.acknowledged); + ( + report.keyboard_input_ack_p50_ms, + report.keyboard_input_ack_p95_ms, + ) = quantiles(&state.input_latency.keyboard); + report.keyboard_input_ack_max_ms = state + .input_latency + .keyboard + .iter() + .copied() + .reduce(f64::max); + ( + report.button_input_ack_p50_ms, + report.button_input_ack_p95_ms, + ) = quantiles(&state.input_latency.buttons); + report.button_input_ack_max_ms = + state.input_latency.buttons.iter().copied().reduce(f64::max); report.input_queue_p95_ms = quantiles(&state.input_latency.queued).1; + ( + report.clipboard_transfer_p50_ms, + report.clipboard_transfer_p95_ms, + ) = quantiles(&state.clipboard_latency.acknowledged); + report.clipboard_transfer_max_ms = state + .clipboard_latency + .acknowledged + .iter() + .copied() + .reduce(f64::max); + report.clipboard_pending_transfers = state.clipboard_latency.pending.len(); + report.clipboard_oldest_pending_age_ms = + state + .clipboard_latency + .pending + .front() + .map(|(_, _, since)| { + (self.started.elapsed().as_millis() as u64).saturating_sub(*since) + }); + report.clipboard_tracking_evicted = state.clipboard_latency.evicted; report.input_ack_tracking_evicted = state.input_latency.evicted; report.pending_input_acks = state.input_latency.pending.len(); - report.oldest_input_ack_age_ms = state.input_latency.pending.front().map(|(_, created)| { - (self.started.elapsed().as_millis() as u64).saturating_sub(*created) - }); + report.oldest_input_ack_age_ms = + state.input_latency.pending.front().map(|(_, created, _)| { + (self.started.elapsed().as_millis() as u64).saturating_sub(*created) + }); report.input_pointer_coalesced = self.input_state.pointer_coalesced.load(Ordering::Relaxed); report.input_events_dropped = self.input_state.input_dropped.load(Ordering::Relaxed); report.input_queue_max_depth = self.input_state.max_depth.load(Ordering::Relaxed); @@ -550,6 +680,7 @@ impl Viewer { encoding: Vec::new(), sending: Vec::new(), input_latency: InputLatency::default(), + clipboard_latency: ClipboardLatency::default(), close: false, interrupted: None, network_stage: "starting".into(), @@ -893,7 +1024,13 @@ impl ApplicationHandler<()> for App { Ok(Some(text)) => { let id = rand::random(); let total = text.len() as u32; + lock(&self.handle.state).clipboard_latency.sent( + id, + total, + self.handle.started.elapsed().as_millis() as u64, + ); tracing::info!( + transfer_id = id, bytes = total, command_paste, "explicit local clipboard text queued" @@ -1030,6 +1167,25 @@ impl ApplicationHandler<()> for App { mod tests { use super::*; #[test] + fn clipboard_measurements_match_size_and_id_without_replaying_duplicates() { + let mut latency = ClipboardLatency::default(); + latency.sent(7, 4096, 100); + assert_eq!(latency.ack(8, 4096, 180), None); + assert_eq!(latency.ack(7, 4095, 190), None); + assert_eq!(latency.pending.len(), 1); + assert_eq!(latency.ack(7, 4096, 200), Some(100.)); + assert_eq!(latency.ack(7, 4096, 250), None); + assert_eq!(latency.acknowledged, vec![100.]); + for id in 10..30 { + latency.sent(id, 4096, 300); + } + assert_eq!(latency.pending.len(), 8); + assert_eq!(latency.evicted, 12); + assert_eq!(latency.ack(10, 4096, 400), None); + latency.pending.clear(); // An old epoch cannot complete a fresh paste. + assert_eq!(latency.ack(29, 4096, 500), None); + } + #[test] fn submission_debt_distinguishes_idle_replacement_and_concurrent_arrival() { let mut debt = SubmissionDebt::default(); assert_eq!(debt.age(90_000), None); // A still screen is not a renderer stall. @@ -1059,23 +1215,48 @@ mod tests { #[test] fn input_latency_matches_sequences_and_keeps_bounded_local_clock_history() { let mut latency = InputLatency::default(); - latency.sent(7, 100, 120); - latency.sent(8, 110, 125); + latency.sent(7, 100, 120, InputClass::Keyboard); + latency.sent(8, 110, 125, InputClass::Button); assert_eq!(latency.ack(8, 310), Some(200.)); assert_eq!(latency.ack(8, 410), None); // duplicate is not another sample assert_eq!(latency.ack(999, 510), None); // unrelated ACK cannot correlate assert_eq!(latency.acknowledged, vec![200.]); assert_eq!(latency.queued, vec![20., 15.]); - assert_eq!(latency.pending.front(), Some(&(7, 100))); + assert_eq!( + latency.pending.front(), + Some(&(7, 100, InputClass::Keyboard)) + ); for seq in 100..1500 { - latency.sent(seq, seq, seq + 3); + latency.sent(seq, seq, seq + 3, InputClass::Other); } assert_eq!(latency.pending.len(), INPUT_QUEUE_CAPACITY); assert_eq!(latency.queued.len(), 1024); - assert_eq!(latency.pending.front(), Some(&(476, 476))); + assert_eq!( + latency.pending.front(), + Some(&(476, 476, InputClass::Other)) + ); assert_eq!(latency.evicted, 377); } + #[test] + fn fast_pointer_acknowledgements_do_not_hide_slow_keyboard_and_buttons() { + let mut latency = InputLatency::default(); + latency.sent(1, 0, 0, InputClass::Keyboard); + assert_eq!(latency.ack(1, 900), Some(900.)); + latency.sent(2, 0, 0, InputClass::Button); + assert_eq!(latency.ack(2, 700), Some(700.)); + for seq in 3..2000 { + latency.sent(seq, 1000, 1000, InputClass::Other); + latency.ack(seq, 1060); + } + assert_eq!(quantiles(&latency.acknowledged).1, Some(60.)); + assert_eq!(quantiles(&latency.keyboard).1, Some(900.)); + assert_eq!(quantiles(&latency.buttons).1, Some(700.)); + assert_eq!(latency.ack(1, 2000), None); + assert_eq!(latency.keyboard.len(), 1); + assert_eq!(latency.buttons.len(), 1); + } + #[tokio::test] async fn pointer_collapse_keeps_click_and_release_order_and_close_wakes_receiver() { let state = Arc::new(InputState { @@ -1186,7 +1367,7 @@ mod tests { fn semantic_burst_acknowledgements_remain_correlated_when_replies_wait() { let mut latency = InputLatency::default(); for seq in 0..800 { - latency.sent(seq, seq, 800); + latency.sent(seq, seq, 800, InputClass::Button); } assert_eq!(latency.pending.len(), 800); assert_eq!(latency.evicted, 0); diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 5e54544..fb728a4 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -32,6 +32,9 @@ const PACING_INTERVAL: Duration = Duration::from_millis(250); /// Capture briefly after accepted input even when a compositor's root DAMAGE /// notification does not describe the redirected application repaint. const INPUT_REFRESH_BURST_MS: u64 = 250; +// Control replies must not inherit a media stream's thirty-second budget. +// After any partial-write failure the session drops SessionSend and resets it. +const CONTROL_REPLY_TIMEOUT: Duration = Duration::from_secs(2); /// Moderate random loss holds the offered rate; severe loss reduces it. const LOSS_HOLD: f64 = 0.02; const LOSS_STEP_DOWN: f64 = 0.10; @@ -1293,14 +1296,7 @@ pub async fn serve_desktop_with( seq, handled_ts_ms: send_clock.now_ms(), }; - if !matches!( - tokio::time::timeout( - FRAME_SEND_TIMEOUT, - write_frame(&mut send.0, &ack) - ) - .await, - Ok(Ok(())) - ) { + if write_control_reply(&mut send.0, &ack).await.is_err() { break; } tracing::trace!(target:"rds_desktop::input_timing", input_seq=seq, @@ -1327,14 +1323,10 @@ pub async fn serve_desktop_with( controls.requested.store(bps, Ordering::Relaxed); } Ok(DesktopControl::Heartbeat { seq, ts_ms }) => { - if !matches!( - tokio::time::timeout( - FRAME_SEND_TIMEOUT, - write_frame(&mut send.0, &DesktopEvent::Heartbeat { seq, ts_ms }) - ) - .await, - Ok(Ok(())) - ) { + if write_control_reply(&mut send.0, &DesktopEvent::Heartbeat { seq, ts_ms }) + .await + .is_err() + { break; } } @@ -1372,13 +1364,26 @@ pub async fn serve_desktop_with( } }; if let Some(text) = text { + let publish_started = Instant::now(); + tracing::info!( + transfer_id = id, + bytes = total, + "desktop clipboard publication started" + ); let owner = clipboard .get_or_insert_with(|| crate::clipboard::Worker::new(session_display)); match tokio::time::timeout(Duration::from_secs(2), owner.publish(text)) .await { Ok(Ok(())) => { - if write_frame( + tracing::info!( + transfer_id = id, + bytes = total, + publish_ms = publish_started.elapsed().as_millis(), + "desktop clipboard publication completed" + ); + let reply_started = Instant::now(); + if write_control_reply( &mut send.0, &DesktopEvent::ClipboardReady { id, bytes: total }, ) @@ -1387,6 +1392,12 @@ pub async fn serve_desktop_with( { break; } + tracing::info!( + transfer_id = id, + bytes = total, + reply_ms = reply_started.elapsed().as_millis(), + "desktop clipboard ready reply written" + ); } result => { tracing::warn!(error=?result,"clipboard publication failed; ending control before paste input"); @@ -1417,6 +1428,30 @@ pub async fn serve_desktop_with( result } +async fn write_control_reply( + send: &mut W, + event: &DesktopEvent, +) -> std::io::Result<()> { + let reply_class = match event { + DesktopEvent::InputAck { .. } => "input_ack", + DesktopEvent::Heartbeat { .. } => "heartbeat", + DesktopEvent::ClipboardReady { .. } => "clipboard_ready", + }; + match tokio::time::timeout(CONTROL_REPLY_TIMEOUT, write_frame(send, event)).await { + Ok(result) => result, + Err(_) => { + tracing::warn!( + reply_class, + "desktop control reply deadline exceeded; ending session" + ); + Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + "desktop control reply deadline exceeded", + )) + } + } +} + /// A canceled control reply must not end with a partial, apparently clean FIN. struct SessionSend(SendStream); @@ -2126,6 +2161,61 @@ async fn send_payload>( mod tests { use super::*; + #[tokio::test(start_paused = true)] + async fn blocked_clipboard_reply_obeys_control_deadline_after_partial_header() { + let (mut send, mut receive) = tokio::io::duplex(2); + let started = tokio::time::Instant::now(); + let result = write_control_reply( + &mut send, + &DesktopEvent::ClipboardReady { id: 7, bytes: 4096 }, + ) + .await; + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::TimedOut); + assert_eq!(started.elapsed(), CONTROL_REPLY_TIMEOUT); + drop(send); // A real serving session resets SessionSend at this boundary. + use tokio::io::AsyncReadExt; + let mut partial = Vec::new(); + receive.read_to_end(&mut partial).await.unwrap(); + assert_eq!(partial.len(), 2, "the control frame was partially written"); + } + + #[tokio::test] + async fn timely_clipboard_reply_preserves_following_input_ack_and_heartbeat() { + let (mut send, mut receive) = tokio::io::duplex(256); + for event in [ + DesktopEvent::ClipboardReady { id: 7, bytes: 4096 }, + DesktopEvent::InputAck { + seq: 8, + handled_ts_ms: 9, + }, + DesktopEvent::Heartbeat { seq: 10, ts_ms: 11 }, + ] { + write_control_reply(&mut send, &event).await.unwrap(); + let received: DesktopEvent = read_frame(&mut receive).await.unwrap(); + match (received, event) { + ( + DesktopEvent::ClipboardReady { id: a, bytes: b }, + DesktopEvent::ClipboardReady { id: c, bytes: d }, + ) => assert_eq!((a, b), (c, d)), + ( + DesktopEvent::InputAck { + seq: a, + handled_ts_ms: b, + }, + DesktopEvent::InputAck { + seq: c, + handled_ts_ms: d, + }, + ) => assert_eq!((a, b), (c, d)), + ( + DesktopEvent::Heartbeat { seq: a, ts_ms: b }, + DesktopEvent::Heartbeat { seq: c, ts_ms: d }, + ) => assert_eq!((a, b), (c, d)), + _ => panic!("control reply changed type or order"), + } + } + } + #[test] fn accepted_input_wake_survives_capture_backpressure_then_expires() { let controls = ProducerControls::new(4_000_000); diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 41bb625..33bcb92 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -15,6 +15,18 @@ window, alongside independent rate-limited warnings. `session_epoch`, `input_queue_depth`, `input_acks_canceled`, `slow_input_acks` and `last_input_ack_ms` complement the existing bounded latency percentiles. Percentiles summarize the retained sample buffer, not a timed rolling interval. +Keyboard and button ACK p50/p95/max are retained independently of motion, so +fast pointer traffic cannot hide a delayed key or click. Each series retains +at most 1024 samples. These remain injection acknowledgements, not application +response or optical display measurements. + +Explicit paste gestures track at most eight transfer IDs and byte counts on +the viewer's local clock. Clipboard-ready replies must match both values; +duplicates, unknown IDs and wrong sizes do not create successful samples. +Reconnection cancels outstanding tracking. Evictions, unmatched replies, +pending age, last/p50/p95/max transfer duration and completed transfers over +250 ms are reported. Pending and completed clipboard delays trigger the same +bounded incident recorder; no clipboard content is retained. ## Validated payload delivery receipts @@ -330,11 +342,20 @@ owns the CLIPBOARD selection and serves UTF8_STRING/TARGETS/TIMESTAMP and ICCCM INCR for larger data, without a clipboard helper process. This is real clipboard publication, not typing text through keyboard-layout substitutions. -One transfer is bounded to 1 MiB of UTF-8, with 32 KiB control chunks, exact -ordered offsets and a five-second assembly deadline. There are at most four +One transfer is bounded to 1 MiB of UTF-8, with at most 32 KiB control chunks +(the native sender uses 16 KiB), exact ordered offsets, a five-second idle +deadline and a thirty-second total assembly deadline. There are at most four active native selection workers and eight outstanding INCR requests per worker. View-only sessions refuse publication. Publication failure ends the control session before subsequent paste input can consume an unrelated old clipboard. +Clipboard publication has its own two-second native deadline. All serving +control replies (input ACK, heartbeat and clipboard-ready) share a two-second +write deadline rather than the thirty-second media budget. A stalled or partial +reply ends and resets that desktop control stream; it is never resumed as a +clean frame and no input or paste is replayed automatically. The transport +connection and unrelated service streams remain outside that session shutdown. +Publication start/completion and reply completion record transfer ID, byte +count and stage duration only. Contents are neither logged nor written to disk. There is no background scan or automatic export of every local clipboard change. Images, files, rich formats, reverse clipboard remain outside this text path. Cmd+V becomes a bounded remote diff --git a/docs/reports/rds-paste-control-backpressure-20261005.md b/docs/reports/rds-paste-control-backpressure-20261005.md new file mode 100644 index 0000000..172bb2b --- /dev/null +++ b/docs/reports/rds-paste-control-backpressure-20261005.md @@ -0,0 +1,51 @@ +# Desktop paste/control backpressure — 2026-10-05 + +This W6.4/W6.7 and O4 increment bounds control replies and improves native input +observation. It does not close installed network, application-response or +sustained stability acceptance. + +## Implemented boundary + +The serving clipboard-ready reply previously awaited an unbounded framed write. +Input ACKs and heartbeat replies inherited the thirty-second media budget. All +three now use a separate two-second write budget. Failure exits the owning +desktop control loop; `SessionSend` resets the stream, so a partial record cannot +be resumed or presented as a clean FIN. Media, authorization, connection and +unrelated service-stream policies remain unchanged. There is no input or paste +replay and no wire-format change. + +Clipboard publication start/completion and ready-reply completion include only +transfer ID, byte count and stage duration. The viewer correlates exact ID/size +with local gesture-to-ready timing, at most eight outstanding entries and 1024 +completed samples. Wrong, duplicate and retired replies cannot create successful +samples. Evictions, cancellations and unmatched replies remain counted. Pending +and completed delays over 250 ms trigger bounded before/after incident windows. +Keyboard and button ACK series are independent of motion and each bounded to +1024 samples; the existing aggregate series remains available. + +## Regression scope + +- A real two-byte-capacity Tokio pipe holds a partially written clipboard reply: + it returns `TimedOut` at two seconds, rather than awaiting indefinitely. +- A timely pipe preserves clipboard-ready, input-ACK and heartbeat records in + order with exact values. +- Nearly 2000 fast motion acknowledgements can push a slow keyboard/button sample + out of the aggregate series; their independent series still report 900/700 ms. +- Clipboard tracking refuses wrong IDs/sizes and duplicates, bounds outstanding + work, counts evictions and rejects replies after epoch cancellation. +- A completed slow paste entirely between diagnostic ticks still records an + incident, as does a pending paste; metadata contains no clipboard payload. + +These synthetic regressions establish the local boundary. They do not establish +that an unbounded write caused a particular live complaint, that a target GUI +consumed the paste, or that submitted GPU frames reached physical scanout. + +## Validation status + +The desktop viewer library's 100 unit tests passed after the implementation. +An initial test build referenced an unavailable test-only serializer; it was +corrected to compare decoded control variants without adding a dependency. +An initial CLI invocation used a nonexistent feature name; subsequent checks +must use the actual `rds-cli/desktop` feature. Formatting, strict workspace/native +clippy, full workspace tests, Linux X11 and installed qualification remain to be +recorded before this increment is considered complete. diff --git a/docs/research.md b/docs/research.md index 6d2ee99..bfd373e 100644 --- a/docs/research.md +++ b/docs/research.md @@ -1,5 +1,22 @@ # Deep research: remote + sync, all-Rust, minimum latency +## 2026-10-05 bounded paste/control observation + +QUIC [stream flow control](https://www.rfc-editor.org/rfc/rfc9000.html#section-4) +can block an application writer independently of average RTT. A higher stream +priority does not remove that wait. [Tokio timeout](https://docs.rs/tokio/latest/tokio/time/fn.timeout.html) +limits the awaited write but cancellation can leave a partial framed record. +RDS therefore resets its owning desktop control stream after a reply failure, +using a separate two-second control budget for ACK, heartbeat and clipboard +metadata; it does not resume or replay that write. Media budgets are unchanged. + +Keyboard/button samples are separate from pointer motion, and clipboard +publication is correlated by ID, size and the viewer's local clock. This +closes an unbounded clipboard-ready write and an observability gap, not a +claim that those defects explain every native freeze. The synthetic blocked +writer and tracking tests remain distinct from installed GUI and network +qualification. + ## 2026-10-05 explicit media receipt boundary The pinned [Noq stopped contract](https://docs.rs/noq/1.3.0/noq/struct.SendStream.html#method.stopped)