diff --git a/crates/rds-core/src/lib.rs b/crates/rds-core/src/lib.rs index 8a31ac3..ff625b1 100644 --- a/crates/rds-core/src/lib.rs +++ b/crates/rds-core/src/lib.rs @@ -37,6 +37,11 @@ pub const PROTOCOL_VERSION: u16 = 3; /// Upper bound for a serialized greeting, guard against abusive peers. pub const MAX_MESSAGE_LEN: u32 = 64 * 1024; +/// Desktop frame STOP_SENDING code: this predecessor is no longer needed +/// after a newer independent picture. It is neither fresh successful delivery +/// nor a broken current reference. Older senders may treat it as generic failure. +pub const DESKTOP_FRAME_OBSOLETE: u32 = 0x5244_5301; + /// First frame on every uni-directional stream (v3): routes the stream /// to the service that owns it. The accepting side runs one /// `accept_uni` demux per connection and hands each stream to the diff --git a/crates/rds-desktop/src/client.rs b/crates/rds-desktop/src/client.rs index fa4f99d..a7aee03 100644 --- a/crates/rds-desktop/src/client.rs +++ b/crates/rds-desktop/src/client.rs @@ -719,6 +719,7 @@ impl GapRepair { enum FrameRead { Complete(FrameHeader, Vec, FrameBudget), + Stale, Rejected, TimedOut, } @@ -736,6 +737,7 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) { let mut admitted = 0u64; let mut completed = 0u64; let mut rejected = 0u64; + let mut stale = 0u64; let mut timed_out = 0u64; let mut gaps = 0u64; let mut received_keys = 0u64; @@ -815,13 +817,14 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) { tracing::warn!(expected_seq=ctx.next_seq.load(Ordering::Relaxed),reorder_budget_ms=wait.as_millis(),request_sent,pending_recovery_age_ms=gap_repair.requested_at_ms.map(|at| now_ms.saturating_sub(at)),"desktop reference gap expired"); } if health.elapsed()>=std::time::Duration::from_secs(5) { - tracing::info!(admitted,completed,rejected,timed_out,gaps,received_keys,in_flight=ctx.receiving.load(Ordering::Relaxed),"desktop receiver health"); + tracing::info!(admitted,completed,rejected,stale,timed_out,gaps,received_keys,in_flight=ctx.receiving.load(Ordering::Relaxed),"desktop receiver health"); health=std::time::Instant::now(); } }, frame = readers.join_next(), if !readers.is_empty() => { match frame { Some(Ok(FrameRead::Complete(header,body,budget)))=>ordered.push(header,(body,budget)), + Some(Ok(FrameRead::Stale))=>{stale+=1;}, Some(Ok(FrameRead::TimedOut))=>{timed_out+=1;gap_repair.clear();delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());}, _=>{rejected+=1;gap_repair.clear();delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());}, } @@ -870,18 +873,27 @@ async fn read_one( let frame_seq = header.seq; let keyframe = header.keyframe; let header_ms = started.elapsed().as_millis(); + if header.seq == u64::MAX + || crate::frame_bytes(header.width as usize, header.height as usize).is_none() + { + return FrameRead::Rejected; + } + let expected_seq = next_seq.load(Ordering::Relaxed); + if header.seq < expected_seq { + let _ = stream.stop(rds_core::DESKTOP_FRAME_OBSOLETE.into()); + tracing::debug!( + frame_seq, + expected_seq, + "desktop obsolete predecessor released" + ); + return FrameRead::Stale; + } let deadline = if keyframe { KEYFRAME_READ_TIMEOUT } else { FRAME_READ_TIMEOUT }; let reading = async { - if header.seq == u64::MAX - || header.seq < next_seq.load(Ordering::Relaxed) - || crate::frame_bytes(header.width as usize, header.height as usize).is_none() - { - return None; - } // Retain metadata-only progress when the future times out. read_to_end // discards its partial buffer on cancellation and hid whether even a // header or any media bytes arrived during observed freezes. diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 8fa38b4..e98af6d 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -60,6 +60,7 @@ struct DeliveryFeedback { late_pending: AtomicU64, failed: AtomicU64, acknowledged: AtomicU64, + obsolete: AtomicU64, last_ack_ms: AtomicU64, producing: AtomicBool, produced: AtomicU64, @@ -119,6 +120,21 @@ struct FrameDelivery { keyframe_pending: Arc, idr: Arc, feedback: Arc, + latest_key_seq: Arc, +} + +#[derive(Debug, PartialEq, Eq)] +enum FrameReceipt { + Delivered, + Obsolete, + Failed, +} + +fn request_frame_repair(latest_key_seq: &AtomicU64, idr: &AtomicBool, failed_seq: u64) { + let key = latest_key_seq.load(Ordering::Acquire); + if key == u64::MAX || key <= failed_seq { + idr.store(true, Ordering::Release); + } } fn delivery_delay_budget(path: Option) -> Duration { @@ -541,6 +557,7 @@ pub async fn serve_desktop_with( // Capture+encode runs on a blocking thread; frames flow to the writer. let (tx, mut rx) = mpsc::channel::(2); let keyframe_pending = Arc::new(AtomicBool::new(false)); + let latest_key_seq = Arc::new(AtomicU64::new(u64::MAX)); let capture_admission = Arc::new(AtomicU64::new(0)); let delivery_feedback = Arc::new(DeliveryFeedback::default()); { @@ -553,6 +570,7 @@ pub async fn serve_desktop_with( let input_refresh_pending = Arc::clone(&controls.input_refresh_pending); let mut producer = config.producer; let keyframe_pending = keyframe_pending.clone(); + let latest_key_seq = latest_key_seq.clone(); let capture_admission = capture_admission.clone(); let feedback = delivery_feedback.clone(); capture.spawn_blocking(move || { @@ -625,6 +643,7 @@ pub async fn serve_desktop_with( feedback.produced.fetch_add(1, Ordering::Relaxed); } if p.header.keyframe && !p.payload.is_empty() { + latest_key_seq.store(p.header.seq, Ordering::Release); keyframe_pending.store(true, Ordering::Release); } // Losing any encoded reference breaks its successors, @@ -750,6 +769,7 @@ pub async fn serve_desktop_with( late_pending, delivery_stalled_ticks = pressure.stalled_ticks, failed_delivery = failed, + obsolete_delivery = feedback.obsolete.load(Ordering::Relaxed), last_ack_ms = feedback.last_ack_ms.load(Ordering::Relaxed), path_rtt_ms = ?path.map(|p| p.rtt.as_millis()), path_via_relay = ?path.map(|p| p.via_relay), @@ -788,25 +808,38 @@ pub async fn serve_desktop_with( let mut superseded = 0u64; let mut acknowledgements = JoinSet::new(); let mut acknowledged = 0u64; + let mut obsolete = 0u64; let mut failed_delivery = 0u64; let mut health = Instant::now(); 'writer: loop { while let Some(result) = acknowledgements.try_join_next() { - if matches!(result, Ok(true)) { - acknowledged += 1; - } else { - failed_delivery += 1; - chain.next = None; - writer_idr.store(true, Ordering::Relaxed); + match result { + Ok((_, FrameReceipt::Delivered)) => acknowledged += 1, + Ok((_, FrameReceipt::Obsolete)) => obsolete += 1, + Ok((seq, FrameReceipt::Failed)) => { + failed_delivery += 1; + chain.failed_receipt(seq); + } + _ => { + failed_delivery += 1; + chain.next = None; + writer_idr.store(true, Ordering::Relaxed); + } } } if acknowledgements.len() >= MAX_PENDING_FRAME_ACKS { - if matches!(acknowledgements.join_next().await, Some(Ok(true))) { - acknowledged += 1; - } else { - failed_delivery += 1; - chain.next = None; - writer_idr.store(true, Ordering::Relaxed); + match acknowledgements.join_next().await { + Some(Ok((_, FrameReceipt::Delivered))) => acknowledged += 1, + Some(Ok((_, FrameReceipt::Obsolete))) => obsolete += 1, + Some(Ok((seq, FrameReceipt::Failed))) => { + failed_delivery += 1; + chain.failed_receipt(seq); + } + _ => { + failed_delivery += 1; + chain.next = None; + writer_idr.store(true, Ordering::Relaxed); + } } continue; } @@ -819,7 +852,7 @@ pub async fn serve_desktop_with( }; if produced.payload.is_empty() { chain.next = None; - writer_idr.store(true, Ordering::Relaxed); + request_frame_repair(&latest_key_seq, &writer_idr, produced.header.seq); continue; } let bps = writer_bitrate.load(Ordering::Relaxed).max(50_000) as f64 / 8.0; @@ -837,7 +870,7 @@ pub async fn serve_desktop_with( // 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); + request_frame_repair(&latest_key_seq, &writer_idr, produced.header.seq); continue; } produced.header.send_ts_ms = writer_clock.now_ms(); @@ -866,6 +899,7 @@ pub async fn serve_desktop_with( keyframe_pending: keyframe_pending.clone(), idr: writer_idr.clone(), feedback: writer_feedback.clone(), + latest_key_seq: latest_key_seq.clone(), }, ) .await @@ -873,6 +907,9 @@ pub async fn serve_desktop_with( SendOutcome::Sent => { sent += 1; } + SendOutcome::Obsolete => { + obsolete += 1; + } SendOutcome::Superseded(newer) => { superseded += 1; pending = Some(newer); @@ -891,6 +928,7 @@ pub async fn serve_desktop_with( sent, superseded, acknowledged, + obsolete, failed_delivery, pending_acknowledgements = acknowledgements.len(), pending_media_frames = capture_admission.load(Ordering::Acquire), @@ -1399,9 +1437,17 @@ mod x11 { #[derive(Default)] struct FrameChain { next: Option, + last_key_seq: Option, } impl FrameChain { + fn failed_receipt(&mut self, seq: u64) { + // The receipt worker already owns the one repair request. An older + // result joining after independent recovery cannot break that chain. + if self.last_key_seq.is_none_or(|key| key <= seq) { + self.next = None; + } + } fn admit(&mut self, produced: &Produced) -> bool { let header = &produced.header; if produced.payload.is_empty() @@ -1412,6 +1458,9 @@ impl FrameChain { return false; } self.next = header.seq.checked_add(1); + if header.keyframe { + self.last_key_seq = Some(header.seq); + } true } } @@ -1441,6 +1490,8 @@ impl Drop for FrameSend { enum SendOutcome { /// Frame fully sent. Sent, + /// Receiver already recovered past this predecessor; continue the writer. + Obsolete, /// A fresher decodable frame supersedes — send it next. Superseded(AdmittedFrame), /// Producer closed mid-send; the final frame was finished. @@ -1449,6 +1500,29 @@ enum SendOutcome { Failed, } +fn obsolete_write(error: &std::io::Error) -> bool { + matches!(error.get_ref().and_then(|cause| cause.downcast_ref::()), + Some(rds_net::WriteError::Stopped(code)) if *code == rds_core::DESKTOP_FRAME_OBSOLETE.into()) +} + +fn obsolete_send( + sending: &mut FrameSend, + produced: &Produced, + delivery: &FrameDelivery, +) -> SendOutcome { + sending.finished = true; + delivery.feedback.obsolete.fetch_add(1, Ordering::Relaxed); + if produced.header.keyframe { + delivery.keyframe_pending.store(false, Ordering::Release); + } + tracing::info!( + frame_seq = produced.header.seq, + payload_bytes = produced.payload.len(), + "desktop obsolete frame stopped during write" + ); + SendOutcome::Obsolete +} + /// Send one frame on its own tagged uni stream, aborting mid-write if /// an independent keyframe lands. Otherwise finish the reference on which /// the next delta may depend, retaining partial-write progress. @@ -1457,7 +1531,7 @@ async fn send_frame( route: rds_core::UniHello, produced: AdmittedFrame, rx: &mut mpsc::Receiver, - acknowledgements: &mut JoinSet, + acknowledgements: &mut JoinSet<(u64, FrameReceipt)>, delivery: FrameDelivery, ) -> SendOutcome { match tokio::time::timeout( @@ -1479,7 +1553,7 @@ async fn send_frame_inner( route: rds_core::UniHello, produced: AdmittedFrame, rx: &mut mpsc::Receiver, - acknowledgements: &mut JoinSet, + acknowledgements: &mut JoinSet<(u64, FrameReceipt)>, delivery: FrameDelivery, ) -> SendOutcome { let AdmittedFrame { produced, permit } = produced; @@ -1503,10 +1577,16 @@ async fn send_frame_inner( // per-connection demux routes on it. Per-session routes keep a stale // stream out of any replacement session's inbox. if let Err(e) = write_frame(&mut *stream, &route).await { + if obsolete_write(&e) { + return obsolete_send(&mut sending, &produced, &delivery); + } tracing::debug!("frame tag write failed: {e}"); return SendOutcome::Failed; } if let Err(e) = write_frame(&mut *stream, &produced.header).await { + if obsolete_write(&e) { + return obsolete_send(&mut sending, &produced, &delivery); + } tracing::debug!("frame header write failed: {e}"); return SendOutcome::Failed; } @@ -1518,11 +1598,21 @@ async fn send_frame_inner( Ok(PayloadOutcome::Superseded(next)) => SendOutcome::Superseded(next), Ok(PayloadOutcome::ProducerEnded) => SendOutcome::Done, Err(e) => { + if obsolete_write(&e) { + return obsolete_send(&mut sending, &produced, &delivery); + } tracing::debug!("frame send failed: {e}"); return SendOutcome::Failed; } }; if let Err(e) = stream.finish() { + // finish() reports an erased ClosedStream. Only a ready, exact peer + // disposition can classify it as obsolete; never wait on other errors. + if matches!(tokio::time::timeout(Duration::ZERO, stream.stopped()).await, + Ok(Ok(Some(code))) if code == rds_core::DESKTOP_FRAME_OBSOLETE.into()) + { + return obsolete_send(&mut sending, &produced, &delivery); + } tracing::debug!("frame finish failed: {e}"); return SendOutcome::Failed; } @@ -1534,6 +1624,7 @@ async fn send_frame_inner( keyframe_pending, idr, feedback, + latest_key_seq, } = delivery; // Retain the reset-on-drop owner until delivery is acknowledged. The // bounded task group is owned by this writer; cancellation resets its @@ -1561,13 +1652,19 @@ async fn send_frame_inner( 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 + FrameReceipt::Delivered + } + Ok(Ok(Some(code))) if code == rds_core::DESKTOP_FRAME_OBSOLETE.into() => { + sending.finished = true; + feedback.obsolete.fetch_add(1, Ordering::Relaxed); + tracing::info!(frame_seq=seq,payload_bytes,ack_ms=started.elapsed().as_millis(),"desktop obsolete frame receipt"); + FrameReceipt::Obsolete } result => { feedback.failed.fetch_add(1, Ordering::Relaxed); - idr.store(true, Ordering::Relaxed); + request_frame_repair(&latest_key_seq, &idr, seq); tracing::warn!(frame_seq=seq,payload_bytes,ack_ms=started.elapsed().as_millis(),outcome=?result,"desktop frame delivery unconfirmed"); - false + FrameReceipt::Failed } }; drop(late); @@ -1578,7 +1675,7 @@ async fn send_frame_inner( if keyframe { keyframe_pending.store(false, Ordering::Release); } - acknowledged + (seq, acknowledged) }); outcome } @@ -2049,6 +2146,151 @@ mod tests { } } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn obsolete_stop_during_payload_does_not_end_writer_or_count_fresh_delivery() { + for (backend, stop_code) in [rds_net::Backend::Iroh, rds_net::Backend::Noq] + .into_iter() + .flat_map(|backend| { + [rds_core::DESKTOP_FRAME_OBSOLETE, 42].map(move |code| (backend, code)) + }) + { + tokio::time::timeout(Duration::from_secs(10), async { + let config = rds_net::EndpointConfig { + backend, + discovery: false, + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + ..Default::default() + }; + let client = rds_net::bind_endpoint(config.clone()).await.unwrap(); + let server = rds_net::bind_endpoint(config).await.unwrap(); + 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 admission = Arc::new(AtomicU64::new(1)); + let feedback = Arc::new(DeliveryFeedback::default()); + let idr = Arc::new(AtomicBool::new(false)); + let produce = AdmittedFrame { + produced: Produced { + header: FrameHeader { + seq: 1, + keyframe: false, + capture_ts_ms: 0, + encode_done_ts_ms: 0, + send_ts_ms: 0, + codec: rds_core::Codec::H264, + width: 32, + height: 32, + }, + payload: Bytes::from(vec![1; 32 * 1024 * 1024]), + }, + permit: CapturePermit(admission.clone()), + }; + let (_tx, mut rx) = mpsc::channel(2); + let mut receipts = JoinSet::new(); + let delivery = FrameDelivery { + keyframe_pending: Arc::new(AtomicBool::new(false)), + idr: idr.clone(), + feedback: feedback.clone(), + latest_key_seq: Arc::new(AtomicU64::new(u64::MAX)), + }; + let (outcome, ()) = tokio::join!( + send_frame( + &a, + rds_core::UniHello::Desktop, + produce, + &mut rx, + &mut receipts, + delivery + ), + async { + let mut stream = b.accept_uni().await.unwrap(); + let _: rds_core::UniHello = read_frame(&mut stream).await.unwrap(); + let h: FrameHeader = read_frame(&mut stream).await.unwrap(); + assert_eq!(h.seq, 1); + stream.stop(stop_code.into()).unwrap(); + } + ); + let obsolete = stop_code == rds_core::DESKTOP_FRAME_OBSOLETE; + assert_eq!( + matches!(outcome, SendOutcome::Failed), + !obsolete, + "only the exact obsolete disposition may continue the writer" + ); + while let Some(result) = receipts.join_next().await { + assert_eq!(result.unwrap(), (1, FrameReceipt::Obsolete)); + } + assert_eq!(admission.load(Ordering::Acquire), 0); + assert_eq!( + feedback.acknowledged.load(Ordering::Acquire), + 0, + "obsolete disposal cannot justify bitrate growth" + ); + assert_eq!(feedback.failed.load(Ordering::Acquire), 0); + assert_eq!( + feedback.obsolete.load(Ordering::Acquire), + u64::from(obsolete) + ); + assert!(!idr.load(Ordering::Acquire)); + let mut next = a.open_uni().await.unwrap(); + next.write_all(b"still usable").await.unwrap(); + next.finish().unwrap(); + let mut stream = b.accept_uni().await.unwrap(); + assert_eq!(stream.read_to_end(32).await.unwrap(), b"still usable"); + a.close(0u32.into(), b"done"); + tokio::join!(client.close(), server.close()); + }) + .await + .expect("obsolete stop did not release resources"); + } + } + + #[test] + fn older_failed_receipt_cannot_invalidate_an_admitted_recovery_or_repeat_repair() { + let mut chain = FrameChain::default(); + let mut frame = Produced { + header: FrameHeader { + seq: 20, + keyframe: true, + capture_ts_ms: 0, + encode_done_ts_ms: 0, + send_ts_ms: 0, + codec: rds_core::Codec::H264, + width: 32, + height: 32, + }, + payload: Bytes::from_static(b"independent picture"), + }; + assert!(chain.admit(&frame)); + // The receipt worker requested repair; the producer consumed that flag + // and the writer already admitted its replacement. The older failed + // receipt may join afterward while that newer key is travelling. + let idr = AtomicBool::new(false); + request_frame_repair(&AtomicU64::new(20), &idr, 10); + chain.failed_receipt(10); + assert!( + !idr.load(Ordering::Relaxed), + "one failed receipt requested a second expensive recovery" + ); + frame.header.seq = 21; + frame.header.keyframe = false; + assert!( + chain.admit(&frame), + "an old failure broke the independently recovered reference chain" + ); + request_frame_repair(&AtomicU64::new(20), &idr, 21); + assert!( + idr.load(Ordering::Relaxed), + "a current failure still requests repair" + ); + chain.failed_receipt(21); + frame.header.seq = 22; + assert!( + !chain.admit(&frame), + "a failed current reference must still invalidate its successors" + ); + } + fn path(sent: u64, lost: u64, rtt_ms: u64, congestion: u64) -> PathStats { PathStats { path_id: 0, diff --git a/crates/rds-desktop/tests/stale_frame_receipts.rs b/crates/rds-desktop/tests/stale_frame_receipts.rs new file mode 100644 index 0000000..c57b983 --- /dev/null +++ b/crates/rds-desktop/tests/stale_frame_receipts.rs @@ -0,0 +1,174 @@ +//! A late predecessor after a newer key is obsolete, not a new broken chain. +use rds_core::{ + Codec, DesktopCaps, DesktopControl, DesktopHello, FrameHeader, HelloAck, StreamHello, UniHello, +}; +use rds_desktop::client::{DesktopSession, SessionOpts}; +use rds_net::{Backend, Connection, EndpointConfig, SendStream, read_frame, write_frame}; +use std::time::Duration; +fn header(seq: u64, keyframe: bool) -> FrameHeader { + FrameHeader { + seq, + keyframe, + capture_ts_ms: 0, + encode_done_ts_ms: 0, + send_ts_ms: 0, + codec: Codec::H264, + width: 32, + height: 32, + } +} +async fn tagged(conn: &Connection, id: [u8; 16]) -> SendStream { + let mut stream = conn.open_uni().await.unwrap(); + write_frame(&mut stream, &UniHello::DesktopFrames { id }) + .await + .unwrap(); + stream +} +async fn complete(conn: &Connection, id: [u8; 16], seq: u64) { + let mut s = tagged(conn, id).await; + write_frame(&mut s, &header(seq, true)).await.unwrap(); + s.write_all(b"encoded relay fixture").await.unwrap(); + s.finish().unwrap(); +} +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn late_header_after_recovery_does_not_request_another_key_but_malformed_does() { + for backend in [Backend::Iroh, Backend::Noq] { + let config = EndpointConfig { + backend, + discovery: false, + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + ..Default::default() + }; + let client_ep = rds_net::bind_endpoint(config.clone()).await.unwrap(); + let server_ep = rds_net::bind_endpoint(config).await.unwrap(); + let (a, b) = tokio::join!(client_ep.connect(server_ep.addr(), rds_core::ALPN), async { + server_ep.accept().await.unwrap().await + }); + let (a, b) = (a.unwrap(), b.unwrap()); + let (client, server) = tokio::join!( + DesktopSession::connect_opts( + &a, + DesktopHello { + display: 0, + max_fps: 30, + codec: Codec::H264, + input_acks: false + }, + SessionOpts { + session: Some(rand::random()), + 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!("tagged 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 (_send, mut recv, id) = server; + let (tx, mut controls) = tokio::sync::mpsc::channel(16); + let reader = tokio::spawn(async move { + while let Ok(message) = read_frame::<_, DesktopControl>(&mut recv).await { + if tx.send(message).await.is_err() { + break; + } + } + }); + complete(&b, id, 0).await; + assert_eq!( + tokio::time::timeout( + Duration::from_secs(2), + client.encoded.as_mut().unwrap().recv() + ) + .await + .unwrap() + .unwrap() + .header + .seq, + 0 + ); + // The predecessor's route is admitted, but its header remains in transit. + let mut late = tagged(&b, id).await; + complete(&b, id, 2).await; + assert_eq!( + tokio::time::timeout( + Duration::from_secs(2), + client.encoded.as_mut().unwrap().recv() + ) + .await + .unwrap() + .unwrap() + .header + .seq, + 2 + ); + write_frame(&mut late, &header(1, false)).await.unwrap(); + assert_eq!( + tokio::time::timeout(Duration::from_secs(2), late.stopped()) + .await + .unwrap() + .unwrap(), + Some(rds_core::DESKTOP_FRAME_OBSOLETE.into()), + "stale header must release its stream with the exact obsolete disposition" + ); + assert!( + tokio::time::timeout(Duration::from_millis(600), controls.recv()) + .await + .is_err(), + "an obsolete predecessor invalidated the already recovered chain" + ); + let mut bad = tagged(&b, id).await; + let mut invalid = header(3, false); + invalid.width = 0; + write_frame(&mut bad, &invalid).await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_secs(2), bad.stopped()) + .await + .unwrap() + .unwrap() + .is_some() + ); + assert!( + matches!( + tokio::time::timeout(Duration::from_secs(2), controls.recv()) + .await + .unwrap(), + Some(DesktopControl::RequestIdr) + ), + "malformed current frame still requires repair" + ); + complete(&b, id, 4).await; + assert_eq!( + tokio::time::timeout( + Duration::from_secs(2), + client.encoded.as_mut().unwrap().recv() + ) + .await + .unwrap() + .unwrap() + .header + .seq, + 4 + ); + drop(client); + reader.abort(); + let _ = reader.await; + a.close(0u32.into(), b"done"); + b.close(0u32.into(), b"done"); + tokio::join!(client_ep.close(), server_ep.close()); + } +} diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 88ab4ec..b97f416 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -409,3 +409,27 @@ ends the activity when the loop returns or unwinds. Its option explicitly allows idle system sleep and does not keep the display awake or prevent lock. This is scoped application activity, not a global power/QoS setting. Installed measurements still decide whether it improves a particular latency episode. + +## Obsolete predecessors and repair + +A valid header older than the receiver's recovered sequence releases its stream +with `STOP_SENDING` code `0x52445301`. The receiver counts it separately from +malformed or timed-out readers and does not request another key. The sender +recognizes only this exact obsolete disposition, including a stop during an +unfinished write. It frees that frame's admission without ending the video +writer, invalidating the recovered chain or counting fresh delivery for bitrate +growth. Unknown stops and actual failures retain their existing failure path. + +No greeting, control message or frame layout changes. Older senders can still +classify the code as generic failure and perform their previous repair behavior; +matching receivers/senders obtain the complete improvement. Limits, malformed +header validation, read/write deadlines and successful-delivery semantics stay +unchanged. Sender and receiver health report obsolete dispositions separately. + +Delivery results also retain their frame sequence. A failed receipt older than +an already admitted independent picture cannot invalidate the recovered chain. +The receipt worker owns its repair request; joining the result does not request +repair a second time. A newer produced key suppresses another request for an +older frame, including reference-admission rejection while that key is queued. +Current failures and unknown worker termination still require repair. This +correlation does not treat a produced picture as successful delivery. diff --git a/docs/reports/rds-obsolete-frame-20261002.md b/docs/reports/rds-obsolete-frame-20261002.md new file mode 100644 index 0000000..f98442a --- /dev/null +++ b/docs/reports/rds-obsolete-frame-20261002.md @@ -0,0 +1,41 @@ +# Obsolete predecessors after recovery — 2026-10-02 + +Source review found two mechanisms capable of prolonging recovery. A predecessor +whose header arrives after a newer independent picture was classified as a +rejected reader and requested another key. Its stream stop was also treated as +failed delivery by the sender; a stop during payload writing ended the video +writer. Neither outcome describes a broken current reference when the receiver +has already recovered past that frame. + +A real Iroh/Noq regression admits the older route, receives a newer key, then +supplies the old header. Old behavior fails because it requests another key. +New behavior releases that reader with a specific obsolete disposition, then +still requires repair for a malformed current header and accepts its replacement. + +A second real-stream regression stops a bounded large payload while it is still +being written. Old behavior ends the writer. New behavior releases admission, +records disposal separately, preserves the connection and counts no fresh ACK +that could justify bitrate growth. An unknown stop code still selects failure; +only the exact obsolete code selects continuation. Both transport backends run +the scenarios, without a display or live user input. + +The implementation uses an application STOP_SENDING code, without changing +serialized wire layouts or authority. Earlier senders retain their generic-stop +behavior. Reader/writer budgets, validation and deadlines are unchanged. Separate +obsolete counters prevent diagnostics from mislabeling disposal as delivery. +This advances W6 reference recovery, not a completed media acceptance checkpoint. + +Installed native follow-up remains required. The predecessor's Full HD ten-minute +observation had decoded-age p95/max185/4079ms, and later ordinary-work logs retained +multi-second pauses despite no slow worker CPU record. Native dispatch-to-window +capture was also confounded by surface occlusion; it cannot establish a network +latency SLO. These negative results are retained and are not closed by regression +success. This defect explains a possible redundant repair path, not every freeze. + +A related writer race also repeated repair: a receipt worker scheduled an IDR, +then its joined failure reset a chain that had already admitted the replacement +and scheduled another IDR. An equivalent-old-semantics regression fails on the +repeated request; corrected sequence-aware handling preserves the next delta +after independent recovery while still invalidating a failed current reference. +The latest produced-key marker coalesces repair for older frames; transport +success counters remain the sole successful-delivery growth signal. diff --git a/docs/research.md b/docs/research.md index 0e0a1b0..7230823 100644 --- a/docs/research.md +++ b/docs/research.md @@ -607,3 +607,15 @@ to distinguish actual computation from elapsed native waiting on Linux/macOS. The fixture combines real CPU work and sleeping on one blocking worker and requires its CPU observation to exclude the wait. These observations do not replace installed latency/quality acceptance. + +## 2026-10-02 obsolete media dispositions + +[QUIC STOP_SENDING](https://www.rfc-editor.org/rfc/rfc9000.html#section-3.5) +provides an application error code for terminating one receive stream. RDS uses +a specific code for a predecessor already superseded by receiver recovery, +distinguishing disposal from successful delivery and a broken current reference. +The sender preserves that distinction both during writes and while awaiting the +transport receipt. Existing peers can retain generic-stop recovery; this does +not require a frame-layout or authorization change. The +[real-stream receipt](reports/rds-obsolete-frame-20261002.md) records failing old +behavior, regression scope and the still-required installed qualification.