From 2b1e9a79bb9f0445e93d1861e96212a6fa9337b9 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sat, 26 Sep 2026 15:18:20 +0500 Subject: [PATCH 1/5] fix(desktop): recover reference chains after encoded frame drops --- crates/rds-desktop/src/session.rs | 278 ++++++++++++++++++------- crates/rds-desktop/tests/session_v2.rs | 18 +- 2 files changed, 209 insertions(+), 87 deletions(-) diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 06f3f7c..956a6bf 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -1,11 +1,10 @@ //! Serving side of a desktop session. //! //! Frame delivery follows the MoQ pattern: every encoded frame goes out on -//! its own uni-directional stream carrying a `FrameHeader`, a fresher -//! queued frame always supersedes a stale one, and a stale frame still -//! in flight is reset mid-send rather than finishing on the wire. -//! Keyframes are never superseded — every delta behind them depends on -//! their landing. Input events, encoder steering and heartbeats arrive +//! its own uni-directional stream carrying a `FrameHeader`. Deltas retain +//! their predecessor; only an independent keyframe can replace a delta +//! still in flight. A sequence gap requires a keyframe before delivery +//! resumes. Input events, encoder steering and heartbeats arrive //! on the bi-directional control stream, which outranks every frame //! stream. @@ -265,17 +264,16 @@ pub async fn serve_desktop_with( } match source.produce(seq, &producer_controls, &clock) { Some(p) => { - // Bounded queue: when full the writer is behind - // and this frame is dropped — its successor - // lands fresher. A keyframe carries the pending - // IDR request though, so dropping one re-arms the - // flag instead of losing the request to - // backpressure. - let was_keyframe = p.header.keyframe; - if tx.try_send(p).is_err() && was_keyframe { + // Losing any encoded reference breaks its successors, + // not only losing an IDR. Keep the two-slot bound and + // ask the producer for an independent replacement. + if tx.try_send(p).is_err() { producer_controls.idr.store(true, Ordering::Relaxed); } - seq += 1; + let Some(next) = seq.checked_add(1) else { + return; + }; + seq = next; } None => return, } @@ -302,15 +300,9 @@ pub async fn serve_desktop_with( }) }; - // Writer task: one uni stream per frame. A continuous producer - // means any frame queued behind an in-progress send is already - // stale — the collapse keeps only the newest (decode-aware: a - // queued keyframe always survives since deltas behind it can't - // decode without it), and a send still in flight when a fresher - // frame arrives is reset mid-write rather than allowed to finish - // (MoQ-style stale reset): the client would drop the tail anyway, - // so its unsent bytes only consume path capacity the fresh frame - // needs. + // Writer task: one uni stream per frame. Collapse stale queued work, but + // never send a delta whose predecessor was discarded. A broken chain + // requests an IDR locally instead of waiting for a client roundtrip. let writer_conn = conn.clone(); let writer_clock = clock.clone(); let writer_bitrate = Arc::clone(&controls.bitrate); @@ -318,12 +310,12 @@ pub async fn serve_desktop_with( workers.spawn(async move { // Token bucket on the paced bitrate: offering faster than the // path sustains only backlogs QUIC's send buffer with frames - // that arrive stale — the collapse cannot reach them once - // buffered. Debt is capped at half a second so a large + // that arrive stale. Debt is capped at half a second so a large // keyframe can't stall the writer. let mut budget = 0.0f64; let mut last = Instant::now(); let mut pending: Option = None; + let mut chain = FrameChain::default(); 'writer: loop { let mut produced = match pending.take() { Some(p) => p, @@ -332,11 +324,10 @@ pub async fn serve_desktop_with( None => break, }, }; - produced = collapse(produced, &mut rx); - // Encode-failure placeholders carry no payload: sending one - // decodes to garbage on the client, while a skipped seq is - // what the client's gap→IDR resync is for. + produced = collapse(produced, &mut rx, &writer_idr); if produced.payload.is_empty() { + chain.next = None; + writer_idr.store(true, Ordering::Relaxed); continue; } let bps = writer_bitrate.load(Ordering::Relaxed).max(50_000) as f64 / 8.0; @@ -348,31 +339,20 @@ pub async fn serve_desktop_with( let wait = ((cost - budget) / bps).min(0.5); tokio::time::sleep(Duration::from_secs_f64(wait)).await; budget = (budget - cost).max(-bps * 0.5); - // Frames produced during the pacing wait are fresher — - // collapse once more before committing to the wire. - produced = collapse(produced, &mut rx); - if produced.payload.is_empty() { - continue; - } + produced = collapse(produced, &mut rx, &writer_idr); } else { budget -= cost; } + // Admission follows the final selection: advancing this before + // the pacing wait would lose track of frames collapsed afterward. + if !chain.admit(&produced) { + writer_idr.store(true, Ordering::Relaxed); + continue; + } produced.header.send_ts_ms = writer_clock.now_ms(); match send_frame(&writer_conn, &produced, &mut rx).await { SendOutcome::Sent => {} SendOutcome::Superseded(newer) => pending = Some(newer), - SendOutcome::ResetStale => { - // The dropped tail broke the delta chain — the next - // produced frame must be an IDR, and the queued - // deltas in front of it are undecodable. - writer_idr.store(true, Ordering::Relaxed); - while let Ok(queued) = rx.try_recv() { - if queued.header.keyframe { - pending = Some(queued); - continue 'writer; - } - } - } SendOutcome::Done | SendOutcome::Failed => break 'writer, } } @@ -699,21 +679,57 @@ mod x11 { } } -/// Drain queued frames newest-wins. Decode-aware: once a keyframe is -/// in the mix it absorbs everything — deltas produced after it cannot -/// decode without it, so the keyframe is kept and later deltas are -/// skipped rather than the other way around. -fn collapse(mut produced: Produced, rx: &mut mpsc::Receiver) -> Produced { - let mut have_keyframe = produced.header.keyframe; - while let Ok(newer) = rx.try_recv() { - if newer.header.keyframe || !have_keyframe { - have_keyframe |= newer.header.keyframe; +/// Prefer a recent independent frame, otherwise the latest queued candidate. +/// FrameChain rejects candidates with a missing reference. If a keyframe is +/// retained while later deltas are discarded, request recovery for that gap too. +fn collapse( + mut produced: Produced, + rx: &mut mpsc::Receiver, + idr: &AtomicBool, +) -> Produced { + // The producer queue has two slots. Bound this drain even if the producer + // refills while we select; selection must not starve actual frame writes. + let mut discarded_reference = false; + for _ in 0..2 { + let Ok(newer) = rx.try_recv() else { break }; + if newer.header.keyframe { produced = newer; + discarded_reference = false; + } else { + discarded_reference = true; + if !produced.header.keyframe { + produced = newer; + } } } + if discarded_reference { + idr.store(true, Ordering::Relaxed); + } produced } +/// Conservative reference contract: every delta may depend on its predecessor. +/// A producer must identify independent keyframes from the encoded bitstream. +#[derive(Default)] +struct FrameChain { + next: Option, +} + +impl FrameChain { + fn admit(&mut self, produced: &Produced) -> bool { + let header = &produced.header; + if produced.payload.is_empty() + || header.seq == u64::MAX + || (!header.keyframe && self.next != Some(header.seq)) + { + self.next = None; + return false; + } + self.next = header.seq.checked_add(1); + true + } +} + /// Reset code for a frame stream abandoned mid-send — the frame went /// stale while still in flight, so its tail is dropped instead of /// consuming path capacity the fresher frame needs. @@ -741,10 +757,6 @@ enum SendOutcome { Sent, /// A fresher decodable frame supersedes — send it next. Superseded(Produced), - /// A stale delta was reset mid-send: the reference chain is broken - /// on the client and only an IDR resyncs it, so queued deltas are - /// worthless and the next produced frame must be a keyframe. - ResetStale, /// Producer closed mid-send; the final frame was finished. Done, /// Transport failure — the writer ends. @@ -752,8 +764,8 @@ enum SendOutcome { } /// Send one frame on its own tagged uni stream, aborting mid-write if -/// a fresher frame lands: an in-flight keyframe is finished (the chain -/// behind it depends on it), a stale delta is reset. +/// an independent keyframe lands. Otherwise finish the reference on which +/// the next delta may depend, retaining partial-write progress. async fn send_frame( conn: &Connection, produced: &Produced, @@ -801,7 +813,7 @@ async fn send_frame_inner( } let outcome = match send_payload(stream, produced, rx).await { Ok(PayloadOutcome::Abandoned(next)) => { - return next.map_or(SendOutcome::ResetStale, SendOutcome::Superseded); + return SendOutcome::Superseded(next); } Ok(PayloadOutcome::Sent) => SendOutcome::Sent, Ok(PayloadOutcome::Superseded(next)) => SendOutcome::Superseded(next), @@ -824,7 +836,7 @@ enum PayloadOutcome { Sent, Superseded(Produced), ProducerEnded, - Abandoned(Option), + Abandoned(Produced), } async fn send_payload( @@ -846,13 +858,11 @@ async fn send_payload( writing.await?; Ok(PayloadOutcome::ProducerEnded) } - Some(newer) if produced.header.keyframe => { + Some(newer) if produced.header.keyframe || !newer.header.keyframe => { writing.await?; Ok(PayloadOutcome::Superseded(newer)) } - Some(newer) => Ok(PayloadOutcome::Abandoned( - newer.header.keyframe.then_some(newer) - )), + Some(newer) => Ok(PayloadOutcome::Abandoned(newer)), }, } } @@ -877,7 +887,7 @@ mod tests { } } - async fn partial_write(closed: bool) { + async fn partial_write(closed: bool, keyframe: bool) { use std::future::{Future, poll_fn}; use std::task::Poll; use tokio::io::AsyncReadExt; @@ -885,7 +895,7 @@ mod tests { tokio::time::timeout(Duration::from_secs(3), async { let (mut writer, mut reader) = tokio::io::duplex(64); let (tx, mut rx) = mpsc::channel(1); - let frame = produced(0, true); + let frame = produced(0, keyframe); let mut wire = vec![0; 17]; { let mut sending = Box::pin(send_payload(&mut writer, &frame, &mut rx)); @@ -929,12 +939,18 @@ mod tests { #[tokio::test] async fn partial_keyframe_supersession_preserves_exact_bytes() { - partial_write(false).await; + partial_write(false, true).await; } #[tokio::test] async fn partial_final_frame_preserves_exact_bytes() { - partial_write(true).await; + partial_write(true, true).await; + partial_write(true, false).await; + } + + #[tokio::test] + async fn partial_delta_finishes_before_its_dependent_successor() { + partial_write(false, false).await; } #[tokio::test] @@ -973,11 +989,11 @@ mod tests { } #[tokio::test] - async fn stale_delta_abandons_without_waiting_for_peer() { + async fn independent_keyframe_replaces_delta_without_waiting_for_peer() { use std::future::{Future, poll_fn}; use std::task::Poll; - for next_keyframe in [false, true] { + { let (mut writer, _reader) = tokio::io::duplex(64); let (tx, mut rx) = mpsc::channel(1); let frame = produced(0, false); @@ -987,7 +1003,7 @@ mod tests { Poll::Ready(()) }) .await; - tx.send(produced(1, next_keyframe)).await.unwrap(); + tx.send(produced(1, true)).await.unwrap(); let result = tokio::time::timeout(Duration::from_secs(1), sending) .await .unwrap() @@ -995,11 +1011,117 @@ mod tests { let PayloadOutcome::Abandoned(next) = result else { panic!("stale delta was not abandoned"); }; - assert_eq!(next.is_some(), next_keyframe); - if let Some(next) = next { - assert_eq!(next.header.seq, 1); + assert_eq!(next.header.seq, 1); + assert!(next.header.keyframe); + } + } + + #[test] + fn reference_chain_waits_for_keyframe_after_gap_empty_or_sequence_end() { + let mut chain = FrameChain::default(); + assert!(!chain.admit(&produced(0, false))); + assert!(chain.admit(&produced(1, true))); + assert!(chain.admit(&produced(2, false))); + assert!(!chain.admit(&produced(4, false))); // Missing reference 3. + assert!(!chain.admit(&produced(5, false))); + assert!(chain.admit(&produced(6, true))); + let mut empty = produced(7, false); + empty.payload = Bytes::new(); + assert!(!chain.admit(&empty)); + assert!(!chain.admit(&produced(8, false))); + assert!(chain.admit(&produced(u64::MAX - 1, true))); + assert!(!chain.admit(&produced(u64::MAX, false))); + assert!(!chain.admit(&produced(0, false))); + } + + #[test] + fn collapsed_references_request_recovery_and_never_admit_a_broken_delta() { + let (tx, mut rx) = mpsc::channel(2); + let idr = AtomicBool::new(false); + let mut chain = FrameChain::default(); + assert!(chain.admit(&produced(0, true))); + tx.try_send(produced(2, false)) + .unwrap_or_else(|_| panic!("fixture queue full")); + let selected = collapse(produced(1, false), &mut rx, &idr); + assert_eq!(selected.header.seq, 2); + assert!(!chain.admit(&selected)); + assert!(idr.swap(false, Ordering::Relaxed)); + + // A queued IDR replaces the broken prefix without dropping its own + // references; no redundant request is needed when it is the last item. + tx.try_send(produced(4, true)) + .unwrap_or_else(|_| panic!("fixture queue full")); + let selected = collapse(produced(3, false), &mut rx, &idr); + assert!(chain.admit(&selected)); + assert!(!idr.load(Ordering::Relaxed)); + + // Retaining a keyframe while shedding a successor also loses a + // reference: a later delta must wait for another independent frame. + tx.try_send(produced(6, false)) + .unwrap_or_else(|_| panic!("fixture queue full")); + let selected = collapse(produced(5, true), &mut rx, &idr); + assert!(chain.admit(&selected)); + assert!(idr.load(Ordering::Relaxed)); + assert!(!chain.admit(&produced(7, false))); + } + + #[cfg(feature = "x11")] + #[test] + fn native_h264_chain_recovers_after_a_lost_reference() { + use crate::{Decoder, EncodedFrame, Encoder, H264Decoder, H264Encoder, RawFrame}; + + let mut encoder = H264Encoder::new(4_000_000, 30.0).unwrap(); + let mut decoder = H264Decoder::new().unwrap(); + let mut chain = FrameChain::default(); + let mut decoded = Vec::new(); + for seq in 0..7 { + let mut bgra = vec![0; 64 * 64 * 4]; + for (i, pixel) in bgra.chunks_exact_mut(4).enumerate() { + let value = if (i / 64 + seq as usize * 4) % 32 < 16 { + 220 + } else { + 40 + }; + pixel.copy_from_slice(&[value, value, value, 255]); + } + if seq == 5 { + encoder.request_idr(); + } + let encoded = encoder + .encode(&RawFrame { + width: 64, + height: 64, + stride: 256, + data: bgra.into(), + }) + .unwrap(); + assert!(!encoded.data.is_empty()); + assert_eq!(encoded.keyframe, seq == 0 || seq == 5); + let mut frame = produced(seq, encoded.keyframe); + frame.header.width = 64; + frame.header.height = 64; + frame.payload = encoded.data; + if seq == 2 { + // Simulate an encoded frame lost to backpressure. + continue; + } + if chain.admit(&frame) { + let raw = decoder + .decode(&EncodedFrame { + codec: frame.header.codec, + keyframe: frame.header.keyframe, + data: frame.payload, + }) + .unwrap() + .expect("admitted H.264 frame must decode"); + assert_eq!((raw.width, raw.height), (64, 64)); + // Check a native decoded pixel after recovery, not just a header. + let expected = if (seq * 4) % 32 < 16 { 220i16 } else { 40 }; + assert!((i16::from(raw.data[0]) - expected).abs() < 12); + decoded.push(seq); } } + assert_eq!(decoded, [0, 1, 5, 6]); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] diff --git a/crates/rds-desktop/tests/session_v2.rs b/crates/rds-desktop/tests/session_v2.rs index d350982..e68c9ae 100644 --- a/crates/rds-desktop/tests/session_v2.rs +++ b/crates/rds-desktop/tests/session_v2.rs @@ -3,8 +3,8 @@ //! Two `rds_net` endpoints, one running `serve_desktop_with` against a //! `SyntheticProducer`, the other a `DesktopSession`. A shared //! `SessionClock` makes `FrameHeader` timestamps directly comparable to -//! client receive times, so viewer-visible latency is measured, not -//! inferred. +//! client receive times. This measures complete header arrival, not decoded +//! pixels or viewer presentation. //! //! The impairment lane runs on the owned transport (`noq`) with //! `impair::ImpairingSocket` wrapping each endpoint's UDP socket: @@ -392,9 +392,9 @@ async fn view_only_and_failed_injection_never_ack_but_keep_control_alive() { } /// C5 impairment + G5 latency gate: 5% loss + 30 ms jitter on a 50 ms -/// base — queue stays bounded, stale frames drop, viewer-visible latency -/// p95 stays in the 150 ms budget and control RTT is unaffected by the -/// video backlog. +/// base — queue stays bounded and control RTT is unaffected by the video +/// backlog. Header-arrival tail bounds below include retransmission; this is +/// not the clean-link 150 ms gate or an input-to-visible measurement. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn impaired_link_latency_gate() { // 60 fps of ~1 KB frames ≈ one datagram per frame — enough samples @@ -460,8 +460,8 @@ async fn impaired_link_latency_gate() { "too few frames arrived: {}", latencies.len() ); - // Protocol queues stay bounded: capture→send wait is the collapse - // + channel time only — must stay near zero even under loss. + // Protocol queues stay bounded: capture→send wait is the bounded + // channel and pacing time — must stay near zero even under loss. let queue_p95 = p95(queue_ms); let lat_p50 = p50(latencies.clone()); assert!( @@ -524,11 +524,11 @@ async fn soak_60fps() { p99_age, rss_peak ); - // Starvation check: the collapse may legitimately shed frames under + // Starvation check: bounded admission may legitimately shed frames under // CPU contention, so the floor is sustained flow, not offered rate — // a stalled pipeline delivers ~0. assert!(ages.len() >= secs as usize * 10, "starved: {}", ages.len()); - // G5: viewer-visible latency ≤150 ms p95 in-process on a clean link. + // G5 synthetic header-arrival age ≤150 ms p95 on a clean loopback link. assert!( p95_age <= 150, "frame age p95 {p95_age}ms exceeds G5 budget" From a0d1119d167a3d2281f88cfda44f3b49f5af11bf Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sat, 26 Sep 2026 15:19:29 +0500 Subject: [PATCH 2/5] docs(desktop): define sender reference recovery contract --- docs/desktop-frame-delivery.md | 30 +++++++++++++++++++++++------- docs/remediation-progress.md | 2 +- 2 files changed, 24 insertions(+), 8 deletions(-) diff --git a/docs/desktop-frame-delivery.md b/docs/desktop-frame-delivery.md index 946b169..3d02f72 100644 --- a/docs/desktop-frame-delivery.md +++ b/docs/desktop-frame-delivery.md @@ -1,11 +1,22 @@ # Desktop frame sender contract -Scope: partial writes and serving-session cancellation in W6.1/W2.5/W2.6. -This does not close codec reference-chain or presentation gates. +Scope: partial writes, reference recovery and serving-session cancellation in +W6.1/W2.5/W2.6. This does not close real-network codec or presentation gates. Each encoded frame has one tagged unidirectional stream. An in-flight keyframe -finishes before a newer frame is selected; a stale delta is abandoned with -RESET. When the producer closes during a send, its final frame finishes. +finishes before a newer frame is selected. A delta also finishes when the next +candidate is another delta, because the successor may depend on it. An +independent keyframe may replace an in-flight delta with RESET. When the producer +closes during a send, its final frame finishes. + +The producer queue remains two slots. Losing any encoded frame to admission +backpressure requests an IDR. Queue collapse selects a recent keyframe when +available and requests recovery for discarded references after that keyframe; +each drain is bounded to two items. The sender admits deltas only after a +keyframe and only at the exact next sequence. Gaps, empty encode results and +sequence exhaustion invalidate the chain. Dependent deltas wait for an IDR, +without first being sent as undecodable work or requiring a client roundtrip. +Selection is checked again after pacing, when additional frames may be dropped. The payload writer retains one pinned `write_all` future across producer events. `write_all` may already have accepted a prefix when its other select @@ -31,7 +42,8 @@ write deadline. No detached worker retains a connection after async teardown. Regression coverage uses a 64-byte duplex buffer and a 4096-byte patterned payload to force a partial write before each producer event. Both supersession and closure previously delivered 4160 bytes; both now preserve exact bytes. -Other tests retain errors after the producer event, abandon stale deltas without +Other tests retain errors after the producer event, preserve a partial delta +before its successor, replace a delta with an independent keyframe without waiting for a reader, and observe RESET after aborting an owned sender on real iroh/noq connections while a subsequent stream still succeeds. @@ -39,8 +51,12 @@ A separate real-session regression observes the producer being dropped after server cancellation with the client connection still alive; this failed before task-group ownership was added and passes on both backends. -Remaining work includes reference-aware queue collapsing with real codec continuity tests, per-session -wire IDs, global encoded/decoded/native allocation accounting, managed viewers, +Native H.264 coverage checks decode and a pixel across a dropped reference and +IDR recovery. This is separate from the synthetic QUIC impairment measurements; +it does not establish real-codec continuity under network loss/reordering. + +Remaining work includes real-codec network continuity and overload/IDR-rate +qualification, per-session wire IDs, global encoded/decoded/native allocation accounting, managed viewers, rendering/input release and platform/network acceptance. See the [client receive contract](desktop-client-lifecycle.md). No wire format, frame priority, dependency, installed binary or published release changes here. diff --git a/docs/remediation-progress.md b/docs/remediation-progress.md index aa7e787..398c299 100644 --- a/docs/remediation-progress.md +++ b/docs/remediation-progress.md @@ -51,7 +51,7 @@ Neither increment closes these product gaps or any wave. | W3.6 | Partial; validated selection and local child failure isolation | Extra paths become eligible on Established; weak policy ownership includes bounded backoff for temporary connection-ID/path-credit exhaustion and candidate-address snapshots. Relay-link route loss and known tunnel closure retire stale relay selection. Mux child send/receive failures are isolated; policy withdraws failed advertisements, excludes failed routes even when last-path close is refused, and closes held connections after all-child loss. Real loopback fixtures preserve open-stream traffic and accept new connections on the surviving child. Policy-observed telemetry exposes sticky event loss and unknown selection. Full path-event resynchronization, lossless retirement metrics, interface/socket recreation, per-service scheduling and physical-network qualification remain open. | | W4.4 | Partial; directional service boundaries | Real iroh/noq agents refuse writes with `SyncRead` and reads with `SyncWrite`, permit authorized transfers, and remain usable after refusal/cancellation. A view-only desktop never calls its input sink; failed injection never emits a success ACK. Linux/macOS account isolation, per-path policy, native seat/focus boundaries, consent and concurrent-role qualification remain open. See [receipt](reports/rds-service-scopes-20260926.md). | | W5.1/W5.3/W5.5 | Partial; native SSH client and standard PTY | `rds-ssh` uses russh 0.63.3 over pinned managed/direct streams. Explicit host pins, key/agent authentication, PTY/exec acknowledgements, terminal restoration, resize, cancellation and complete exit/output handling have Linux regression coverage and an OpenSSH interop fixture. SSH-specific fixed telemetry names are accepted by Vector; JSON uses a separate private file and terminal console logging pauses during SSH. GDS host/account provisioning, certificates/MFA, native macOS, broker/reattachment and mixed-load/network qualification remain open. See [contract](ssh.md). | -| W6.1 | Partial; sender partial-write correctness | A pinned payload write preserves progress across supersession/producer closure; write/FIN errors fail, and a frame guard resets on cancellation/error under one 30-second deadline. The serving future owns capture/writer/pacing tasks and resets its control stream on cancellation; normal exit joins async siblings. Failing-before byte/capture regressions and real iroh/noq RESET coverage pass. Reference-aware queue dropping, wire session IDs and native presentation/network gates remain open; see [contract](desktop-frame-delivery.md). | +| W6.1 | Partial; sender byte correctness and reference recovery | Pinned writes preserve progress; deltas finish before dependent successors and may be replaced by an independent keyframe. Queue loss requests IDR, and a sender sequence guard refuses broken delta chains after selection/pacing. Failure/cancellation resets streams under one deadline; serving owns child tasks. Failing-before byte/lifecycle regressions, native codec recovery and real iroh/noq RESET coverage pass. Real-codec network/overload qualification, wire session IDs and native presentation gates remain open; see [contract](desktop-frame-delivery.md). | | W6.2 | Partial; client receive ownership and limits | Four encoded frames per session/eight per process, two blocking decoder calls, owned task groups and cancellation, full handshake deadlines, dimension/sequence checks and rate-limited keyframe recovery are implemented. Real iroh/noq lifecycle and native codec regressions passed locally. Per-session wire IDs, global decoded/native memory accounting, renderer and native platform/network acceptance remain open; see [contract](desktop-client-lifecycle.md). | | W6.4 | Partial; ACK semantics and native X11 mapping | ACK follows successful backend acceptance. A bounded session-owned input worker reuses its sink. X11 maps evdev keys/buttons, bounds coordinates and fractional scroll, uses relative motion correctly, selects the explicit screen, checks native errors and releases owned holds on drop. Four failing-before native regressions and two-screen Xvfb checks cover the increment; the Linux CI lane requires actual native execution. Exclusive seat/controller ownership, focus races, custom layouts/IME, dynamic geometry and input-to-visible qualification remain open; see [contract](x11-input.md). | | W10.1/W10.2 | Partial; O1 foundation and O2 admin/source increment | Shared bounded Rust telemetry and Vector/OpenObserve pipeline; opt-in authenticated loopback metrics on agent/relay/server, aggregate source observations and old public metrics removal. Durable policy/catalog observations and effective agent revocation revision/lease are implemented. Finer queue/task and upstream adapter coverage, phase/reason correlation, support bundles, private rollout, independent liveness and overhead/platform qualification remain O2–O6; see [contract](observability.md). | From f8c703a3cfbb7687f2d2532b72682820848ebcb3 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sat, 26 Sep 2026 15:19:55 +0500 Subject: [PATCH 3/5] test(desktop): use fixed pixel chunks in codec recovery fixture --- crates/rds-desktop/src/session.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 956a6bf..3bf5f3f 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -1076,7 +1076,7 @@ mod tests { let mut decoded = Vec::new(); for seq in 0..7 { let mut bgra = vec![0; 64 * 64 * 4]; - for (i, pixel) in bgra.chunks_exact_mut(4).enumerate() { + for (i, pixel) in bgra.as_chunks_mut::<4>().0.iter_mut().enumerate() { let value = if (i / 64 + seq as usize * 4) % 32 < 16 { 220 } else { From b3d910d79cb93d6cb4954e7a171af2afef5a98d2 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sat, 26 Sep 2026 15:28:13 +0500 Subject: [PATCH 4/5] docs(desktop): record reference recovery qualification --- .../rds-desktop-references-20260926-data.json | 41 ++++++++++++++ .../rds-desktop-references-20260926.md | 55 +++++++++++++++++++ 2 files changed, 96 insertions(+) create mode 100644 docs/reports/rds-desktop-references-20260926-data.json create mode 100644 docs/reports/rds-desktop-references-20260926.md diff --git a/docs/reports/rds-desktop-references-20260926-data.json b/docs/reports/rds-desktop-references-20260926-data.json new file mode 100644 index 0000000..5d7e859 --- /dev/null +++ b/docs/reports/rds-desktop-references-20260926-data.json @@ -0,0 +1,41 @@ +{ + "schema_version": 1, + "commit": "4221c59a53713de716f4e781a858a3587f991621", + "dirty": false, + "os": "linux", + "architecture": "x86_64", + "profile": "test (debug)", + "features": "rds-desktop default; dev harness enables transport-noq", + "command": "RDS_SOAK_SECS=60 CARGO_BUILD_JOBS=2 cargo test --locked -p rds-desktop --test session_v2 -- --nocapture --test-threads=1", + "toolchain": "rustc 1.98.1 (48a229cea 2026-09-01)", + "lockfile_sha256": "78723e059a78f49d715542632620884b1a279b0db82acd379259b5344f20acdb", + "test_binary_sha256": "79975cd92bf9b3163db368f2c5a09c9adbcf39c7921bfcacdece0419048ab811", + "repetitions": 1, + "failures": 0, + "skips": 0, + "tests_passed": 7, + "measurement": "synthetic complete frame-header arrival, not decoded/presented pixels", + "clean_iroh_loopback": { + "duration_seconds": 60, + "frames": 3601, + "age_p95_ms": 3, + "age_p99_ms": 5, + "process_rss_peak_kib": 50060 + }, + "impaired_noq_loopback": { + "frames": 420, + "age_p95_ms": 110, + "age_p99_ms": 116, + "queue_p95_ms": 37, + "wire_p95_ms": 80, + "control_rtt_p95_ms": 159, + "server_forwarded": 1276, + "server_dropped": 67, + "server_bytes": 355132, + "client_forwarded": 1165, + "client_dropped": 64, + "client_bytes": 57085 + }, + "impairment": "Impairment::lossy() at recorded commit, applied under both endpoints UDP sockets", + "acceptance_limits": "existing session_v2 thresholds unchanged; no native desktop or WAN acceptance" +} diff --git a/docs/reports/rds-desktop-references-20260926.md b/docs/reports/rds-desktop-references-20260926.md new file mode 100644 index 0000000..9c6f6c5 --- /dev/null +++ b/docs/reports/rds-desktop-references-20260926.md @@ -0,0 +1,55 @@ +# Desktop reference-chain qualification — 2026-09-26 + +Scope: W6.1 increment; no wave, input-to-visible or release gate is closed. +Qualified source: `4221c59a53713de716f4e781a858a3587f991621` (clean tree). +See the [sender contract](../desktop-frame-delivery.md). + +A blocked delta payload was abandoned when its dependent successor arrived. +The new 64-byte-backpressure regression failed before the change and now checks +that the full 4096-byte reference reaches the reader exactly once. A successor +delta waits for that reference; an independent keyframe can replace it with +RESET. Final frames still finish when the producer closes. + +Queue admission now requests IDR after losing any encoded frame. Queue collapse +is bounded and requests recovery for discarded references, including deltas +after a retained keyframe. A sender sequence guard rejects deltas after a gap, +empty encode result or sequence exhaustion until an independent frame arrives. +Admission occurs after the final post-pacing selection. Queue capacity remains +two; no wire or dependency change is introduced. + +## Validation and measurements + +Linux x86_64, Rust 1.98.1, two build jobs, one local Cargo operation at a time: + +- Formatting and workspace clippy, default and X11, passed with warnings denied. +- Workspace: 555 passed, zero failed, three ignored across 96 test/doc targets. +- X11 unit suite: 28 passed, zero failed; the native display-repainting fixture + remains explicitly ignored here and required in the separate Linux Xvfb lane. +- Native H.264 encodes seven images with actual IDR/delta flags. After dropping + one reference, the guard admits sequences 0, 1, 5 and 6, rejects dependent + deltas 3/4, and checks a decoded pixel before and after the requested IDR. +- Existing real iroh/noq stream RESET and serving/client cancellation tests pass. +- A clean-source serial run passed all seven session scenarios. A 60-second + iroh loopback delivered 3601 synthetic headers, age p95 3 ms/p99 5 ms, process + peak RSS 50060 KiB. The impaired noq lane delivered 420 headers, age p95 + 110 ms/p99 116 ms, queue p95 37 ms and control RTT p95 159 ms. Both directions + recorded actual seeded datagram drops. + +[Reproduction data](rds-desktop-references-20260926-data.json) contains exact +source, command, toolchain, binary/lockfile hashes and impairment counters. +These are single-run synthetic header measurements, not decoded/presented +latency, a comparative improvement, a native-memory ceiling or WAN acceptance. +The native codec fixture is a separate local test, not real-codec QUIC impairment. + +An intermediate FIFO-only prototype preserved references but failed the existing +queue-age threshold: p95 114 ms against 100 ms. It was replaced by bounded +collapse with explicit recovery/sequence admission. The final policy passed a +focused impairment test and the complete suites above. No threshold was relaxed +and no failed test was retried unchanged to manufacture a pass. X11 clippy also +caught a test-fixture chunk API lint, which was corrected before final checks. + +Remaining: real H.264 network loss/reordering/RESET continuity, long-duration +overload and IDR-rate behavior, pacing/grant convergence, desktop wire session +IDs, managed viewer/rendering, native memory accounting and physical-platform +input-to-visible acceptance. GitHub integration is independently required before +merge. No installed agent or published preview is changed. From 314ab2d1a6bc92cda255fe6eb87b62534aae2257 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sat, 26 Sep 2026 15:51:03 +0500 Subject: [PATCH 5/5] test(relay): bound and diagnose owned relay qualification phases --- crates/rds-relay/tests/owned_e2e.rs | 100 +++++++++++------- .../rds-desktop-references-20260926.md | 15 +++ 2 files changed, 78 insertions(+), 37 deletions(-) diff --git a/crates/rds-relay/tests/owned_e2e.rs b/crates/rds-relay/tests/owned_e2e.rs index 9f1cb43..c06f934 100644 --- a/crates/rds-relay/tests/owned_e2e.rs +++ b/crates/rds-relay/tests/owned_e2e.rs @@ -150,36 +150,51 @@ async fn attached_relay_after_silent_direct(dual_stack: bool) { #[tokio::test] async fn relay_forwards_handshake_and_datagrams() { + async fn phase(name: &str, work: impl Future) -> T { + eprintln!("owned relay fixture: starting {name}"); + let result = tokio::time::timeout(std::time::Duration::from_secs(5), work) + .await + .unwrap_or_else(|_| panic!("owned relay fixture timed out during {name}")); + eprintln!("owned relay fixture: completed {name}"); + result + } + let _ = tracing_subscriber::fmt() .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) .try_init(); - let relay = rds_relay::server::serve( - EndpointConfig { - backend: Backend::Noq, - secret_key: Some(key(0)), - bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], - ..Default::default() - }, - Vec::new(), + let relay = phase( + "relay startup", + rds_relay::server::serve( + EndpointConfig { + backend: Backend::Noq, + secret_key: Some(key(0)), + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + ..Default::default() + }, + Vec::new(), + ), ) .await .expect("relay up"); - let a = endpoint(1, relay.endpoint_addr()).await; - let b = endpoint(2, relay.endpoint_addr()).await; + let a = phase("attach A", endpoint(1, relay.endpoint_addr())).await; + let b = phase("attach B", endpoint(2, relay.endpoint_addr())).await; assert_eq!(relay.endpoints(), 2, "both endpoints attached"); // Exercise the directory boundary too: owned relay locators have their own // scheme and public identity, and must survive signed announce/resolve. - let directory = rds_discovery::service::serve( - "127.0.0.1:0".parse().unwrap(), - std::sync::Arc::new(rds_discovery::MemoryStore::default()), - rds_discovery::service::ServiceConfig::open_ephemeral(), + let directory = phase( + "directory startup", + rds_discovery::service::serve( + "127.0.0.1:0".parse().unwrap(), + std::sync::Arc::new(rds_discovery::MemoryStore::default()), + rds_discovery::service::ServiceConfig::open_ephemeral(), + ), ) .await .unwrap(); let client = rds_discovery::client::Client::new(directory.addr()); - let _announce = rds_net::announce( + let announce = rds_net::announce( b.clone(), rds_net::AnnounceConfig { issuer: rds_discovery::RecordIssuer::memory(ed25519_dalek::SigningKey::from_bytes( @@ -206,45 +221,56 @@ async fn relay_forwards_handshake_and_datagrams() { let target = relay_only(resolved); // Accept must be polled while connect is in flight: the server's // endpoint only drives the handshake once the incoming is taken. - let accept_b = tokio::spawn({ - let b = b.clone(); - async move { b.accept().await.expect("incoming").await } - }); - let conn_a = a.connect(target, ALPN).await.expect("connect over relay"); + let (conn_a, conn_b) = phase("peer handshake", async { + tokio::join!(a.connect(target, ALPN), async { + b.accept().await.expect("incoming").await + }) + }) + .await; + let conn_a = conn_a.expect("connect over relay"); assert_eq!(conn_a.remote_id(), b.id()); - let conn_b = accept_b.await.unwrap().expect("accept"); + let conn_b = conn_b.expect("accept"); assert_eq!(conn_b.remote_id(), a.id()); // Datagrams both ways — payload is opaque outer-QUIC to the relay. conn_a.send_datagram(b"hello".to_vec().into()).unwrap(); - let got = conn_b.read_datagram().await.unwrap(); + let got = phase("datagram A to B", conn_b.read_datagram()) + .await + .unwrap(); assert_eq!(&got[..], b"hello"); conn_b.send_datagram(b"world".to_vec().into()).unwrap(); - let got = conn_a.read_datagram().await.unwrap(); + let got = phase("datagram B to A", conn_a.read_datagram()) + .await + .unwrap(); assert_eq!(&got[..], b"world"); // Streams ride the same tunnel. - let (mut send, mut recv) = conn_a.open_bi().await.unwrap(); - send.write_all(b"stream-data").await.unwrap(); - send.finish().unwrap(); - let (mut bs, mut br) = conn_b.accept_bi().await.unwrap(); - let mut buf = vec![0u8; 11]; - br.read_exact(&mut buf).await.unwrap(); - assert_eq!(&buf, b"stream-data"); - bs.write_all(b"ack").await.unwrap(); - bs.finish().unwrap(); - let ack = recv.read_to_end(usize::MAX).await.unwrap(); - assert_eq!(&ack, b"ack"); + phase("reliable stream roundtrip", async { + let (mut send, mut recv) = conn_a.open_bi().await.unwrap(); + send.write_all(b"stream-data").await.unwrap(); + send.finish().unwrap(); + let (mut bs, mut br) = conn_b.accept_bi().await.unwrap(); + let mut buf = vec![0u8; 11]; + br.read_exact(&mut buf).await.unwrap(); + assert_eq!(&buf, b"stream-data"); + bs.write_all(b"ack").await.unwrap(); + bs.finish().unwrap(); + let ack = recv.read_to_end(3).await.unwrap(); + assert_eq!(&ack, b"ack"); + }) + .await; let (forwarded, dropped, bytes) = relay.stats(); assert!(forwarded > 0, "relay must have forwarded datagrams"); assert!(bytes > 0); assert_eq!(dropped, 0); - a.close().await; - b.close().await; - relay.close().await.unwrap(); + drop(announce); + phase("close A", a.close()).await; + phase("close B", b.close()).await; + phase("close relay", relay.close()).await.unwrap(); + phase("close directory", directory.close()).await.unwrap(); } #[tokio::test] diff --git a/docs/reports/rds-desktop-references-20260926.md b/docs/reports/rds-desktop-references-20260926.md index 9c6f6c5..aaa0b5b 100644 --- a/docs/reports/rds-desktop-references-20260926.md +++ b/docs/reports/rds-desktop-references-20260926.md @@ -48,6 +48,21 @@ focused impairment test and the complete suites above. No threshold was relaxed and no failed test was retried unchanged to manufacture a pass. X11 clippy also caught a test-fixture chunk API lint, which was corrected before final checks. +The first GitHub Linux run at `b3d910d79cb93d6cb4954e7a171af2afef5a98d2` +([run 36235900357](https://github.com/NDDev-OpenNetwork/remote-device-sync/actions/runs/36235900357)) +stalled in the pre-existing `relay_forwards_handshake_and_datagrams` fixture. +The live job log showed its five sibling scenarios completed and this scenario +still running after more than 17 minutes of job time. Its unbounded awaits did +not identify the stalled phase; this is not evidence of a specific relay-runtime +root cause. The fixture now reports startup/attach/handshake/datagram/stream/close +phases with five-second deadlines, polls connect/accept together without a +detached task, bounds its ACK read and explicitly closes the directory. All +original payload/identity/forwarding assertions remain. One full local run and +a fixed ten-run batch passed all six scenarios (66 executions). This bounds and +improves diagnosis of the qualification test; it does not prove loss-free +DATAGRAM delivery or repair an unlocalized runtime fault. Final GitHub checks +are required for the revised source, not inferred from earlier green jobs. + Remaining: real H.264 network loss/reordering/RESET continuity, long-duration overload and IDR-rate behavior, pacing/grant convergence, desktop wire session IDs, managed viewer/rendering, native memory accounting and physical-platform