diff --git a/crates/rds-desktop/src/client.rs b/crates/rds-desktop/src/client.rs index 8e2ebb3..e0820f7 100644 --- a/crates/rds-desktop/src/client.rs +++ b/crates/rds-desktop/src/client.rs @@ -44,6 +44,45 @@ const FRAME_STREAM_TIMEOUT: std::time::Duration = std::time::Duration::from_secs const FRAME_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(3); const KEYFRAME_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(8); +fn reorder_wait(rtt_ms: u64) -> std::time::Duration { + // A missing stream tag may itself be in transit. Two observed round trips + // tolerate WAN reordering; the cap keeps a truly missing reference bounded. + let ms = if rtt_ms == u64::MAX { + 250 + } else { + rtt_ms.saturating_mul(2).clamp(100, 1000) + }; + std::time::Duration::from_millis(ms) +} + +#[derive(Default)] +struct HeartbeatProbes { + pending: std::collections::VecDeque<(u64, u64, std::time::Instant)>, +} + +impl HeartbeatProbes { + fn sent(&mut self, seq: u64, ts_ms: u64) { + // Caller timestamps are opaque correlation values. In managed mode + // they belong to the viewer, whose clock survives session reconnects. + self.pending.retain(|(s, t, _)| (*s, *t) != (seq, ts_ms)); + if self.pending.len() >= 64 { + self.pending.pop_front(); + } + self.pending + .push_back((seq, ts_ms, std::time::Instant::now())); + } + + fn echoed(&mut self, seq: u64, ts_ms: u64) -> Option { + let index = self + .pending + .iter() + .position(|(s, t, _)| (*s, *t) == (seq, ts_ms))?; + self.pending + .remove(index) + .map(|(_, _, sent)| sent.elapsed()) + } +} + /// Decode state owned by one session. Codec reference frames are chain /// state: a decoder shared across sessions would cross-contaminate /// streams, so it lives here and is dropped with the session. @@ -365,8 +404,13 @@ impl DesktopSession { // Construct ownership before spawning. Dropping the supervisor aborts // its JoinSet, which in turn drops streams and the frame-reader JoinSet. let mut tasks = JoinSet::new(); + let probes = Arc::new(tokio::sync::Mutex::new(HeartbeatProbes::default())); + let sending_probes = probes.clone(); tasks.spawn(async move { while let Some(msg) = ctrl_rx.recv().await { + if let DesktopControl::Heartbeat { seq, ts_ms } = &msg { + sending_probes.lock().await.sent(*seq, *ts_ms); + } if !matches!( tokio::time::timeout(FRAME_STREAM_TIMEOUT, write_frame(&mut send.0, &msg)) .await, @@ -378,15 +422,13 @@ impl DesktopSession { }); let rtt_marker = control_rtt_ms.clone(); - let event_clock = clock.clone(); tasks.spawn(async move { loop { match read_frame::<_, DesktopEvent>(&mut recv).await { - Ok(ev @ DesktopEvent::Heartbeat { ts_ms, .. }) => { - rtt_marker.store( - event_clock.now_ms().saturating_sub(ts_ms), - Ordering::Relaxed, - ); + Ok(ev @ DesktopEvent::Heartbeat { seq, ts_ms }) => { + if let Some(rtt) = probes.lock().await.echoed(seq, ts_ms) { + rtt_marker.store(rtt.as_millis() as u64, Ordering::Relaxed); + } events_tx.send(ev); } Ok(ev) => { @@ -406,6 +448,7 @@ impl DesktopSession { ctrl: ctrl_tx.clone(), clock: clock.clone(), receiving: receiving.clone(), + control_rtt_ms: control_rtt_ms.clone(), }, )); let task = tokio::spawn(async move { @@ -642,6 +685,7 @@ struct ReceiveContext { ctrl: mpsc::Sender, clock: SessionClock, receiving: Arc, + control_rtt_ms: Arc, } enum FrameRead { @@ -726,8 +770,13 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) { tokio::select! { biased; _ = repair.tick() => { - if ordered.expire(std::time::Duration::from_millis(100)) { + // An admitted reference has its own bounded reader deadline. + // Before its tag arrives, allow a measured WAN reorder window + // rather than generating an IDR at the old fixed 100 ms. + let wait = reorder_wait(ctx.control_rtt_ms.load(Ordering::Relaxed)); + if readers.is_empty() && ordered.expire(wait) { gaps+=1; + tracing::warn!(expected_seq=ctx.next_seq.load(Ordering::Relaxed),reorder_budget_ms=wait.as_millis(),"desktop reference gap expired"); delivery.invalidate(); delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms()); } @@ -867,6 +916,38 @@ pub async fn run_desktop_client( mod tests { use super::*; + #[test] + fn heartbeat_measurements_are_bounded_and_require_exact_correlation() { + let mut probes = HeartbeatProbes::default(); + for seq in 0..128 { + probes.sent(seq, 1_000_000 + seq); + } + assert_eq!(probes.pending.len(), 64); + assert!(probes.echoed(0, 1_000_000).is_none()); + assert!(probes.echoed(127, 0).is_none()); + assert!(probes.echoed(126, 1_000_127).is_none()); + assert!(probes.echoed(127, 1_000_127).is_some()); + assert!(probes.echoed(127, 1_000_127).is_none()); + probes.sent(64, 1_000_064); + assert_eq!( + probes + .pending + .iter() + .filter(|(s, t, _)| (*s, *t) == (64, 1_000_064)) + .count(), + 1 + ); + } + + #[test] + fn reference_wait_uses_rtt_and_caps_untrusted_or_missing_measurements() { + assert_eq!(reorder_wait(0).as_millis(), 100); + assert_eq!(reorder_wait(200).as_millis(), 400); + assert_eq!(reorder_wait(50_000).as_millis(), 1000); + assert_eq!(reorder_wait(u64::MAX - 1).as_millis(), 1000); + assert_eq!(reorder_wait(u64::MAX).as_millis(), 250); + } + #[test] fn exhausted_control_sequences_never_wrap() { let sequence = AtomicU64::new(u64::MAX - 1); diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 711c12e..e4bc41f 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -37,6 +37,35 @@ const MAX_PENDING_FRAME_ACKS: usize = 3; const FRAME_ACK_TIMEOUT: Duration = Duration::from_secs(5); const KEYFRAME_ACK_TIMEOUT: Duration = Duration::from_secs(10); +// QUIC path counters may remain clean while a reliable relay queues media. +// Observe actual frame delivery as well, without retaining frame payloads. +#[derive(Default)] +struct DeliveryFeedback { + delayed: AtomicU64, + failed: AtomicU64, + acknowledged: AtomicU64, + last_ack_ms: AtomicU64, + producing: AtomicBool, + produced: AtomicU64, + codec_skips: AtomicU64, + last_produced_ms: AtomicU64, +} + +#[derive(Clone)] +struct FrameDelivery { + keyframe_pending: Arc, + idr: Arc, + feedback: Arc, +} + +fn delivery_delay_budget(path: Option) -> Duration { + path.map_or(Duration::from_millis(250), |p| { + p.rtt + .saturating_mul(3) + .clamp(Duration::from_millis(250), Duration::from_secs(1)) + }) +} + /// Monotonic clock shared by producer and writer so `FrameHeader` /// timestamps are comparable within one session. #[derive(Clone)] @@ -184,6 +213,8 @@ pub struct BitrateController { last_congestion: u64, primed: bool, last_path: Option, + delivery_hold_ticks: u8, + delivery_cut_cooldown_ticks: u8, } impl BitrateController { @@ -198,6 +229,8 @@ impl BitrateController { last_congestion: 0, primed: false, last_path: None, + delivery_hold_ticks: 0, + delivery_cut_cooldown_ticks: 0, } } @@ -264,6 +297,44 @@ impl BitrateController { self.current = next; next } + + // The serving session supplements path samples with frame ACKs. A late + // frame reduces offered load before its hard reset deadline. Hold that + // reduction for five seconds; coalesce a burst of correlated receipts for + // 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. + fn step_with_delivery( + &mut self, + path: Option, + deadline_misses: u64, + impaired: bool, + delivered: bool, + ) -> u64 { + let previous = self.current; + let proposed = self.step(path, deadline_misses); + self.delivery_cut_cooldown_ticks = self.delivery_cut_cooldown_ticks.saturating_sub(1); + self.current = if impaired { + self.delivery_hold_ticks = 20; + if self.delivery_cut_cooldown_ticks == 0 { + self.delivery_cut_cooldown_ticks = 4; + // RTT/loss and a delayed receipt may report the same event. + // Apply the stronger response once, never multiply both cuts. + let media_cut = (previous / 10 * 7 + previous % 10 * 7 / 10).max(self.floor); + proposed.min(media_cut) + } else { + proposed.min(previous) + } + } else if self.delivery_hold_ticks > 0 { + self.delivery_hold_ticks -= 1; + proposed.min(previous) + } else if !delivered { + proposed.min(previous) + } else { + proposed.min(previous.saturating_add((previous / 100).max(1))) + }; + self.current + } } /// Serve one desktop session on an already-accepted stream pair. @@ -325,6 +396,7 @@ pub async fn serve_desktop_with( let (tx, mut rx) = mpsc::channel::(2); let keyframe_pending = Arc::new(AtomicBool::new(false)); let capture_admission = Arc::new(AtomicU64::new(0)); + let delivery_feedback = Arc::new(DeliveryFeedback::default()); { let clock = clock.clone(); let bitrate = Arc::clone(&controls.bitrate); @@ -334,6 +406,7 @@ pub async fn serve_desktop_with( let mut producer = config.producer; let keyframe_pending = keyframe_pending.clone(); let capture_admission = capture_admission.clone(); + let feedback = delivery_feedback.clone(); capture.spawn_blocking(move || { let producer_controls = ProducerControls { bitrate, @@ -386,11 +459,21 @@ pub async fn serve_desktop_with( if paused { source.resume_after_backpressure(); } - match source.produce(seq, &producer_controls, &clock) { + feedback.producing.store(true, Ordering::Relaxed); + let result = source.produce(seq, &producer_controls, &clock); + feedback.producing.store(false, Ordering::Relaxed); + match result { Some(p) => { if p.payload.is_empty() && source.preserves_reference() { + feedback.codec_skips.fetch_add(1, Ordering::Relaxed); continue; } + if !p.payload.is_empty() { + feedback + .last_produced_ms + .store(clock.now_ms(), Ordering::Relaxed); + feedback.produced.fetch_add(1, Ordering::Relaxed); + } if p.header.keyframe && !p.payload.is_empty() { keyframe_pending.store(true, Ordering::Release); } @@ -417,15 +500,23 @@ pub async fn serve_desktop_with( }) }; - // Pacing: sample path counters + deadline misses into the controller, + // Pacing: sample path counters + actual media delivery + deadline misses, // which writes the bitrate the producer reads each frame. { let conn = conn.clone(); let bitrate = Arc::clone(&controls.bitrate); let requested = Arc::clone(&controls.requested); let misses = Arc::clone(&controls.deadline_misses); + let feedback = delivery_feedback.clone(); + let admission = capture_admission.clone(); + let pending_key = keyframe_pending.clone(); + let progress_clock = clock.clone(); let mut controller = BitrateController::new(4_000_000, ceiling); workers.spawn(async move { + let mut last_delayed = 0; + let mut last_failed = 0; + let mut last_acknowledged = 0; + let mut health = Instant::now(); let mut tick = tokio::time::interval(PACING_INTERVAL); tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { @@ -439,12 +530,26 @@ pub async fn serve_desktop_with( let missed = misses.swap(0, Ordering::Relaxed); let previous = controller.current(); let path = conn.current_path_stats(); - let bps = controller.step(path, missed); + let delayed = feedback.delayed.load(Ordering::Relaxed); + let failed = feedback.failed.load(Ordering::Relaxed); + let acknowledged = feedback.acknowledged.load(Ordering::Relaxed); + let delayed_frames = delayed.saturating_sub(last_delayed); + let failed_frames = failed.saturating_sub(last_failed); + let bps = controller.step_with_delivery( + path, + missed, + delayed_frames > 0 || failed_frames > 0, + acknowledged > last_acknowledged, + ); + (last_delayed, last_failed, last_acknowledged) = (delayed, failed, acknowledged); if bps != previous { tracing::debug!( previous_bps = previous, bitrate_bps = bps, deadline_misses = missed, + delayed_frames, + failed_frames, + last_ack_ms = feedback.last_ack_ms.load(Ordering::Relaxed), path_rtt_ms = ?path.map(|p| p.rtt.as_millis()), path_id = ?path.map(|p| p.path_id), path_via_relay = ?path.map(|p| p.via_relay), @@ -455,6 +560,26 @@ pub async fn serve_desktop_with( ); } bitrate.store(bps.min(u64::from(u32::MAX)), Ordering::Relaxed); + // Independent of frame sends: during a freeze, distinguish + // native production, codec skips and delivery backpressure. + if health.elapsed() >= Duration::from_secs(5) { + let produced = feedback.produced.load(Ordering::Relaxed); + tracing::info!( + producing = feedback.producing.load(Ordering::Relaxed), + produced, + codec_skips = feedback.codec_skips.load(Ordering::Relaxed), + last_produced_age_ms = ?(produced > 0).then(|| progress_clock.now_ms().saturating_sub(feedback.last_produced_ms.load(Ordering::Relaxed))), + pending_media_frames = admission.load(Ordering::Acquire), + keyframe_pending = pending_key.load(Ordering::Acquire), + bitrate_bps = bps, + delayed_delivery = delayed, + failed_delivery = failed, + last_ack_ms = feedback.last_ack_ms.load(Ordering::Relaxed), + path_rtt_ms = ?path.map(|p| p.rtt.as_millis()), + "desktop production and delivery health" + ); + health = Instant::now(); + } } }) }; @@ -466,6 +591,7 @@ pub async fn serve_desktop_with( let writer_clock = clock.clone(); let writer_bitrate = Arc::clone(&controls.bitrate); let writer_idr = Arc::clone(&controls.idr); + let writer_feedback = delivery_feedback.clone(); let frame_route = config.frame_route.unwrap_or(rds_core::UniHello::Desktop); workers.spawn(async move { // Token bucket on the paced bitrate: offering faster than the @@ -554,8 +680,11 @@ pub async fn serve_desktop_with( produced, &mut rx, &mut acknowledgements, - keyframe_pending.clone(), - writer_idr.clone(), + FrameDelivery { + keyframe_pending: keyframe_pending.clone(), + idr: writer_idr.clone(), + feedback: writer_feedback.clone(), + }, ) .await { @@ -585,6 +714,8 @@ pub async fn serve_desktop_with( pending_media_frames = capture_admission.load(Ordering::Acquire), keyframe_pending = keyframe_pending.load(Ordering::Acquire), bitrate_bps = writer_bitrate.load(Ordering::Relaxed), + delayed_delivery = writer_feedback.delayed.load(Ordering::Relaxed), + last_ack_ms = writer_feedback.last_ack_ms.load(Ordering::Relaxed), "desktop sender health" ); health = Instant::now(); @@ -1128,20 +1259,11 @@ async fn send_frame( produced: AdmittedFrame, rx: &mut mpsc::Receiver, acknowledgements: &mut JoinSet, - keyframe_pending: Arc, - idr: Arc, + delivery: FrameDelivery, ) -> SendOutcome { match tokio::time::timeout( FRAME_SEND_TIMEOUT, - send_frame_inner( - conn, - route, - produced, - rx, - acknowledgements, - keyframe_pending, - idr, - ), + send_frame_inner(conn, route, produced, rx, acknowledgements, delivery), ) .await { @@ -1159,8 +1281,7 @@ async fn send_frame_inner( produced: AdmittedFrame, rx: &mut mpsc::Receiver, acknowledgements: &mut JoinSet, - keyframe_pending: Arc, - idr: Arc, + delivery: FrameDelivery, ) -> SendOutcome { let AdmittedFrame { produced, permit } = produced; let mut sending = match conn.open_uni().await { @@ -1209,19 +1330,40 @@ async fn send_frame_inner( let seq = produced.header.seq; let keyframe = produced.header.keyframe; let payload_bytes = produced.payload.len(); + let delay_budget = delivery_delay_budget(conn.current_path_stats()); + let FrameDelivery { + keyframe_pending, + idr, + feedback, + } = delivery; // Retain the reset-on-drop owner until delivery is acknowledged. The // bounded task group is owned by this writer; cancellation resets its // outstanding frames without closing unrelated connection services. acknowledgements.spawn(async move { let started = Instant::now(); let deadline = if keyframe {KEYFRAME_ACK_TIMEOUT} else {FRAME_ACK_TIMEOUT}; - let acknowledged = match tokio::time::timeout(deadline, sending.stream.stopped()).await { + let result = { + let receipt = tokio::time::timeout(deadline, sending.stream.stopped()); + tokio::pin!(receipt); + tokio::select! { + result = &mut receipt => result, + _ = tokio::time::sleep(delay_budget) => { + feedback.delayed.fetch_add(1, Ordering::Relaxed); + tracing::warn!(frame_seq=seq,payload_bytes,delay_budget_ms=delay_budget.as_millis(),"desktop frame delivery delayed"); + receipt.await + } + } + }; + feedback.last_ack_ms.store(started.elapsed().as_millis() as u64, Ordering::Relaxed); + let acknowledged = match result { Ok(Ok(None)) => { sending.finished = true; + feedback.acknowledged.fetch_add(1, Ordering::Relaxed); tracing::trace!(frame_seq=seq,payload_bytes,ack_ms=started.elapsed().as_millis(),"desktop frame transport acknowledged"); true } result => { + feedback.failed.fetch_add(1, Ordering::Relaxed); idr.store(true, Ordering::Relaxed); tracing::warn!(frame_seq=seq,payload_bytes,ack_ms=started.elapsed().as_millis(),outcome=?result,"desktop frame delivery unconfirmed"); false @@ -1643,6 +1785,78 @@ mod tests { assert_eq!(bps, 8_000_000, "clean windows must reach the ceiling"); } + #[test] + fn delayed_media_reduces_load_despite_clean_relay_counters_and_recovers_cautiously() { + let mut c = BitrateController::new(4_000_000, 8_000_000); + let mut sample = path(1000, 0, 80, 0); + sample.via_relay = true; + c.step_with_delivery(Some(sample), 0, false, true); + sample.sent += 100; + let reduced = c.step_with_delivery(Some(sample), 0, true, false); + assert_eq!(reduced, 2_800_000); + // Fresh ACKs must not undo the cut during its five-second hold. + for _ in 0..20 { + sample.sent += 100; + assert_eq!(c.step_with_delivery(Some(sample), 0, false, true), reduced); + } + // Packet activity without a media receipt cannot justify growth. + sample.sent += 100; + assert_eq!(c.step_with_delivery(Some(sample), 0, false, false), reduced); + sample.sent += 100; + assert_eq!( + c.step_with_delivery(Some(sample), 0, false, true), + 2_828_000 + ); + assert!(c.step_with_delivery(None, 0, true, false) < 2_828_000); + } + + #[test] + fn repeated_media_failure_honors_the_existing_floor_and_ceiling() { + let mut c = BitrateController::new(200_000, 300_000); + for _ in 0..30 { + c.step_with_delivery(None, 0, true, false); + } + assert_eq!(c.current(), 100_000); + c.steer(1_000_000); + assert_eq!(c.current(), 300_000); + } + + #[test] + fn one_delivery_burst_does_not_compound_frame_and_path_penalties() { + let mut c = BitrateController::new(4_000_000, 8_000_000); + c.step_with_delivery(Some(path(1000, 0, 80, 0)), 0, false, true); + // A path event and a delayed receipt describe the same congestion. + let reduced = c.step_with_delivery(Some(path(1100, 0, 160, 1)), 0, true, false); + assert_eq!(reduced, 2_800_000, "one event must not apply two 30% cuts"); + // Up to three outstanding frames can report the same burst on + // successive pacing ticks. Keep their first cut rather than cubing it. + for i in 0..3 { + assert_eq!( + c.step_with_delivery(Some(path(1200 + i * 100, 0, 160, 1)), 0, true, false), + reduced, + "correlated receipts over-penalized image quality" + ); + } + // Continued pressure after a full second must still reduce load. + assert_eq!( + c.step_with_delivery(Some(path(1600, 0, 160, 1)), 0, true, false), + 1_960_000 + ); + } + + #[test] + fn media_delay_budget_accounts_for_propagation_but_remains_bounded() { + assert_eq!(delivery_delay_budget(None), Duration::from_millis(250)); + assert_eq!( + delivery_delay_budget(Some(path(0, 0, 200, 0))), + Duration::from_millis(600) + ); + assert_eq!( + delivery_delay_budget(Some(path(0, 0, 9000, 0))), + Duration::from_secs(1) + ); + } + #[test] fn controller_honors_floor_and_misses() { let mut c = BitrateController::new(200_000, 8_000_000); diff --git a/crates/rds-desktop/tests/client_lifecycle.rs b/crates/rds-desktop/tests/client_lifecycle.rs index 9e28501..c23f9a4 100644 --- a/crates/rds-desktop/tests/client_lifecycle.rs +++ b/crates/rds-desktop/tests/client_lifecycle.rs @@ -10,6 +10,10 @@ use rds_net::{ Backend, Connection, Endpoint, EndpointConfig, RecvStream, SendStream, read_frame, write_frame, }; +// These scenarios intentionally consume the process-wide eight-reader budget. +// Isolate fixtures while retaining concurrency within each real connection. +static SESSION_TEST: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); + async fn pair(backend: Backend) -> (Endpoint, Endpoint, Connection, Connection) { let config = EndpointConfig { backend, @@ -94,6 +98,7 @@ async fn rejected_header(b: &Connection, session: [u8; 16], h: FrameHeader) { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn completing_a_delta_first_preserves_the_inflight_keyframe() { + let _isolation = SESSION_TEST.lock().await; for backend in [Backend::Iroh, Backend::Noq] { let (client_ep, server_ep, a, b) = pair(backend).await; let (mut client, control_send, control_recv, id) = session(&a, &b).await; @@ -107,7 +112,7 @@ async fn completing_a_delta_first_preserves_the_inflight_keyframe() { delta.write_all(b"later delta").await.unwrap(); delta.finish().unwrap(); assert!( - tokio::time::timeout(Duration::from_millis(30), client.frame_headers.recv()) + tokio::time::timeout(Duration::from_millis(350), client.frame_headers.recv()) .await .is_err() ); @@ -128,8 +133,160 @@ async fn completing_a_delta_first_preserves_the_inflight_keyframe() { } } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_progressing_delta_reference_survives_the_reorder_window() { + let _isolation = SESSION_TEST.lock().await; + for backend in [Backend::Iroh, Backend::Noq] { + let (client_ep, server_ep, a, b) = pair(backend).await; + let (mut client, control_send, control_recv, id) = session(&a, &b).await; + let mut keyframe = tagged(&b, id).await; + write_frame(&mut keyframe, &header(0)).await.unwrap(); + keyframe.write_all(b"initial keyframe").await.unwrap(); + keyframe.finish().unwrap(); + assert_eq!(client.frame_headers.recv().await.unwrap().seq, 0); + + let mut reference = tagged(&b, id).await; + let mut h = header(1); + h.keyframe = false; + write_frame(&mut reference, &h).await.unwrap(); + reference.write_all(b"reference prefix").await.unwrap(); + until(|| client.receive_stats().in_flight == 1).await; + let mut successor = tagged(&b, id).await; + h.seq = 2; + write_frame(&mut successor, &h).await.unwrap(); + successor.write_all(b"successor").await.unwrap(); + successor.finish().unwrap(); + assert!( + tokio::time::timeout(Duration::from_millis(350), client.frame_headers.recv()) + .await + .is_err() + ); + reference.write_all(b" reference tail").await.unwrap(); + reference.finish().unwrap(); + for seq in [1, 2] { + assert_eq!( + tokio::time::timeout(Duration::from_secs(2), client.frame_headers.recv()) + .await + .unwrap() + .unwrap() + .seq, + seq, + "a later completed frame discarded its admitted reference" + ); + } + drop((client, control_send, control_recv)); + tokio::join!(client_ep.close(), server_ep.close()); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn wan_reference_header_can_arrive_after_the_lan_reorder_window() { + let _isolation = SESSION_TEST.lock().await; + for backend in [Backend::Iroh, Backend::Noq] { + let (client_ep, server_ep, a, b) = pair(backend).await; + let (client, streams) = tokio::join!( + DesktopSession::connect_opts( + &a, + DesktopHello { + display: 0, + codec: Codec::H264, + max_fps: 30, + input_acks: false + }, + SessionOpts { + session: Some([73; 16]), + relay_encoded: true, + ..Default::default() + }, + ), + async { + let (mut send, mut recv) = b.accept_bi().await.unwrap(); + let StreamHello::DesktopV2 { session, .. } = read_frame(&mut recv).await.unwrap() + else { + panic!("isolated desktop expected") + }; + write_frame( + &mut send, + &HelloAck::Desktop(DesktopCaps { + displays: vec![], + codecs: vec![Codec::H264], + }), + ) + .await + .unwrap(); + (send, recv, session) + } + ); + let mut client = client.unwrap(); + let (mut control_send, mut control_recv, id) = streams; + // The managed viewer's clock predates this desktop session after + // reconnect. Echo timestamps remain caller-owned on the wire. + client + .send_control(DesktopControl::Heartbeat { + seq: 734, + ts_ms: 1_000_000, + }) + .await + .unwrap(); + let DesktopControl::Heartbeat { seq, ts_ms } = read_frame(&mut control_recv).await.unwrap() + else { + panic!("heartbeat expected") + }; + tokio::time::sleep(Duration::from_millis(200)).await; + write_frame(&mut control_send, &DesktopEvent::Heartbeat { seq, ts_ms }) + .await + .unwrap(); + until(|| { + client + .control_rtt() + .is_some_and(|rtt| rtt >= Duration::from_millis(200)) + }) + .await; + + let mut key = tagged(&b, id).await; + write_frame(&mut key, &header(0)).await.unwrap(); + key.write_all(b"initial").await.unwrap(); + key.finish().unwrap(); + assert_eq!(client.frame_headers.recv().await.unwrap().seq, 0); + let mut successor = tagged(&b, id).await; + let mut h = header(2); + h.keyframe = false; + write_frame(&mut successor, &h).await.unwrap(); + successor.write_all(b"successor").await.unwrap(); + successor.finish().unwrap(); + // No reference reader has been admitted yet. A WAN-sized gap in tag + // arrival must not cause an IDR before its missing stream arrives. + assert!( + tokio::time::timeout( + Duration::from_millis(250), + read_frame::<_, DesktopControl>(&mut control_recv) + ) + .await + .is_err() + ); + let mut reference = tagged(&b, id).await; + h.seq = 1; + write_frame(&mut reference, &h).await.unwrap(); + reference.write_all(b"reference").await.unwrap(); + reference.finish().unwrap(); + for expected in [1, 2] { + assert_eq!( + tokio::time::timeout(Duration::from_secs(2), client.frame_headers.recv()) + .await + .unwrap() + .unwrap() + .seq, + expected + ); + } + drop((client, control_send, control_recv)); + tokio::join!(client_ep.close(), server_ep.close()); + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn slow_keyframe_completes_while_delta_deadline_stays_short() { + let _isolation = SESSION_TEST.lock().await; for backend in [Backend::Iroh, Backend::Noq] { let (client_ep, server_ep, a, b) = pair(backend).await; let (mut client, control_send, control_recv, id) = session(&a, &b).await; @@ -168,6 +325,7 @@ async fn slow_keyframe_completes_while_delta_deadline_stays_short() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn readers_are_bounded_and_session_owns_all_streams() { + let _isolation = SESSION_TEST.lock().await; for backend in [Backend::Iroh, Backend::Noq] { tokio::time::timeout(Duration::from_secs(20), async { let (client_ep, server_ep, a, b) = pair(backend).await; diff --git a/crates/rds-desktop/tests/frame_delivery_budget.rs b/crates/rds-desktop/tests/frame_delivery_budget.rs index 596cb92..e5ae6fc 100644 --- a/crates/rds-desktop/tests/frame_delivery_budget.rs +++ b/crates/rds-desktop/tests/frame_delivery_budget.rs @@ -108,8 +108,12 @@ async fn endpoint() -> (Endpoint, Gate) { struct CountingSource { inner: rds_desktop::SyntheticProducer, calls: Arc, + bitrate: Arc, } impl rds_desktop::FrameProducer for CountingSource { + fn resume_after_backpressure(&mut self) { + self.inner.resume_after_backpressure(); + } fn produce( &mut self, seq: u64, @@ -117,6 +121,10 @@ impl rds_desktop::FrameProducer for CountingSource { clock: &rds_desktop::SessionClock, ) -> Option { self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + self.bitrate.store( + controls.bitrate.load(std::sync::atomic::Ordering::Relaxed), + std::sync::atomic::Ordering::SeqCst, + ); self.inner.produce(seq, controls, clock) } } @@ -130,13 +138,17 @@ async fn blocked_acknowledgements_bound_capture_and_resume_without_closing_conne use std::sync::atomic::{AtomicU64, Ordering}; tokio::time::timeout(Duration::from_secs(10), async { - let (server, _) = endpoint().await; - let (client, gate) = endpoint().await; + // Pause before server datagrams are emitted. This creates delayed + // media without deliberately losing packets or receiver ACKs. + let (server, gate) = endpoint().await; + let (client, _) = endpoint().await; let (a, b) = tokio::join!(client.connect(server.addr(), rds_core::ALPN), async { server.accept().await.unwrap().await }); let (a, b) = (a.unwrap(), b.unwrap()); let calls = Arc::new(AtomicU64::new(0)); + let bitrate = Arc::new(AtomicU64::new(0)); + let observed_bitrate = bitrate.clone(); let produced = calls.clone(); let (start, ready) = tokio::sync::oneshot::channel(); let serving = tokio::spawn(async move { @@ -165,6 +177,7 @@ async fn blocked_acknowledgements_bound_capture_and_resume_without_closing_conne producer: Some(Box::new(CountingSource { inner: SyntheticProducer::new(240, 64, 64, 1024).keyframe_every(1), calls: produced, + bitrate: observed_bitrate, })), frame_route: Some(UniHello::DesktopFrames { id: session }), ..Default::default() @@ -196,7 +209,8 @@ async fn blocked_acknowledgements_bound_capture_and_resume_without_closing_conne } tokio::time::sleep(Duration::from_millis(100)).await; let held = calls.load(Ordering::SeqCst); - tokio::time::sleep(Duration::from_millis(300)).await; + let initial_bitrate = bitrate.load(Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(800)).await; assert_eq!( held, 1, "a recovery keyframe must finish before encoding successors" @@ -214,6 +228,10 @@ async fn blocked_acknowledgements_bound_capture_and_resume_without_closing_conne last = encoded.recv().await.unwrap().header.seq; } assert!(last > first, "delivery must resume on the same connection"); + assert!( + bitrate.load(Ordering::SeqCst) < initial_bitrate, + "delayed media ACKs must reduce the encoder target even with no reported packet loss" + ); drop(session); serving.abort(); let _ = serving.await; diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 848be5f..6ebee59 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -37,8 +37,28 @@ requests `AutoNoVsync` and one frame of latency, with the backend's supported fallback. Submission timing does not prove when a physical pixel becomes visible. Completed encoded frames are briefly ordered before decode so a delta that finishes first cannot discard its in-flight reference keyframe. Ordering -retains at most three successor bodies; a gap lasting 100 ms requests bounded -IDR recovery. The existing global encoded/decode limits still apply. +retains at most three successor bodies. The missing-reference timer uses two measured control RTTs, bounded to +100–1000 ms (250 ms before measurement). It requests IDR recovery only after +admitted readers finish or reach their own bounded deadlines. Heartbeat RTT +uses bounded local send/echo correlations; caller timestamps stay opaque, so +reconnecting a managed desktop does not compare two different clock origins. Jitter does not discard successors of a reference still +being received. The existing global encoded/decode limits still apply. + +The sender also measures frame delivery receipts independently of QUIC packet +loss counters, which may look clean while a reliable relay queues traffic. +A receipt delayed beyond three sampled path RTTs (bounded to 250–1000 ms), or +a failed delivery, reduces the encoder target. A five-second recovery hold and +growth of at most 1% per 250 ms sample with fresh media receipts prevent an +immediate return to the same backlog. Media reductions coalesce correlated +receipts over one second. When path and media observations describe the same +sample, the stronger reduction applies once, preserving quality while continued +pressure still reduces load. Existing bitrate bounds, frame deadlines, +three-frame admission and reference-preserving live encoder updates still apply. +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 +age, admission/keyframe waits and path RTT; +diagnostics contain frame metadata, never pixels or clipboard contents. The viewer reopens an interrupted desktop channel with a fresh session route. When the managed connection disappeared, it reconnects the same pinned peer diff --git a/docs/reports/rds-delivery-feedback-20261001.md b/docs/reports/rds-delivery-feedback-20261001.md new file mode 100644 index 0000000..f20b24d --- /dev/null +++ b/docs/reports/rds-delivery-feedback-20261001.md @@ -0,0 +1,123 @@ +# Desktop delivery feedback and reference ordering — 2026-10-01 + +This increment advances W6.4/W6.6 and the W10 diagnostic work. It does not close +hardware encoding, physical glass-to-glass latency, or the complete impairment +and platform gates. + +## Failure and change + +Live diagnostics showed frame ACK deadlines and a partly received recovery +keyframe while reliable-relay packet loss counters remained clean. The encoder +continued targeting its ceiling. Separately, the receiver expired completed +successors after 100 ms even when an admitted reference reader was still active, +causing avoidable reference loss and new large keyframes. + +Ordering now retains at most three successors until admitted readers complete +or reach their existing bounded deadlines. A missing reference with no active reader +uses two measured control RTTs, bounded to 100–1000 ms, with a 250 ms initial +fallback. Reader/payload/decode limits are unchanged. + +The sender observes each owned frame stream's delivery receipt. A soft delay +threshold of three sampled path RTTs, bounded to 250–1000 ms, reports an impaired +media sample without resetting that frame. Hard ACK/read deadlines are unchanged. +An impaired sample reduces target bitrate by 30%; a five-second hold prevents +immediate reversal. Correlated media receipts coalesce over one second; a +simultaneous path penalty is combined by taking the stronger response once. +Continued pressure after that second can reduce load again. Recovery requires +a new successful frame receipt and grows +by at most 1% per 250 ms sample. Existing floor, grant ceiling, three-frame +capture admission and live reference-preserving encoder updates remain in force. +No protocol, identity, permission, or dependency change is introduced. + +Normal sender health records delayed-delivery counts and the latest receipt +wait duration. Delay warnings identify sequence and payload length only; no +screen or clipboard content enters diagnostics. Every receipt task and its +reset-on-cancel stream remain owned by the existing bounded writer task group. + +## Verification + +The real Iroh/Noq regression with a delayed delta reference failed on the prior +implementation: its successor was discarded and delivery timed out. With the +change it delivers both reference and successor in order. The existing delayed +initial-keyframe case now waits 350 ms, beyond the faulty 100 ms window. Both +cases, bounded hostile readers, and independent key/delta deadlines pass. + +The gated real-UDP test verifies one held keyframe, bounded capture and +same-connection resumption. Its first rate assertion also passed the prior +controller and did not isolate media feedback. The refined fixture forwards +cadence-resume to its wrapped producer and pauses server egress before datagrams +are emitted, rather than dropping receiver ACKs. With the production pacing +call restored to the old path-only controller, the rate assertion fails; with +media feedback enabled it passes. Both outcomes are retained. Controller tests +also cover clean relay counters plus delayed media, the recovery hold, no growth +without fresh ACKs, bounded growth, and floor/ceiling preservation. The +process-wide hostile-reader fixtures are serialized so independent scenarios do +not consume each other's global budget; concurrency and budget assertions within +each scenario are retained. + +At the initial production commit, Mac whole all-feature workspace tests passed: +**775 passed, 0 failed, 2 ignored** +across 107 result groups. Strict all-feature/all-target workspace clippy, +formatting and cargo-deny advisories/bans/licenses/sources pass. The two ignored +cases remain outside this increment's acceptance. The fixture refinement changed +only tests. Subsequent media-burst adaptation is qualified separately below. + +At that initial production commit, Linux whole all-feature workspace tests passed: +**777 passed, 0 failed, +11 ignored**, across 107 result groups, together with strict all-feature +workspace/all-target clippy and formatting. Ignored native/display/account cases +are not counted as passes. These checks ran on the production commit; the +subsequent test-fixture refinement has its separate targeted receipt. + +The longer run exposed excessive rate reductions from correlated receipts and +overlapping path/media signals. A regression reproduced two 30% cuts for one +sample (4 Mbit/s became 1.96 instead of 2.8 Mbit/s). The controller now combines +the responses once and coalesces media cuts for one second. The regression also +verifies that sustained pressure after that second still lowers offered load. +The full desktop test suite passes the new code; updated whole-workspace and +installed-device observations will be recorded after qualification. + +Installed-device observations belong to the private estate. The prior longer +run contained an unplanned recovery and remains failed stability evidence. A +new installed run is being qualified; a locked local console is excluded from +visible-pixel acceptance. No full milestone gate is marked closed here. + +A real Iroh/Noq test also delays the missing reference tag beyond 100 ms while +measuring a 200 ms control RTT. The former fixed timer requested an unnecessary +IDR before that reference arrived. The receive loop now uses a bounded RTT-based +window and logs the expected sequence and expiration budget on actual gaps. +No payload or identity is added to these diagnostics. The initial test harness +mistakenly requested a legacy route while expecting V2; that failed setup was +retained separately and excluded from the corrected failed-before evidence. + +A separate five-second production health record remains active when the writer +has no new frame. It reports native production activity, intentional codec skip +count, latest produced-frame age, admission/keyframe waits and path RTT. These +metadata distinguish capture/encode inactivity, intentional skips and transport +backpressure during the observed long pauses. Hard budgets and wire remain +unchanged; no image, clipboard content, credential or endpoint identity is logged. + +At the WAN/progress-diagnostics revision, Mac whole all-feature workspace tests +pass **778/0/2 ignored**,107 result groups; Linux whole recheck passes +**780/0/11 ignored**,107 groups. Strict all-feature/all-target workspace clippy +and formatting pass both. One first-run Linux owned-relay closure assertion +failed under host load; the unchanged focused case and full recheck pass, and +both raw outcomes remain recorded. Current both-OS CI also passes. Ignored +platform/account cases remain outside these totals and acceptance claims. + +The independent logs exposed a clock-origin bug in managed heartbeats: verbatim +viewer timestamps survived reconnect, while the desktop receiver subtracted +its own new session clock. RTT saturated to zero and therefore used the100 ms +reorder floor. Heartbeat timing now correlates up to64 sent sequence/timestamp +pairs with local monotonic instants. Echoes preserve the caller's wire values; +unmatched or duplicate echoes cannot invent an RTT. A real managed-style +heartbeat test with a different timestamp origin failed before and passes with +the fix. Bounded-correlation tests cover eviction, exact matching and duplicates. + +After the heartbeat clock-origin repair, both whole all-feature suites pass: +Mac **779 passed, 0 failed, 2 ignored**; Linux **781 passed, 0 failed, 11 ignored**, +107 result groups each. Strict all-feature/all-target clippy and formatting pass. +The scoped real-device qualification remains separate: prior failed motion runs +are retained and the corrected build has a fresh visible native/mixed-workload +run underway. No broad milestone or physical-latency closure is inferred from +these checks.