From 6f33a9ae04c1ebe3ccb1e518b48fc8ae50bb5d87 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Fri, 2 Oct 2026 19:58:28 +0500 Subject: [PATCH] fix(desktop): isolate managed events from media backpressure --- crates/rds-agent/tests/local_manager.rs | 16 + crates/rds-cli/src/desktop.rs | 50 ++- crates/rds-client/src/local/desktop.rs | 398 +++++++++++++++++- crates/rds-client/src/local/desktop_events.rs | 273 ++++++++++++ crates/rds-client/src/local/mod.rs | 118 +++++- crates/rds-core/src/local.rs | 89 ++++ crates/rds-desktop/src/render/viewer.rs | 5 + docs/architecture.md | 12 +- docs/local-sessions.md | 27 ++ docs/native-viewer.md | 31 +- docs/reports/rds-event-isolation-20261002.md | 66 +++ docs/research.md | 20 + 12 files changed, 1056 insertions(+), 49 deletions(-) create mode 100644 crates/rds-client/src/local/desktop_events.rs create mode 100644 docs/reports/rds-event-isolation-20261002.md diff --git a/crates/rds-agent/tests/local_manager.rs b/crates/rds-agent/tests/local_manager.rs index 82732c3..004eafe 100644 --- a/crates/rds-agent/tests/local_manager.rs +++ b/crates/rds-agent/tests/local_manager.rs @@ -1139,6 +1139,22 @@ async fn managed_desktop_reports_remote_refusal_without_leaking() { matches!(result, Err(Error::Rejected(ErrorCode::Remote))), "expected clean remote refusal, got {result:?}" ); + let separated = tokio::time::timeout( + Duration::from_secs(10), + client.desktop_profile_separated( + Some(session), + rds_core::DesktopHello { + display: 0, + max_fps: 30, + codec: rds_core::Codec::H264, + input_acks: true, + }, + 1080, + ), + ) + .await + .expect("separated desktop refusal hung"); + assert!(matches!(separated, Err(Error::Rejected(ErrorCode::Remote)))); // The refused open must not park a stream permit: open_tcp still // has its full data budget — every slot but the lane reserved // for control traffic. diff --git a/crates/rds-cli/src/desktop.rs b/crates/rds-cli/src/desktop.rs index 4d7384c..3824a3a 100644 --- a/crates/rds-cli/src/desktop.rs +++ b/crates/rds-cli/src/desktop.rs @@ -346,13 +346,10 @@ mod native { started: Instant, ) -> anyhow::Result { view.stage("opening desktop"); - let mut channel = client - .desktop_profile( - Some(session), - hello(options), - Some(options.resolution.height()), - ) + let (mut channel, mut events) = client + .desktop_profile_separated(Some(session), hello(options), options.resolution.height()) .await?; + view.managed_events_separated(); extent(view, &channel.caps, options.display)?; view.status("Waiting for screen"); let control = channel.control_handle(); @@ -365,20 +362,8 @@ mod native { loop { match channel.recv().await? { None => break Ok(false), - Some(ManagedMessage::Event(rds_core::DesktopEvent::Heartbeat { - ts_ms, - .. - })) => view - .control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), - Some(ManagedMessage::Event(rds_core::DesktopEvent::InputAck { - seq, .. - })) => view.input_ack(seq), - Some(ManagedMessage::Event(rds_core::DesktopEvent::ClipboardReady { - bytes, - .. - })) => { - view.clipboard_ready(bytes); - tracing::info!(bytes, "remote clipboard ready"); + Some(ManagedMessage::Event(_)) => { + anyhow::bail!("unexpected event on separated video channel") } Some(ManagedMessage::Frame(frame)) => { view.stage("decoding"); @@ -410,10 +395,33 @@ mod native { } } }; + let event_observation = async { + loop { + match events.recv().await.transpose()? { + Some(rds_core::DesktopEvent::Heartbeat { ts_ms, .. }) => view + .control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), + Some(rds_core::DesktopEvent::InputAck { seq, .. }) => view.input_ack(seq), + Some(rds_core::DesktopEvent::ClipboardReady { bytes, .. }) => { + view.clipboard_ready(bytes); + tracing::info!(bytes, "remote clipboard ready"); + } + None => { + tracing::warn!("managed desktop event channel ended"); + return Ok(false); + } + } + } + }; + let incoming = async { + tokio::select! { + result = event_observation => result, + result = media => result, + } + }; let controls = control::pump(input, &control, &last_frame, started, |message| { view.input_sent(message); }); - let result = control::run(controls, media, stop).await; + let result = control::run(controls, incoming, stop).await; // A winning leg may cancel a partially written control on the other // leg. EOF closes the manager's desktop; never append Finished to a // potentially incomplete frame. Unrelated manager sessions survive. diff --git a/crates/rds-client/src/local/desktop.rs b/crates/rds-client/src/local/desktop.rs index 046885e..c0d9a3e 100644 --- a/crates/rds-client/src/local/desktop.rs +++ b/crates/rds-client/src/local/desktop.rs @@ -33,6 +33,85 @@ async fn write_payload(writer: &mut W, payload: &[u8]) -> writer.write_all(payload).await } +/// A slow video writer cannot hold remote input acknowledgements or heartbeat +/// echoes. All retained legs end with the desktop owner or event subscriber. +pub(super) async fn serve_separated( + stream: &mut UnixStream, + session: rds_desktop::client::DesktopSession, + registration: super::desktop_events::Registration, +) -> io::Result<()> { + let (reader, writer) = stream.split(); + serve_separated_io(reader, writer, session, registration).await +} + +async fn serve_separated_io( + mut reader: R, + mut writer: W, + mut session: rds_desktop::client::DesktopSession, + mut registration: super::desktop_events::Registration, +) -> io::Result<()> { + let ctrl = session.control_sender(); + let encoded = session.encoded.as_mut().ok_or_else(invalid)?; + let up = async { + loop { + match read_frame::<_, DesktopUp>(&mut reader).await { + Ok(DesktopUp::Control(control)) => ctrl + .send(control) + .await + .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "control closed"))?, + Ok(DesktopUp::Finished) => return Ok(()), + Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(()), + Err(e) => return Err(e), + } + } + }; + tokio::pin!(up); + // Keep this read alive across the attachment fence: canceling a pending + // framed read and continuing it later would lose its partial parse state. + let event_stop = registration.stop.clone(); + tokio::select! { + result = &mut up => return result, + result = registration.attached() => result?, + _ = event_stop.cancelled() => return Ok(()), + } + let media = async { + while let Some(frame) = encoded.recv().await { + if frame.payload.len() > MAX_DESKTOP_PAYLOAD { + return Err(invalid()); + } + write_frame( + &mut writer, + &DesktopDown::Frame { + header: frame.header, + }, + ) + .await?; + write_payload(&mut writer, &frame.payload).await?; + } + Ok(()) + }; + let events = async { + while let Some(event) = session.events.recv().await { + registration.sender.try_send(event).map_err(|_| { + // Metadata only; never discard an ACK and keep a healthy status. + tracing::warn!("desktop event subscriber unavailable or full"); + io::Error::new( + io::ErrorKind::BrokenPipe, + "desktop event subscriber unavailable or full", + ) + })?; + } + Ok(()) + }; + tokio::select! { + biased; + _ = registration.stop.cancelled() => Ok(()), + result = up => result, + result = events => result, + result = media => result, + } +} + /// Read one u32-length-prefixed encoded payload, bounded like the remote /// frame reader. async fn read_payload(reader: &mut R) -> io::Result> { @@ -127,6 +206,9 @@ pub enum ManagedMessage { Event(rds_core::DesktopEvent), } +/// Dedicated event receive queue; drop the owning desktop to abort its reader. +pub type ManagedEvents = tokio::sync::mpsc::Receiver>; + /// Shared control-plane state behind `ManagedDesktop`/`ManagedControl`: /// the socket's write half and the viewer-side sequence counters. /// Serializing writes through the mutex keeps postcard frames atomic. @@ -212,6 +294,7 @@ impl ManagedControl { pub struct ManagedDesktop { messages: tokio::sync::mpsc::Receiver>>, reading: tokio::task::JoinHandle<()>, + event_reading: Option>, state: Arc, /// The manager session this channel is pinned to. pub session: SessionId, @@ -222,6 +305,9 @@ pub struct ManagedDesktop { impl Drop for ManagedDesktop { fn drop(&mut self) { self.reading.abort(); + if let Some(reading) = &self.event_reading { + reading.abort(); + } } } @@ -274,6 +360,7 @@ impl ManagedDesktop { Self { messages, reading, + event_reading: None, state: Arc::new(ControlState { writer: Mutex::new(writer), display, @@ -286,6 +373,35 @@ impl ManagedDesktop { } } + pub(super) fn separated( + stream: UnixStream, + mut events: UnixStream, + session: SessionId, + caps: DesktopCaps, + display: u32, + ) -> (Self, ManagedEvents) { + let mut desktop = Self::new(stream, session, caps, display); + let (sender, receiver) = tokio::sync::mpsc::channel(super::desktop_events::EVENT_CAPACITY); + desktop.event_reading = Some(tokio::spawn(async move { + // Own the whole event socket: dropping its write half would look + // like caller departure to the same-UID event service. + loop { + let event = match read_frame::<_, DesktopDown>(&mut events).await { + Ok(DesktopDown::Event(event)) => Ok(event), + Ok(DesktopDown::Finished) => break, + Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break, + Ok(DesktopDown::Frame { .. }) => Err(invalid()), + Err(e) => Err(e), + }; + let failed = event.is_err(); + if sender.send(event).await.is_err() || failed { + break; + } + } + })); + (desktop, receiver) + } + /// A cloneable control half for tasks that send while `recv` runs — /// mirroring `DesktopSession::control_sender`. pub fn control_handle(&self) -> ManagedControl { @@ -539,6 +655,10 @@ mod tests { } async fn loopback_relay() -> Relay { + loopback_relay_with_payload(1024).await + } + + async fn loopback_relay_with_payload(payload_bytes: usize) -> Relay { use rds_core::{HelloAck, StreamHello, UniHello}; use rds_desktop::{SessionConfig, SyntheticProducer, serve_desktop_with}; @@ -578,7 +698,7 @@ mod tests { SessionConfig { input_sink: Some(Box::new(NoopInput)), producer: Some(Box::new( - SyntheticProducer::new(30, 640, 480, 1024).keyframe_every(5), + SyntheticProducer::new(30, 640, 480, payload_bytes).keyframe_every(5), )), frame_route: route, ..Default::default() @@ -723,4 +843,280 @@ mod tests { .unwrap(); relay.server_task.abort(); } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn separate_events_progress_while_video_body_is_not_read() { + use super::super::desktop_events; + use std::time::Duration; + let relay = loopback_relay_with_payload(1024 * 1024).await; + let registry = Arc::new(desktop_events::Registry::default()); + let registration = registry.register().unwrap(); + let route = registration.route; + let subscription = registry.take(route).unwrap(); + let stop = registration.stop.clone(); + let (video_server, mut video_client) = tokio::io::duplex(64); + let (mut event_server, mut event_client) = UnixStream::pair().unwrap(); + let events = + tokio::spawn( + async move { desktop_events::serve(&mut event_server, subscription).await }, + ); + let media = tokio::spawn(async move { + { + let (read, write) = tokio::io::split(video_server); + serve_separated_io(read, write, relay.session, registration).await + } + }); + let first = tokio::time::timeout( + Duration::from_secs(5), + read_frame::<_, DesktopDown>(&mut video_client), + ) + .await + .unwrap() + .unwrap(); + assert!(matches!(first, DesktopDown::Frame { .. })); + let bytes = video_client.read_u32().await.unwrap(); + assert!( + bytes > 64 && bytes <= 1024 * 1024, + "body must exceed the deterministic video capacity" + ); + // Retain the whole large body unread. On the old single socket, an + // event cannot be parsed before draining exactly these payload bytes. + for seq in 1..=8 { + write_frame( + &mut video_client, + &DesktopUp::Control(DesktopControl::Heartbeat { + seq, + ts_ms: seq + 100, + }), + ) + .await + .unwrap(); + let event = tokio::time::timeout( + Duration::from_secs(2), + read_frame::<_, DesktopDown>(&mut event_client), + ) + .await + .unwrap() + .unwrap(); + assert!( + matches!(event, DesktopDown::Event(rds_core::DesktopEvent::Heartbeat { seq: echoed, ts_ms }) if echoed == seq && ts_ms == seq + 100) + ); + assert!(!media.is_finished(), "video owner unexpectedly ended"); + } + // Event subscriber departure ends its own blocked video writer. + drop(event_client); + tokio::time::timeout(Duration::from_secs(2), events) + .await + .unwrap() + .unwrap() + .unwrap(); + tokio::time::timeout(Duration::from_secs(2), media) + .await + .unwrap() + .unwrap() + .unwrap(); + assert!(stop.is_cancelled()); + relay.server_task.abort(); + } + + #[tokio::test] + async fn separated_readers_keep_events_independent_of_full_media_queue() { + let (video, mut peer) = UnixStream::pair().unwrap(); + let (events, mut event_peer) = UnixStream::pair().unwrap(); + let caps = DesktopCaps { + displays: vec![], + codecs: vec![rds_core::Codec::H264], + }; + let (mut channel, mut events) = + ManagedDesktop::separated(video, events, SessionId([7; 16]), caps, 0); + // Force the FIFO reader's existing one-message channel to fill and + // block. Independent event parsing must remain available. + for seq in 0..3 { + write_frame( + &mut peer, + &DesktopDown::Frame { + header: header(seq), + }, + ) + .await + .unwrap(); + write_payload(&mut peer, &[0xff; 16]).await.unwrap(); + } + write_frame( + &mut event_peer, + &DesktopDown::Event(rds_core::DesktopEvent::InputAck { + seq: 77, + handled_ts_ms: 0, + }), + ) + .await + .unwrap(); + let event = tokio::time::timeout(std::time::Duration::from_secs(1), events.recv()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert!(matches!( + event, + rds_core::DesktopEvent::InputAck { seq: 77, .. } + )); + for seq in 0..3 { + assert!( + matches!(channel.recv().await.unwrap(), Some(ManagedMessage::Frame(frame)) if frame.header.seq == seq) + ); + } + drop(channel); + assert!( + tokio::time::timeout(std::time::Duration::from_secs(1), events.recv()) + .await + .unwrap() + .is_none() + ); + } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn separated_client_uses_authenticated_manager_and_retains_shared_peer() { + use super::super::{Client, Prepared, Server}; + use rds_core::{ + HelloAck, StreamHello, UniHello, + local::{Command, Reply}, + }; + use std::{os::unix::fs::DirBuilderExt, time::Duration}; + let backends = { + #[cfg(feature = "transport-noq")] + { + vec![rds_net::Backend::Iroh, rds_net::Backend::Noq] + } + #[cfg(not(feature = "transport-noq"))] + { + vec![rds_net::Backend::Iroh] + } + }; + for backend in backends { + let config = || rds_net::EndpointConfig { + backend, + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + discovery: false, + ..Default::default() + }; + let remote = rds_net::bind_endpoint(config()).await.unwrap(); + let local = rds_net::bind_endpoint(config()).await.unwrap(); + let remote_task = tokio::spawn({ + let remote = remote.clone(); + async move { + let conn = remote.accept().await.unwrap().await.unwrap(); + let mut workers = tokio::task::JoinSet::new(); + loop { + tokio::select! { + Some(result) = workers.join_next(), if !workers.is_empty() => { result.unwrap(); }, + streams = conn.accept_bi() => { + let (mut send, mut recv) = match streams { Ok(pair) => pair, Err(_) => break }; + let conn = conn.clone(); + workers.spawn(async move { + match read_frame::<_, StreamHello>(&mut recv).await.unwrap() { + StreamHello::Ping { nonce } => { + write_frame(&mut send, &HelloAck::Ok).await.unwrap(); + send.write_all(&nonce.to_be_bytes()).await.unwrap(); + send.finish().unwrap(); + } + StreamHello::DesktopV3 { session, hello, output_height } => { + assert_eq!(output_height, 1080); + write_frame(&mut send, &HelloAck::Desktop(DesktopCaps { + displays: vec![], codecs: vec![rds_core::Codec::H264], + })).await.unwrap(); + rds_desktop::serve_desktop_with(conn, send, recv, hello, rds_desktop::SessionConfig { + frame_route: Some(UniHello::DesktopFrames { id: session }), + input_sink: Some(Box::new(NoopInput)), + producer: Some(Box::new(rds_desktop::SyntheticProducer::new(30, 640, 480, 1024))), + ..Default::default() + }).await.unwrap(); + } + other => panic!("unexpected service {other:?}"), + } + }); + } + } + } + workers.shutdown().await; + } + }); + let root = std::path::Path::new("/tmp") + .canonicalize() + .unwrap() + .join(format!( + "rds-event-route-{}-{}", + std::process::id(), + rand::random::() + )); + std::fs::DirBuilder::new() + .mode(0o700) + .create(&root) + .unwrap(); + let mut manager = Server::start( + Some(Prepared::bind(&root).await.unwrap()), + local.clone(), + None, + ); + let client = Client::new(&root); + let Reply::Connected(id) = client + .request(Command::Connect { + target: rds_net::Ticket::of(&remote).to_string(), + grant: None, + }) + .await + .unwrap() + else { + panic!("missing session") + }; + let (desktop, mut events) = client + .desktop_profile_separated( + Some(id), + rds_core::DesktopHello { + display: 0, + max_fps: 30, + codec: rds_core::Codec::H264, + input_acks: true, + }, + 1080, + ) + .await + .unwrap(); + let controls = desktop.control_handle(); + let seq = controls + .send_input(InputKind::KeyDown { code: 30 }) + .await + .unwrap(); + let ack = tokio::time::timeout(Duration::from_secs(2), events.recv()) + .await + .unwrap() + .unwrap() + .unwrap(); + assert!(matches!(ack, rds_core::DesktopEvent::InputAck { seq: ack, .. } if ack == seq)); + controls.heartbeat().await.unwrap(); + assert!(matches!( + tokio::time::timeout(Duration::from_secs(2), events.recv()) + .await + .unwrap() + .unwrap() + .unwrap(), + rds_core::DesktopEvent::Heartbeat { .. } + )); + drop(desktop); + assert!( + tokio::time::timeout(Duration::from_secs(2), events.recv()) + .await + .unwrap() + .is_none() + ); + assert!( + matches!(client.request(Command::Ping { session: Some(id), nonce: 91 }).await.unwrap(), Reply::Pong { session, .. } if session == id) + ); + manager.close().await.unwrap(); + local.close().await; + remote.close().await; + tokio::time::timeout(Duration::from_secs(2), remote_task) + .await + .unwrap() + .unwrap(); + std::fs::remove_dir_all(root).unwrap(); + } + } } diff --git a/crates/rds-client/src/local/desktop_events.rs b/crates/rds-client/src/local/desktop_events.rs new file mode 100644 index 0000000..332ca79 --- /dev/null +++ b/crates/rds-client/src/local/desktop_events.rs @@ -0,0 +1,273 @@ +//! Same-UID event routes owned by one desktop, never by a peer connection. +use std::{ + collections::BTreeMap, + io, + sync::{Arc, Mutex, Weak}, + time::Duration, +}; + +use rds_core::{ + DesktopEvent, + local::{DesktopDown, ErrorCode, SessionId}, +}; +use rds_net::write_frame; +use tokio::{ + io::AsyncReadExt, + net::UnixStream, + sync::{mpsc, oneshot}, +}; +use tokio_util::sync::CancellationToken; + +const ATTACH_TIMEOUT: Duration = Duration::from_secs(5); +const WRITE_TIMEOUT: Duration = Duration::from_secs(2); +pub(super) const EVENT_CAPACITY: usize = 128; + +struct Pending { + events: mpsc::Receiver, + ready: oneshot::Sender<()>, + stop: CancellationToken, +} + +#[derive(Default)] +pub(super) struct Registry(Mutex>>); + +impl Registry { + pub(super) fn register(self: &Arc) -> Result { + let mut routes = self.0.lock().map_err(|_| ErrorCode::Internal)?; + if routes.len() >= super::MAX_STREAMS { + return Err(ErrorCode::Capacity); + } + // Never replace an existing route, even on a random-id collision. + let route = (0..8) + .map(|_| SessionId(rand::random())) + .find(|id| !routes.contains_key(id)) + .ok_or(ErrorCode::Capacity)?; + let (sender, events) = mpsc::channel(EVENT_CAPACITY); + let (ready, attached) = oneshot::channel(); + let stop = CancellationToken::new(); + routes.insert( + route, + Some(Pending { + events, + ready, + stop: stop.clone(), + }), + ); + Ok(Registration { + route, + sender, + attached, + stop, + registry: Arc::downgrade(self), + }) + } + + pub(super) fn take(&self, route: SessionId) -> Result { + let pending = self + .0 + .lock() + .map_err(|_| ErrorCode::Internal)? + .get_mut(&route) + .and_then(Option::take) + .ok_or(ErrorCode::NotFound)?; + Ok(Subscription { + events: pending.events, + ready: Some(pending.ready), + stop: pending.stop, + }) + } +} + +pub(super) struct Registration { + pub(super) route: SessionId, + pub(super) sender: mpsc::Sender, + attached: oneshot::Receiver<()>, + pub(super) stop: CancellationToken, + registry: Weak, +} + +impl Registration { + pub(super) async fn attached(&mut self) -> io::Result<()> { + tokio::time::timeout(ATTACH_TIMEOUT, &mut self.attached) + .await + .map_err(|_| { + io::Error::new( + io::ErrorKind::TimedOut, + "desktop event attachment timed out", + ) + })? + .map_err(|_| { + io::Error::new(io::ErrorKind::BrokenPipe, "desktop event attachment ended") + }) + } +} + +impl Drop for Registration { + fn drop(&mut self) { + self.stop.cancel(); + if let Some(registry) = self.registry.upgrade() { + registry + .0 + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .remove(&self.route); + } + } +} + +pub(super) struct Subscription { + events: mpsc::Receiver, + ready: Option>, + stop: CancellationToken, +} + +impl Drop for Subscription { + fn drop(&mut self) { + self.stop.cancel(); + } +} + +/// EOF or unexpected caller bytes terminate this desktop, never its shared peer. +/// The ready fence is after the successful subscription reply write. +pub(super) async fn serve(stream: &mut UnixStream, subscription: Subscription) -> io::Result<()> { + let (reader, writer) = stream.split(); + serve_io(reader, writer, subscription).await +} + +async fn serve_io( + mut reader: R, + mut writer: W, + mut subscription: Subscription, +) -> io::Result<()> { + if let Some(ready) = subscription.ready.take() { + ready + .send(()) + .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "desktop owner ended"))?; + } + let mut probe = [0]; + let events = async { + while let Some(event) = subscription.events.recv().await { + let started = std::time::Instant::now(); + let result = tokio::time::timeout( + WRITE_TIMEOUT, + write_frame(&mut writer, &DesktopDown::Event(event)), + ) + .await + .unwrap_or_else(|_| { + Err(io::Error::new( + io::ErrorKind::TimedOut, + "desktop event write timed out", + )) + }); + if let Err(error) = result { + tracing::warn!(error_kind = ?error.kind(), elapsed_ms = started.elapsed().as_millis() as u64, + "desktop event write failed"); + return Err(error); + } + } + Ok(()) + }; + tokio::select! { + biased; + _ = subscription.stop.cancelled() => Ok(()), + bytes = reader.read(&mut probe) => match bytes? { + 0 => Ok(()), + _ => Err(io::Error::new(io::ErrorKind::InvalidData, "unexpected desktop event input")), + }, + result = events => result, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn routes_are_bounded_single_claim_and_owner_scoped() { + let registry = Arc::new(Registry::default()); + let mut owners = Vec::new(); + for _ in 0..super::super::MAX_STREAMS { + owners.push(registry.register().unwrap()); + } + assert!(matches!(registry.register(), Err(ErrorCode::Capacity))); + let route = owners[0].route; + let subscription = registry.take(route).unwrap(); + assert!(matches!(registry.take(route), Err(ErrorCode::NotFound))); + assert!( + matches!(registry.register(), Err(ErrorCode::Capacity)), + "claimed routes must remain reserved until their owner ends" + ); + let stop = owners[0].stop.clone(); + drop(owners.remove(0)); + assert!(stop.is_cancelled()); + drop(subscription); + let remaining = owners[0].route; + assert!( + !owners[0].stop.is_cancelled(), + "another desktop was canceled" + ); + drop(owners); + assert!(matches!(registry.take(remaining), Err(ErrorCode::NotFound))); + assert!(registry.0.lock().unwrap().is_empty()); + assert_eq!(Arc::strong_count(®istry), 1); + } + + #[test] + fn abandoned_subscription_cancels_its_owner_before_ready() { + let registry = Arc::new(Registry::default()); + let mut owner = registry.register().unwrap(); + let subscription = registry.take(owner.route).unwrap(); + assert!(matches!( + owner.attached.try_recv(), + Err(oneshot::error::TryRecvError::Empty) + )); + drop(subscription); + assert!(owner.stop.is_cancelled()); + assert!(matches!( + owner.attached.try_recv(), + Err(oneshot::error::TryRecvError::Closed) + )); + } + + #[tokio::test] + async fn subscription_rejects_extra_input_and_releases_its_route() { + use tokio::io::AsyncWriteExt; + let registry = Arc::new(Registry::default()); + let mut owner = registry.register().unwrap(); + let subscription = registry.take(owner.route).unwrap(); + let (mut server, mut client) = UnixStream::pair().unwrap(); + let task = tokio::spawn(async move { serve(&mut server, subscription).await }); + owner.attached().await.unwrap(); + client.write_all(&[1]).await.unwrap(); + let error = tokio::time::timeout(Duration::from_secs(2), task) + .await + .unwrap() + .unwrap() + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + assert!(owner.stop.is_cancelled()); + } + #[tokio::test] + async fn blocked_event_writer_has_a_deadline_and_cancels_only_its_owner() { + let registry = Arc::new(Registry::default()); + let mut owner = registry.register().unwrap(); + let other = registry.register().unwrap(); + let subscription = registry.take(owner.route).unwrap(); + let (server, _unread_client) = tokio::io::duplex(1); + let (read, write) = tokio::io::split(server); + let task = tokio::spawn(async move { serve_io(read, write, subscription).await }); + owner.attached().await.unwrap(); + owner + .sender + .try_send(DesktopEvent::Heartbeat { seq: 1, ts_ms: 9 }) + .unwrap(); + let error = tokio::time::timeout(Duration::from_secs(4), task) + .await + .unwrap() + .unwrap() + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::TimedOut); + assert!(owner.stop.is_cancelled()); + assert!(!other.stop.is_cancelled()); + } +} diff --git a/crates/rds-client/src/local/mod.rs b/crates/rds-client/src/local/mod.rs index 663d90d..bd53577 100644 --- a/crates/rds-client/src/local/mod.rs +++ b/crates/rds-client/src/local/mod.rs @@ -179,6 +179,7 @@ async fn run( let mut workers = JoinSet::new(); let streams = Arc::new(Semaphore::new(MAX_STREAMS)); let transfers = Arc::new(Semaphore::new(sync::MAX_TRANSFERS)); + let events = Arc::new(desktop_events::Registry::default()); let mut sample = tokio::time::interval(Duration::from_secs(1)); sample.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); let result = loop { @@ -204,10 +205,11 @@ async fn run( let directory = directory.clone(); let streams = streams.clone(); let transfers = transfers.clone(); + let events = events.clone(); workers.spawn(async move { // Responses carry typed reasons; never log the request or // an upstream error containing a ticket/grant/private path. - let _ = serve(stream, shared, endpoint, directory, streams, transfers).await; + let _ = serve(stream, shared, endpoint, directory, streams, transfers, events).await; }); } } @@ -219,14 +221,20 @@ async fn run( } mod desktop; +mod desktop_events; -pub use desktop::{ManagedControl, ManagedDesktop, ManagedMessage, RelayedFrame}; +pub use desktop::{ManagedControl, ManagedDesktop, ManagedEvents, ManagedMessage, RelayedFrame}; struct Output { reply: Reply, reservation: Option, tcp: Option<(RequestStreams, OwnedSemaphorePermit)>, - desktop: Option<(rds_desktop::client::DesktopSession, OwnedSemaphorePermit)>, + desktop: Option<( + rds_desktop::client::DesktopSession, + OwnedSemaphorePermit, + Option, + )>, + events: Option<(desktop_events::Subscription, OwnedSemaphorePermit)>, } impl Output { @@ -236,6 +244,7 @@ impl Output { reservation: None, tcp: None, desktop: None, + events: None, } } } @@ -247,6 +256,7 @@ async fn serve( directory: Option, streams: Arc, transfers: Arc, + events: Arc, ) -> Result<(), Error> { let request: Request = tokio::time::timeout(PRELUDE_TIMEOUT, read_frame(&mut stream)) .await @@ -265,7 +275,7 @@ async fn serve( tokio::select! { biased; _ = stream.read(&mut probe) => return Ok(()), - result = tokio::time::timeout(timeout, execute(request.command, &shared, &endpoint, directory, streams, transfers)) => result.unwrap_or(Err(ErrorCode::Timeout)), + result = tokio::time::timeout(timeout, execute(request.command, &shared, &endpoint, directory, streams, transfers, &events)) => result.unwrap_or(Err(ErrorCode::Timeout)), } }; let response = Response { @@ -286,8 +296,16 @@ async fn serve( // data/FIN instead of resetting a successfully completed upload. drop(tcp.release()); } - if let Some((session, _permit)) = output.desktop { - desktop::serve(&mut stream, session).await?; + if let Some((session, _permit, registration)) = output.desktop { + match registration { + Some(registration) => { + desktop::serve_separated(&mut stream, session, registration).await? + } + None => desktop::serve(&mut stream, session).await?, + } + } + if let Some((subscription, _permit)) = output.events { + desktop_events::serve(&mut stream, subscription).await?; } } Ok(()) @@ -309,6 +327,7 @@ async fn execute( directory: Option, streams: Arc, transfers: Arc, + events: &Arc, ) -> Result { match command { Command::Sync { session, operation } => { @@ -398,6 +417,7 @@ async fn execute( reservation: Some(reservation), tcp: None, desktop: None, + events: None, }) } Command::Renew { session, grant } => { @@ -421,6 +441,7 @@ async fn execute( reservation: Some(reservation), tcp: None, desktop: None, + events: None, }) } Command::Select { session } => { @@ -467,14 +488,32 @@ async fn execute( reservation: None, tcp: Some((RequestStreams::new(pair), permit)), desktop: None, + events: None, + }) + } + Command::DesktopEvents { route } => { + let permit = streams + .try_acquire_owned() + .map_err(|_| ErrorCode::Capacity)?; + let subscription = events.take(route)?; + Ok(Output { + reply: Reply::DesktopEventsOpened { route }, + reservation: None, + tcp: None, + desktop: None, + events: Some((subscription, permit)), }) } Command::Desktop { session, ref hello } + | Command::DesktopSeparated { + session, ref hello, .. + } | Command::DesktopProfile { session, ref hello, .. } => { let output_height = match &command { - Command::DesktopProfile { output_height, .. } => Some(*output_height), + Command::DesktopProfile { output_height, .. } + | Command::DesktopSeparated { output_height, .. } => Some(*output_height), _ => None, }; let permit = streams @@ -497,14 +536,28 @@ async fn execute( ) .await .map_err(|_| ErrorCode::Remote)?; - Ok(Output { - reply: Reply::DesktopOpened { + let registration = if matches!(command, Command::DesktopSeparated { .. }) { + Some(events.register()?) + } else { + None + }; + let reply = match ®istration { + Some(registration) => Reply::DesktopSeparatedOpened { session, caps: remote.caps().clone(), + route: registration.route, }, + None => Reply::DesktopOpened { + session, + caps: remote.caps().clone(), + }, + }; + Ok(Output { + reply, reservation: None, tcp: None, - desktop: Some((remote, permit)), + desktop: Some((remote, permit, registration)), + events: None, }) } } @@ -526,7 +579,11 @@ impl Client { async fn exchange(&self, command: Command) -> Result<(Reply, UnixStream), Error> { let body = matches!( command, - Command::OpenTcp { .. } | Command::Desktop { .. } | Command::DesktopProfile { .. } + Command::OpenTcp { .. } + | Command::Desktop { .. } + | Command::DesktopProfile { .. } + | Command::DesktopSeparated { .. } + | Command::DesktopEvents { .. } ); let timeout = if matches!(command, Command::Sync { .. }) { SYNC_TIMEOUT + Duration::from_secs(10) @@ -565,7 +622,11 @@ impl Client { pub async fn request(&self, command: Command) -> Result { if matches!( command, - Command::OpenTcp { .. } | Command::Desktop { .. } | Command::DesktopProfile { .. } + Command::OpenTcp { .. } + | Command::Desktop { .. } + | Command::DesktopProfile { .. } + | Command::DesktopSeparated { .. } + | Command::DesktopEvents { .. } ) { return Err(Error::Protocol); } @@ -611,6 +672,39 @@ impl Client { } } + /// Independently received events and bounded FIFO video on same-UID sockets. + /// Requires the additive separated-channel extension; no legacy fallback. + pub async fn desktop_profile_separated( + &self, + session: Option, + hello: rds_core::DesktopHello, + output_height: u32, + ) -> Result<(ManagedDesktop, ManagedEvents), Error> { + let display = hello.display; + let (reply, stream) = self + .exchange(Command::DesktopSeparated { + session, + hello: Box::new(hello), + output_height, + }) + .await?; + let Reply::DesktopSeparatedOpened { + session, + caps, + route, + } = reply + else { + return Err(Error::Protocol); + }; + let (reply, events) = self.exchange(Command::DesktopEvents { route }).await?; + if !matches!(reply, Reply::DesktopEventsOpened { route: opened } if opened == route) { + return Err(Error::Protocol); + } + Ok(ManagedDesktop::separated( + stream, events, session, caps, display, + )) + } + pub async fn snapshot(&self) -> Result { match self.request(Command::List).await? { Reply::Snapshot(snapshot) => Ok(snapshot), diff --git a/crates/rds-core/src/local.rs b/crates/rds-core/src/local.rs index 81c9a4b..97bbdbe 100644 --- a/crates/rds-core/src/local.rs +++ b/crates/rds-core/src/local.rs @@ -153,6 +153,17 @@ pub enum Command { hello: Box, output_height: u32, }, + /// Additive v5 extension: video/control on this socket, events on a + /// separately attached same-UID socket. Older managers refuse it. + DesktopSeparated { + session: Option, + hello: Box, + output_height: u32, + }, + /// Claim one process-lifetime event route exactly once. + DesktopEvents { + route: SessionId, + }, } /// Never Debug: paths may contain private information. @@ -222,6 +233,14 @@ pub enum Reply { session: SessionId, stats: SyncStats, }, + DesktopSeparatedOpened { + session: SessionId, + caps: crate::DesktopCaps, + route: SessionId, + }, + DesktopEventsOpened { + route: SessionId, + }, } /// Stable reasons, with no upstream error, address, path or credential payload. @@ -383,6 +402,76 @@ mod tests { } } + #[test] + fn separated_desktop_extension_preserves_existing_discriminants() { + let id = SessionId([3; 16]); + let hello = || { + Box::new(crate::DesktopHello { + display: 0, + max_fps: 60, + codec: crate::Codec::H264, + input_acks: true, + }) + }; + for (command, tag) in [ + (Command::Ticket, 7), + ( + Command::Desktop { + session: Some(id), + hello: hello(), + }, + 10, + ), + ( + Command::DesktopProfile { + session: Some(id), + hello: hello(), + output_height: 1080, + }, + 11, + ), + ( + Command::DesktopSeparated { + session: Some(id), + hello: hello(), + output_height: 1080, + }, + 12, + ), + (Command::DesktopEvents { route: id }, 13), + ] { + let bytes = postcard::to_stdvec(&command).unwrap(); + assert_eq!(bytes[0], tag); + postcard::from_bytes::(&bytes).unwrap(); + } + let caps = || crate::DesktopCaps { + displays: vec![], + codecs: vec![crate::Codec::H264], + }; + for (reply, tag) in [ + ( + Reply::DesktopOpened { + session: id, + caps: caps(), + }, + 6, + ), + ( + Reply::DesktopSeparatedOpened { + session: id, + caps: caps(), + route: id, + }, + 9, + ), + (Reply::DesktopEventsOpened { route: id }, 10), + ] { + let bytes = postcard::to_stdvec(&reply).unwrap(); + assert_eq!(bytes[0], tag); + postcard::from_bytes::(&bytes).unwrap(); + } + } + proptest::proptest! { #[test] fn local_decoders_never_panic(bytes in proptest::collection::vec(proptest::prelude::any::(), 0..1024), text in ".*") { diff --git a/crates/rds-desktop/src/render/viewer.rs b/crates/rds-desktop/src/render/viewer.rs index 053df7f..6b46407 100644 --- a/crates/rds-desktop/src/render/viewer.rs +++ b/crates/rds-desktop/src/render/viewer.rs @@ -116,6 +116,7 @@ pub struct ViewerReport { pub encode_to_send_p95_ms: Option, pub control_rtt_ms: Option, pub input_acks: u64, + pub managed_events_separated: bool, pub input_ack_p50_ms: Option, pub input_ack_p95_ms: Option, pub input_queue_p95_ms: Option, @@ -312,6 +313,9 @@ impl ViewerHandle { tracing::trace!(target: "rds_desktop::input_timing", input_seq=event.seq,event_class,event_created_ms=event.event_ts_ms,input_sent_ms=sent_ms,queue_ms=sent_ms.saturating_sub(event.event_ts_ms),"native input dispatched"); } } + pub fn managed_events_separated(&self) { + lock(&self.state).report.managed_events_separated = true; + } pub fn input_ack(&self, seq: u64) { let mut state = lock(&self.state); state.report.input_acks += 1; @@ -515,6 +519,7 @@ impl App { match result { Err(error) => self.fail(event_loop, error), Ok(DrawOutcome::Presented) => { + lock(&self.handle.state).render_stage = "presented".into(); if let Some(frame) = pending { let mut state = lock(&self.handle.state); state.report.frames_submitted += 1; diff --git a/docs/architecture.md b/docs/architecture.md index 2151925..f6f537b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -81,11 +81,13 @@ a wgpu surface and one pending BGRA image. The CLI owns the cancelable network worker and reopens a desktop session on the same authenticated peer after loss. Per-session wire routing and broader native media acceptance remain separate. -Managed native input/heartbeat dispatch and its decoded-progress watchdog are -retained async legs independent of receive/decode waits. They share the existing -serialized control writer, with bounded writes and EOF teardown after either -leg ends. Incoming event observation still shares the bounded media IPC route; -see [the native contract](native-viewer.md) for the tested boundary. +Managed native input/heartbeat dispatch, incoming control-event observation and +ordered media/decode have independent retained futures. The additive local +separated-desktop API attaches one same-UID event socket using a bounded, +single-claim, process-lifetime route. Video FIFO backpressure cannot block its +input ACK/heartbeat observation. Both sockets share desktop termination, without +closing the manager's unrelated peer streams. Existing combined APIs and remote +wire formats remain; see [the native contract](native-viewer.md). Owned path policy now uses [validated eligibility](path-selection.md): only the handshake path is seeded; application-opened candidates stay Backup until an diff --git a/docs/local-sessions.md b/docs/local-sessions.md index 7d9dab2..5368091 100644 --- a/docs/local-sessions.md +++ b/docs/local-sessions.md @@ -296,3 +296,30 @@ for native/direct desktop use, unchanged. Each step needs code, failure tests, platform evidence and an updated receipt. All remediation waves remain open. This sequence preserves every other task in the [full remediation plan](remediation-plan.md). + + +## Separated desktop events (additive local v5 extension) + +`DesktopSeparated` replies with `DesktopSeparatedOpened` and an opaque random +route; `DesktopEvents` claims that route exactly once and returns +`DesktopEventsOpened`. Both requests pass the existing same-UID socket identity +and directory checks. All existing command/reply discriminants stay unchanged; +older managers reject the new commands. No remote ALPN, grant or frame format +changes. The native viewer and local agent need a coordinated update; the +legacy combined/headless APIs remain available explicitly. + +The first socket carries ordered encoded video and upward controls. The second +carries only `DesktopDown::Event` messages and clean EOF; a frame there is a +protocol error. The attachment-ready fence follows the successful subscription +reply. An unattached owner waits at most five seconds. The bounded route registry +holds no peer connection and its guard removes/cancels a route on every exit. +Routes cannot be claimed twice or reused after cancellation/manager restart. + +Each desktop/event socket uses a stream permit (64 total), and the existing +96-worker limit also applies. Queues hold 128 small events, one queued encoded +message and one being read. Encoded FIFO and 32 MiB payload bounds are unchanged. +The event writer has a two-second per-message limit; queue overflow fails the +desktop explicitly, rather than dropping an ACK and claiming healthy progress. +Both owned client readers abort on desktop drop. Subscriber EOF, unexpected +caller bytes, video/control EOF or remote end cancels only this desktop. The +manager's unrelated TCP/sync sessions stay live. diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 1d59b36..391cf17 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -6,16 +6,27 @@ decoding stays on bounded blocking workers. The serving device needs a real capture/input backend; the current Linux implementation uses X11/XTEST. macOS capture, VideoToolbox and Wayland serving remain separate work. -Managed native control dispatch has its own retained async future, polled -concurrently with the receive/decode future. Input writes, heartbeat and the -15-second decoded-progress watchdog continue while a blocking codec worker is -pending. Each control write remains bounded to two seconds; codec calls retain -the existing five-second caller bound and global worker limits. Decoder state -and encoded ordering are unchanged. Either leg's completion or cancellation -ends the desktop IPC by EOF; it never appends a control frame after canceling a -potentially partial write. Incoming acknowledgement observation can still wait -behind decode in the bounded IPC receive path. See the -[dispatch regression evidence](reports/rds-managed-control-20261002.md). +Managed native dispatch and acknowledgement observation use separate retained +async legs from ordered receive/decode. The agent's separated local extension +keeps video and upward controls on the first same-UID Unix socket and transfers +input ACKs, heartbeat echoes and clipboard-ready metadata over a second socket. +A blocked video-body write or a busy decoder cannot hold event observation. +`managed_events_separated` in the viewer report identifies this installed mode. + +Each control write remains bounded to two seconds; codec calls retain their +five-second caller bound and global worker limits. The 15-second decoded-progress +watchdog still requires actual decoded frames. This change does not keep a +nonfunctional video session healthy simply because heartbeats continue. +Decoder state, FIFO encoded ordering and the one-message media queue remain. +Either socket's departure ends this desktop by EOF, with both owned reader tasks +aborted; no frame is appended after canceling a potentially partial write. +Shared TCP/sync connections and other desktops are not closed. + +Legacy managed/headless APIs retain their original combined socket. The native +viewer requires the additive v5 separated-channel commands; an older agent +refuses them without silently reverting to the coupled route. Update the local +agent and viewer together. Remote desktop framing and grants are unchanged. +See [the independent-event qualification](reports/rds-event-isolation-20261002.md). The viewer uses ordinary OS window stacking and can move behind other applications. An occluded Metal surface may pause presentation. Returning focus diff --git a/docs/reports/rds-event-isolation-20261002.md b/docs/reports/rds-event-isolation-20261002.md new file mode 100644 index 0000000..9604134 --- /dev/null +++ b/docs/reports/rds-event-isolation-20261002.md @@ -0,0 +1,66 @@ +# Managed desktop event isolation + +Scope: W2.4/W6.7 managed native latency and diagnostic remediation. Physical +input-to-pixel, loaded Full HD quality and sustained-network acceptance remain +open. This receipt does not close a milestone or predict WAN performance. + +The remote control stream had its own priority, but the local bridge serialized +its events after video bodies on one Unix socket. The viewer then awaited native +decode before observing events from that combined queue. A healthy remote input +ACK or heartbeat could therefore wait on application media backpressure. + +The new additive local v5 extension attaches an independent same-UID event +socket using one bounded, random, single-claim route. Video and upward controls +stay on the original socket. Agent control/event/video futures and native +control/event/decode futures are retained independently. Either socket departure +ends only that desktop; the shared peer remains usable. Source geometry, +encoded FIFO, one-message media queue, 32 MiB payload cap and global decoder +permits remain. Each socket uses the existing stream/worker admission; event +queues hold 128 entries, attachment has five seconds and event writes two seconds. +An unavailable/full event sink fails explicitly. There is no dropped-ACK success. + +Existing commands/replies retain their wire tags; the original combined/headless +APIs stay explicit. An older manager refuses the extension; upgrade the local +agent and native viewer together. Remote ALPN, desktop framing and grants do not +change. A report flag identifies the separated managed mode. A Presented redraw +without a new CPU image now sets the stage to presented, avoiding a stale +acquiring-surface diagnostic despite an idle AppKit main thread. + +## Regression evidence + +- A real QUIC synthetic desktop feeds the agent bridge while a deterministic + 64-byte video pipe retains its frame body unread. Eight heartbeat echoes must + cross the independent event socket before any payload is drained. Closing that + subscriber must end the blocked video owner. The effective synthetic body is + selected by the sender bitrate, rather than assumed to equal its size ceiling. +- Real Unix client reader queues retain three ordered encoded frames with the + one-message FIFO full. An independent input ACK is observed without draining + media. Desktop drop closes the owned event reader even if a control handle lives. +- Authenticated local manager requests and the actual new Client API run against + both Iroh/Noq loopback peers. Input ACK/heartbeat arrive independently; dropping + the desktop keeps the same peer usable for a subsequent Ping. +- Route capacity, single claim, owner-only cleanup, aborted subscription and + unexpected event-socket input have direct lifecycle/negative assertions. +- Explicit serialization tests preserve existing command/reply discriminants. + +Initial development checks retained a borrow error, a misspelled test event +field and an incorrect fixed synthetic-body-size assumption; those are not +passing evidence. Broader platform, lint, release and installed qualification +are recorded with the final candidate. Raw logs are retained privately. No +incomplete lane is claimed green. + +The initial macOS five-crate all-feature check passed 309 tests, zero failures, +one explicit OpenSSH/account ignore in 42 groups. Its first broad lint attempt +stopped with OS disk exhaustion; the failure is retained. Only regenerable +incremental build cache was removed, and subsequent checks disable incremental +cache growth. Final verification includes the separated remote-refusal path +and its unchanged full TCP budget assertion. Platform/source finalization and +installed-device acceptance remain separate. + +After the route-reservation review, final client tests passed 21/0/0. Claimed +routes remain reserved until their owner ends, so even an active-id collision +cannot replace a different registration. A deterministic one-byte event pipe +proves the two-second event-write deadline, owner cancellation and preservation +of a second desktop. The updated remote-refusal/full-TCP-budget test also passed +on both backend variants. Default and all-feature workspace lint passed; final +all-feature lint/release and Linux CI are being qualified against the candidate. diff --git a/docs/research.md b/docs/research.md index 4fa750a..94f9304 100644 --- a/docs/research.md +++ b/docs/research.md @@ -643,3 +643,23 @@ 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. + + +## 2026-10-02 independent managed control-event observation + +The maintained [WebRTC pacer design](https://webrtc.googlesource.com/src/+/refs/heads/main/modules/pacing/g3doc/index.md) +separates queues for different traffic classes instead of allowing a buffered +media track to hold another class. [QUIC stream prioritization](https://www.rfc-editor.org/rfc/rfc9000.html#section-2.3) +is an application responsibility; independent remote streams do not preserve +that independence if the application serializes them again onto one local pipe. + +RDS's managed bridge did that serialization: video-body writes and native decode +waits could delay input ACK/heartbeat observation despite the remote control +stream's higher priority. The separated local event extension removes those +two application waits while preserving existing admission, decode and ordering +bounds. It does not bypass QUIC congestion/flow control or guarantee physical +network failover. The [measurement boundary](reports/rds-event-isolation-20261002.md) +records blocked-body, FIFO and lifecycle regressions separately from installed +native quality and latency acceptance. The old Iroh 0.96 network-change regression +reported in upstream's January release note is not assumed to exist in pinned +1.2; dependency changes require current source evidence and same-harness parity.