From a8e78e19c5795b96df530154671229d25296bbf6 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 11:59:34 +0500 Subject: [PATCH 1/2] =?UTF-8?q?feat(local):=20managed=20desktop=20viewer?= =?UTF-8?q?=20channel=20=E2=80=94=20local=20wire=20v5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Local IPC moves to version 5: Command::Desktop opens a remote DesktopV2 session on the pinned managed connection in the new relay mode, replies DesktopOpened, then the socket speaks DesktopDown/DesktopUp — postcard control messages plus u32-length-prefixed raw encoded payloads bounded at 32 MiB, deliberately beside the 64 KiB frame bound. rds-desktop gains SessionOpts::relay_encoded: the session keeps sequence/stale-frame and transport IDR discipline but publishes EncodedDelivery (header + encoded Bytes) to a bounded tap instead of decoding, so the manager never links a codec and stays buildable headless. control_sender()/send_control() forward viewer controls verbatim; the viewer-side RelayDecoder reapplies wait-for-keyframe and broken-chain discipline with the same 500 ms NeedIdr rate limit the in-session path uses. rds-client serves the channel: the pump forwards encoded frames, events and controls until viewer Finished/EOF, remote end or body error, then drops the session and releases the shared stream permit. Client::desktop returns ManagedDesktop — split-socket receive plus a cloneable ManagedControl that serializes postcard writes so concurrent senders cannot interleave. rds-cli defaults `rds desktop` and adds `session desktop` to the managed path with RelayDecoder decode and IDR forwarding; --direct keeps native in-process sessions. Tests: wire v5 enum round-trips and proptest decoders; payload bounds, truncation, Finished/EOF and concurrent-sender unit tests; three real-loopback serve e2e tests (frames + heartbeat echo + clean finish, remote drop, caller EOF); a relay-mode transport e2e; and a real-agent managed open asserting clean refusal without permit leaks. Refs: remediation-plan W2.4 (viewer manager API) --- Cargo.lock | 1 + crates/rds-agent/tests/local_manager.rs | 52 ++ crates/rds-cli/src/managed.rs | 101 +++- crates/rds-client/Cargo.toml | 1 + crates/rds-client/src/local/desktop.rs | 660 ++++++++++++++++++++++++ crates/rds-client/src/local/mod.rs | 69 ++- crates/rds-core/src/local.rs | 132 ++++- crates/rds-desktop/src/client.rs | 119 +++++ crates/rds-desktop/tests/session_v2.rs | 77 +++ 9 files changed, 1204 insertions(+), 8 deletions(-) create mode 100644 crates/rds-client/src/local/desktop.rs diff --git a/Cargo.lock b/Cargo.lock index 43f1384..6bd4a10 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3454,6 +3454,7 @@ dependencies = [ "postcard", "rand 0.10.3", "rds-core", + "rds-desktop", "rds-discovery", "rds-net", "rds-observe", diff --git a/crates/rds-agent/tests/local_manager.rs b/crates/rds-agent/tests/local_manager.rs index 6f08dca..d459744 100644 --- a/crates/rds-agent/tests/local_manager.rs +++ b/crates/rds-agent/tests/local_manager.rs @@ -955,3 +955,55 @@ async fn abandoned_silent_tcp_bodies_release_capacity_but_half_close_preserves_r "abandoned IPC bodies retained all 64 stream slots" ); } + +/// `Client::desktop` exercises the whole managed path — connect, permit, +/// stream open — and reports a clean refusal when the peer cannot serve +/// desktop. Without the agent `desktop` feature the service gate refuses; +/// with it but headless, capture fails — either way the channel closes +/// and no stream slot leaks. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn managed_desktop_reports_remote_refusal_without_leaking() { + for backend in backends() { + let root = Scratch::new(); + let path = root.0.join("control"); + let local = bind_endpoint(config(backend)).await.unwrap(); + let prepared = Prepared::bind(&path).await.unwrap(); + let mut server = Server::start(Some(prepared), local.clone(), None); + let client = Client::new(&path); + let mut tasks = JoinSet::new(); + let (peer_ep, _tcp) = peer(backend, &local, b'D', &mut tasks).await; + let session = connect(&client, Ticket::of(&peer_ep).to_string()).await; + let result = tokio::time::timeout( + Duration::from_secs(10), + client.desktop( + Some(session), + rds_core::DesktopHello { + display: 0, + max_fps: 30, + codec: rds_core::Codec::H264, + input_acks: false, + }, + ), + ) + .await + .expect("managed desktop open hung"); + assert!( + matches!(result, Err(Error::Rejected(ErrorCode::Remote))), + "expected clean remote refusal, got {result:?}" + ); + // The refused open must not park a stream permit: open_tcp still + // has its full budget. + let mut held = Vec::new(); + for _ in 0..64 { + held.push( + client + .open_tcp(session, _tcp.clone()) + .await + .expect("stream slots leaked"), + ); + } + drop(held); + server.close().await.unwrap(); + tasks.abort_all(); + } +} diff --git a/crates/rds-cli/src/managed.rs b/crates/rds-cli/src/managed.rs index b2d447e..ae2502a 100644 --- a/crates/rds-cli/src/managed.rs +++ b/crates/rds-cli/src/managed.rs @@ -86,6 +86,16 @@ enum Action { #[arg(long, default_value = "64")] max_connections: std::num::NonZeroU16, }, + /// Open a managed desktop channel through the selected or pinned + /// session; the agent relays encoded frames and the CLI decodes. + Desktop { + #[arg(long)] + session: Option, + #[arg(long, default_value = "0")] + display: u32, + #[arg(long, default_value = "30")] + max_fps: u32, + }, } pub async fn run(options: Options, directory: PathBuf) -> anyhow::Result<()> { @@ -158,6 +168,14 @@ pub async fn run(options: Options, directory: PathBuf) -> anyhow::Result<()> { let session = client.selected(session).await?; return forward(&client, session, bind, remote, max_connections).await; } + Action::Desktop { + session, + display, + max_fps, + } => { + let session = client.selected(session).await?; + return desktop(&client, session, display, max_fps).await; + } }; match client.request(command).await? { Reply::Connected(id) => println!("{id}"), @@ -226,11 +244,7 @@ pub async fn run_default( ); target } - super::Command::Desktop { .. } => { - anyhow::bail!( - "desktop does not yet have a manager API; use --direct with a separate --key-file" - ); - } + super::Command::Desktop { target, .. } => target, _ => anyhow::bail!("unsupported managed command"), }; let grant = read_grant(grant).await?; @@ -297,6 +311,11 @@ pub async fn run_default( } => { forward(&client, session, bind, remote, max_connections).await?; } + super::Command::Desktop { + display, max_fps, .. + } => { + desktop(&client, session, display, max_fps).await?; + } _ => anyhow::bail!("unsupported managed command"), } Ok(()) @@ -348,3 +367,75 @@ async fn forward( result = tokio::signal::ctrl_c() => result.map_err(Into::into), } } + +/// Managed desktop viewer: the agent relays encoded frames, the CLI owns +/// decode and resync. Mirrors the direct viewer's stats output; control +/// events (heartbeat echoes, input acks) ride the same channel. +#[cfg(feature = "desktop")] +async fn desktop( + client: &Client, + session: SessionId, + display: u32, + max_fps: u32, +) -> anyhow::Result<()> { + use rds_client::local::ManagedMessage; + use rds_desktop::client::{RelayDecoder, RelayOutcome}; + + let mut channel = client + .desktop( + Some(session), + rds_core::DesktopHello { + display, + max_fps, + codec: rds_core::Codec::H264, + input_acks: false, + }, + ) + .await?; + println!("desktop caps: {:?}", channel.caps); + let mut decoder = RelayDecoder::new(); + let mut count = 0u64; + let start = std::time::Instant::now(); + loop { + let message = tokio::select! { + message = channel.recv() => match message? { + Some(message) => message, + None => break, + }, + _ = tokio::signal::ctrl_c() => { + channel.finish().await.ok(); + return Ok(()); + } + }; + match message { + ManagedMessage::Frame(frame) => match decoder.push(&frame.header, frame.payload) { + RelayOutcome::Frame(raw) => { + count += 1; + if count.is_multiple_of(30) { + let secs = start.elapsed().as_secs_f64(); + println!( + "decoded {count} frames, {:.1} fps, last {}x{}", + count as f64 / secs, + raw.width, + raw.height + ); + } + } + RelayOutcome::Pending => {} + RelayOutcome::NeedIdr => channel.request_idr().await?, + }, + ManagedMessage::Event(_) => {} + } + } + Ok(()) +} + +#[cfg(not(feature = "desktop"))] +async fn desktop( + _client: &Client, + _session: SessionId, + _display: u32, + _max_fps: u32, +) -> anyhow::Result<()> { + anyhow::bail!("rds built without desktop support; enable the `desktop` feature") +} diff --git a/crates/rds-client/Cargo.toml b/crates/rds-client/Cargo.toml index fc86190..530c6eb 100644 --- a/crates/rds-client/Cargo.toml +++ b/crates/rds-client/Cargo.toml @@ -16,6 +16,7 @@ blake3.workspace = true postcard = { workspace = true, features = ["use-std"] } rand.workspace = true rds-core.workspace = true +rds-desktop.workspace = true rds-discovery.workspace = true rds-net.workspace = true rds-observe.workspace = true diff --git a/crates/rds-client/src/local/desktop.rs b/crates/rds-client/src/local/desktop.rs new file mode 100644 index 0000000..41007ca --- /dev/null +++ b/crates/rds-client/src/local/desktop.rs @@ -0,0 +1,660 @@ +//! Managed desktop channel: the session manager holds the remote +//! `DesktopSession` and relays encoded frames and control events over IPC. +//! The viewer decodes — the manager never links a codec. +//! +//! Wire shape mirrors the remote one: postcard `DesktopDown`/`DesktopUp` +//! messages via `write_frame`, and each `Frame` header followed by a +//! big-endian u32 length plus raw encoded payload (≤ `MAX_DESKTOP_PAYLOAD`). + +use std::io; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Instant; + +use rds_core::local::{DesktopDown, DesktopUp, MAX_DESKTOP_PAYLOAD, SessionId}; +use rds_core::{DesktopCaps, DesktopControl, FrameHeader, InputEvent, InputKind}; +use rds_net::{read_frame, write_frame}; +use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}; +use tokio::net::UnixStream; +use tokio::net::unix::{OwnedReadHalf, OwnedWriteHalf}; +use tokio::sync::Mutex; + +fn invalid() -> io::Error { + io::Error::new(io::ErrorKind::InvalidData, "invalid local desktop body") +} + +/// Write one encoded payload with its u32 length prefix. +async fn write_payload(writer: &mut W, payload: &[u8]) -> io::Result<()> { + let len: u32 = payload + .len() + .try_into() + .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "payload too large"))?; + writer.write_all(&len.to_be_bytes()).await?; + writer.write_all(payload).await +} + +/// Read one u32-length-prefixed encoded payload, bounded like the remote +/// frame reader. +async fn read_payload(reader: &mut R) -> io::Result> { + let mut len_buf = [0u8; 4]; + reader.read_exact(&mut len_buf).await?; + let len = u32::from_be_bytes(len_buf); + if len as usize > MAX_DESKTOP_PAYLOAD { + return Err(invalid()); + } + let mut payload = vec![0u8; len as usize]; + reader.read_exact(&mut payload).await?; + Ok(payload) +} + +/// Agent side: pump the pinned remote session's encoded frames and control +/// events down the socket, and viewer controls up, until either side ends. +/// Caller EOF, viewer `Finished`, or remote-session end all close the +/// channel; returning drops the session and aborts its remote legs. +pub(super) async fn serve( + stream: &mut UnixStream, + mut session: rds_desktop::client::DesktopSession, +) -> io::Result<()> { + // A relay must be opened with `relay_encoded`; otherwise there is + // nothing to forward — fail rather than park the channel forever. + let ctrl = session.control_sender(); + if session.encoded.is_none() { + return Err(invalid()); + } + let (mut reader, mut writer) = stream.split(); + + 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), + } + } + }; + // Borrows `session`'s receivers for the pump's lifetime; when `serve` + // returns the session drops and aborts its remote legs. + let down = async { + let encoded = session.encoded.as_mut().unwrap(); + let events = &mut session.events; + let mut open = 2u8; + while open > 0 { + tokio::select! { + frame = encoded.recv() => match frame { + Some(f) => { + if f.payload.len() > MAX_DESKTOP_PAYLOAD { + return Err(invalid()); + } + write_frame(&mut writer, &DesktopDown::Frame { header: f.header }).await?; + write_payload(&mut writer, &f.payload).await?; + } + None => open -= 1, + }, + event = events.recv() => match event { + Some(e) => write_frame(&mut writer, &DesktopDown::Event(e)).await?, + None => open -= 1, + }, + } + } + write_frame(&mut writer, &DesktopDown::Finished).await + }; + tokio::select! { + result = up => result, + result = down => result, + } +} + +/// One relayed encoded frame as the viewer receives it. +#[derive(Debug)] +pub struct RelayedFrame { + pub header: FrameHeader, + pub payload: Vec, +} + +/// One message off a managed desktop channel. +#[derive(Debug)] +pub enum ManagedMessage { + /// One encoded frame — feed a `rds_desktop::client::RelayDecoder`. + Frame(RelayedFrame), + /// A control-plane event (input acks, heartbeat echoes). + Event(rds_core::DesktopEvent), +} + +/// 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. +struct ControlState { + writer: Mutex, + display: u32, + input_seq: AtomicU64, + heartbeat_seq: AtomicU64, + /// Base for locally-stamped event timestamps (the viewer's own clock). + started: Instant, +} + +impl ControlState { + async fn write_control(&self, control: DesktopControl) -> io::Result<()> { + let mut writer = self.writer.lock().await; + write_frame(&mut *writer, &DesktopUp::Control(control)).await + } +} + +/// Cloneable control half of a managed desktop channel — the analogue of +/// `DesktopSession::control_sender`. Input tasks hold one while the owner +/// keeps receiving frames. +#[derive(Clone)] +pub struct ManagedControl { + state: Arc, +} + +impl ManagedControl { + /// Send a fully-formed control message verbatim — callers that manage + /// their own sequencing use this; everyone else prefers the typed + /// helpers. + pub async fn control(&self, control: DesktopControl) -> io::Result<()> { + self.state.write_control(control).await + } + + /// Queue one input event; sequence, timestamp and target display are + /// filled in from channel state, mirroring `DesktopSession::send_input`. + pub async fn send_input(&self, kind: InputKind) -> io::Result { + let seq = self + .state + .input_seq + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |s| s.checked_add(1)) + .map_err(|_| io::Error::other("input sequence exhausted"))?; + self.control(DesktopControl::Input(InputEvent { + seq, + event_ts_ms: self.state.started.elapsed().as_millis() as u64, + display_id: self.state.display, + kind, + })) + .await?; + Ok(seq) + } + + /// Liveness probe; the remote echoes it as `DesktopEvent::Heartbeat`. + pub async fn heartbeat(&self) -> io::Result { + let seq = self + .state + .heartbeat_seq + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |s| s.checked_add(1)) + .map_err(|_| io::Error::other("heartbeat sequence exhausted"))?; + self.control(DesktopControl::Heartbeat { + seq, + ts_ms: self.state.started.elapsed().as_millis() as u64, + }) + .await?; + Ok(seq) + } + + /// Ask the remote encoder for a fresh keyframe. + pub async fn request_idr(&self) -> io::Result<()> { + self.control(DesktopControl::RequestIdr).await + } + + /// Request an encoder bitrate in bits per second. + pub async fn set_bitrate(&self, bps: u32) -> io::Result<()> { + self.control(DesktopControl::SetBitrate(bps)).await + } +} + +/// Client handle on a managed desktop channel. Owns the authenticated IPC +/// socket; `None` from `recv` means the remote session ended cleanly or the +/// manager went away. +pub struct ManagedDesktop { + reader: OwnedReadHalf, + state: Arc, + /// The manager session this channel is pinned to. + pub session: SessionId, + /// Negotiated capabilities the remote agent reported. + pub caps: DesktopCaps, +} + +impl std::fmt::Debug for ManagedDesktop { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ManagedDesktop") + .field("session", &self.session) + .field("caps", &self.caps) + .finish_non_exhaustive() + } +} + +impl ManagedDesktop { + pub(super) fn new( + stream: UnixStream, + session: SessionId, + caps: DesktopCaps, + display: u32, + ) -> Self { + let (reader, writer) = stream.into_split(); + Self { + reader, + state: Arc::new(ControlState { + writer: Mutex::new(writer), + display, + input_seq: AtomicU64::new(0), + heartbeat_seq: AtomicU64::new(0), + started: Instant::now(), + }), + session, + caps, + } + } + + /// A cloneable control half for tasks that send while `recv` runs — + /// mirroring `DesktopSession::control_sender`. + pub fn control_handle(&self) -> ManagedControl { + ManagedControl { + state: Arc::clone(&self.state), + } + } + + /// Next manager→viewer message; `None` after the channel ends. EOF + /// between messages means the manager went away — also `None` — while + /// EOF inside a frame payload stays an error. + pub async fn recv(&mut self) -> io::Result> { + match read_frame::<_, DesktopDown>(&mut self.reader).await { + Ok(DesktopDown::Frame { header }) => { + let payload = read_payload(&mut self.reader).await?; + Ok(Some(ManagedMessage::Frame(RelayedFrame { + header, + payload, + }))) + } + Ok(DesktopDown::Event(event)) => Ok(Some(ManagedMessage::Event(event))), + Ok(DesktopDown::Finished) => Ok(None), + Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => Ok(None), + Err(e) => Err(e), + } + } + + /// Shorthand for `control_handle().send_input(...)`. + pub async fn send_input(&self, kind: InputKind) -> io::Result { + self.control_handle().send_input(kind).await + } + + /// Shorthand for `control_handle().heartbeat()`. + pub async fn heartbeat(&self) -> io::Result { + self.control_handle().heartbeat().await + } + + /// Shorthand for `control_handle().request_idr()`. + pub async fn request_idr(&self) -> io::Result<()> { + self.control_handle().request_idr().await + } + + /// Shorthand for `control_handle().set_bitrate(...)`. + pub async fn set_bitrate(&self, bps: u32) -> io::Result<()> { + self.control_handle().set_bitrate(bps).await + } + + /// Announce a clean end to the manager; dropping without it is also a + /// clean end (caller EOF). A `ManagedControl` that outlives this handle + /// can still buffer writes until the manager closes its side, then + /// sees an error. + pub async fn finish(self) -> io::Result<()> { + let mut writer = self.state.writer.lock().await; + write_frame(&mut *writer, &DesktopUp::Finished).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn header(seq: u64) -> FrameHeader { + FrameHeader { + seq, + capture_ts_ms: 1, + encode_done_ts_ms: 2, + send_ts_ms: 3, + keyframe: seq == 0, + codec: rds_core::Codec::H264, + width: 640, + height: 480, + } + } + + fn pair() -> (ManagedDesktop, UnixStream) { + let (viewer, peer) = UnixStream::pair().unwrap(); + let caps = DesktopCaps { + displays: vec![], + codecs: vec![rds_core::Codec::H264], + }; + ( + ManagedDesktop::new(viewer, SessionId([0x11; 16]), caps, 0), + peer, + ) + } + + #[tokio::test] + async fn payload_round_trip_and_bounds() { + let (mut a, mut b) = UnixStream::pair().unwrap(); + write_payload(&mut a, b"encoded-payload").await.unwrap(); + assert_eq!(read_payload(&mut b).await.unwrap(), b"encoded-payload"); + write_payload(&mut a, &[]).await.unwrap(); + assert!(read_payload(&mut b).await.unwrap().is_empty()); + // A length prefix above the cap fails before the body is read. + a.write_u32((MAX_DESKTOP_PAYLOAD + 1) as u32).await.unwrap(); + assert_eq!( + read_payload(&mut b).await.unwrap_err().kind(), + io::ErrorKind::InvalidData + ); + // A truncated payload is an error, not a partial frame. + a.write_u32(8).await.unwrap(); + a.write_all(b"abc").await.unwrap(); + drop(a); + assert!(read_payload(&mut b).await.is_err()); + } + + #[tokio::test] + async fn recv_reads_frame_event_finished_and_eof() { + let (mut channel, mut peer) = pair(); + let frame = header(0); + write_frame(&mut peer, &DesktopDown::Frame { header: frame }) + .await + .unwrap(); + write_payload(&mut peer, b"h264").await.unwrap(); + write_frame( + &mut peer, + &DesktopDown::Event(rds_core::DesktopEvent::Heartbeat { seq: 4, ts_ms: 9 }), + ) + .await + .unwrap(); + write_frame(&mut peer, &DesktopDown::Finished) + .await + .unwrap(); + match channel.recv().await.unwrap() { + Some(ManagedMessage::Frame(f)) => { + assert_eq!(f.header.seq, 0); + assert_eq!(f.payload, b"h264"); + } + other => panic!("expected frame, got {other:?}"), + } + match channel.recv().await.unwrap() { + Some(ManagedMessage::Event(rds_core::DesktopEvent::Heartbeat { seq, .. })) => { + assert_eq!(seq, 4) + } + other => panic!("expected event, got {other:?}"), + } + assert!(channel.recv().await.unwrap().is_none()); + // After Finished the peer is done; a bare EOF is also a clean end. + let (mut channel, peer) = pair(); + drop(peer); + assert!(channel.recv().await.unwrap().is_none()); + } + + #[tokio::test] + async fn recv_propagates_truncated_payload() { + let (mut channel, mut peer) = pair(); + write_frame(&mut peer, &DesktopDown::Frame { header: header(1) }) + .await + .unwrap(); + peer.write_u32(16).await.unwrap(); + peer.write_all(b"short").await.unwrap(); + drop(peer); + // EOF inside a payload is corruption, not a clean end. + assert!(channel.recv().await.is_err()); + } + + #[tokio::test] + async fn finish_and_controls_reach_the_manager() { + let (channel, mut peer) = pair(); + let control = channel.control_handle(); + control.request_idr().await.unwrap(); + control.set_bitrate(2_000_000).await.unwrap(); + match read_frame::<_, DesktopUp>(&mut peer).await.unwrap() { + DesktopUp::Control(DesktopControl::RequestIdr) => {} + other => panic!("expected RequestIdr, got {other:?}"), + } + match read_frame::<_, DesktopUp>(&mut peer).await.unwrap() { + DesktopUp::Control(DesktopControl::SetBitrate(bps)) => { + assert_eq!(bps, 2_000_000) + } + other => panic!("expected SetBitrate, got {other:?}"), + } + channel.finish().await.unwrap(); + match read_frame::<_, DesktopUp>(&mut peer).await.unwrap() { + DesktopUp::Finished => {} + other => panic!("expected Finished, got {other:?}"), + } + // A control handle outliving the channel errors once the peer is + // gone; a live peer would still accept buffered writes. + drop(peer); + assert!(control.request_idr().await.is_err()); + } + + #[tokio::test] + async fn control_serializes_concurrent_senders() { + let (channel, mut peer) = pair(); + let a = channel.control_handle(); + let b = channel.control_handle(); + let first = tokio::spawn(async move { + for _ in 0..20 { + a.request_idr().await.unwrap(); + } + }); + for _ in 0..20 { + b.request_idr().await.unwrap(); + } + first.await.unwrap(); + // Interleaved writers must still yield 40 well-formed frames. + for _ in 0..40 { + match read_frame::<_, DesktopUp>(&mut peer).await.unwrap() { + DesktopUp::Control(DesktopControl::RequestIdr) => {} + other => panic!("expected RequestIdr, got {other:?}"), + } + } + } + + /// Synthetic input sink — `serve_desktop_with` probes the host's real + /// backend when `input_sink` is `None`, which does not exist headless. + struct NoopInput; + impl rds_desktop::InputSink for NoopInput { + fn inject(&mut self, _: &InputEvent) -> Result<(), rds_desktop::DesktopError> { + Ok(()) + } + } + + /// Loopback QUIC endpoints plus a synthetic-frame serving side — + /// the same shape `session_v2` uses, small enough for a unit file. + /// The endpoints stay owned by the returned guard: an endpoint drop + /// closes its connections, so the session dies with it otherwise. + struct Relay { + session: rds_desktop::client::DesktopSession, + server_task: tokio::task::JoinHandle<()>, + _server_ep: rds_net::Endpoint, + _client_ep: rds_net::Endpoint, + } + + async fn loopback_relay() -> Relay { + use rds_core::{HelloAck, StreamHello, UniHello}; + use rds_desktop::{SessionConfig, SyntheticProducer, serve_desktop_with}; + + let config = || rds_net::EndpointConfig { + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + discovery: false, + ..Default::default() + }; + let server_ep = rds_net::bind_endpoint(config()).await.unwrap(); + let client_ep = rds_net::bind_endpoint(config()).await.unwrap(); + let target = server_ep.addr(); + let server_ep_guard = server_ep.clone(); + let server_task = tokio::spawn(async move { + let conn = server_ep.accept().await.unwrap().await.unwrap(); + let (mut send, mut recv) = conn.accept_bi().await.unwrap(); + let (hello, route) = match read_frame::<_, StreamHello>(&mut recv).await.unwrap() { + StreamHello::DesktopV2 { session, hello } => { + (hello, Some(UniHello::DesktopFrames { id: session })) + } + StreamHello::Desktop(hello) => (hello, None), + other => panic!("unexpected hello {other:?}"), + }; + write_frame( + &mut send, + &HelloAck::Desktop(DesktopCaps { + displays: vec![], + codecs: vec![rds_core::Codec::H264], + }), + ) + .await + .unwrap(); + serve_desktop_with( + conn, + send, + recv, + hello, + SessionConfig { + input_sink: Some(Box::new(NoopInput)), + producer: Some(Box::new( + SyntheticProducer::new(30, 640, 480, 1024).keyframe_every(5), + )), + frame_route: route, + ..Default::default() + }, + ) + .await + .unwrap(); + }); + let conn = client_ep.connect(target, rds_core::ALPN).await.unwrap(); + let session = rds_desktop::client::DesktopSession::connect_opts( + &conn, + rds_core::DesktopHello { + display: 0, + max_fps: 30, + codec: rds_core::Codec::H264, + input_acks: false, + }, + rds_desktop::client::SessionOpts { + session: Some(rand::random()), + relay_encoded: true, + ..Default::default() + }, + ) + .await + .unwrap(); + Relay { + session, + server_task, + _server_ep: server_ep_guard, + _client_ep: client_ep, + } + } + + /// Read one down-message plus its raw payload when it is a frame. + async fn recv_down(viewer: &mut UnixStream) -> io::Result { + match read_frame::<_, DesktopDown>(viewer).await? { + down @ DesktopDown::Frame { .. } => { + read_payload(viewer).await?; + Ok(down) + } + down => Ok(down), + } + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn loopback_relay_publishes_encoded() { + let mut relay = loopback_relay().await; + let encoded = relay.session.encoded.as_mut().expect("relay tap"); + let delivery = tokio::time::timeout(std::time::Duration::from_secs(10), encoded.recv()) + .await + .expect("no encoded frame") + .expect("encoded tap closed"); + assert!(!delivery.payload.is_empty()); + // Frame headers publish on the shared tap in relay mode too. + let header = tokio::time::timeout( + std::time::Duration::from_secs(5), + relay.session.frame_headers.recv(), + ) + .await + .expect("no frame header") + .expect("header tap closed"); + assert_eq!(header.seq, delivery.header.seq); + relay.server_task.abort(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn serve_relays_frames_echoes_controls_and_finishes() { + let relay = loopback_relay().await; + let (mut ipc, mut viewer) = UnixStream::pair().unwrap(); + let pump = tokio::spawn(async move { serve(&mut ipc, relay.session).await }); + // Encoded frames arrive as postcard header + raw payload. + for _ in 0..8 { + match recv_down(&mut viewer).await.unwrap() { + DesktopDown::Frame { .. } => {} + other => panic!("expected frame, got {other:?}"), + } + } + // A viewer control crosses the pump into the remote session and + // its heartbeat echo comes back down. + write_frame( + &mut viewer, + &DesktopUp::Control(DesktopControl::Heartbeat { seq: 77, ts_ms: 5 }), + ) + .await + .unwrap(); + loop { + match recv_down(&mut viewer).await.unwrap() { + DesktopDown::Event(rds_core::DesktopEvent::Heartbeat { seq, .. }) => { + assert_eq!(seq, 77); + break; + } + DesktopDown::Frame { .. } => {} + other => panic!("unexpected {other:?}"), + } + } + write_frame(&mut viewer, &DesktopUp::Finished) + .await + .unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), pump) + .await + .expect("pump did not stop on Finished") + .unwrap() + .unwrap(); + relay.server_task.abort(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn serve_ends_when_remote_session_drops() { + let relay = loopback_relay().await; + let (mut ipc, mut viewer) = UnixStream::pair().unwrap(); + let pump = tokio::spawn(async move { serve(&mut ipc, relay.session).await }); + // Read a frame so the pump is mid-stream when the remote dies. + match recv_down(&mut viewer).await.unwrap() { + DesktopDown::Frame { .. } => {} + other => panic!("expected frame, got {other:?}"), + } + relay.server_task.abort(); + // Remote conn closure ends both taps; the viewer sees Finished. + loop { + match recv_down(&mut viewer).await.unwrap() { + DesktopDown::Finished => break, + DesktopDown::Frame { .. } | DesktopDown::Event(_) => {} + } + } + tokio::time::timeout(std::time::Duration::from_secs(5), pump) + .await + .expect("pump did not stop on remote end") + .unwrap() + .unwrap(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn serve_ends_on_caller_eof() { + let relay = loopback_relay().await; + let (mut ipc, viewer) = UnixStream::pair().unwrap(); + let pump = tokio::spawn(async move { serve(&mut ipc, relay.session).await }); + drop(viewer); + tokio::time::timeout(std::time::Duration::from_secs(5), pump) + .await + .expect("pump did not stop on caller EOF") + .unwrap() + .unwrap(); + relay.server_task.abort(); + } +} diff --git a/crates/rds-client/src/local/mod.rs b/crates/rds-client/src/local/mod.rs index 51c67a4..dcd7b32 100644 --- a/crates/rds-client/src/local/mod.rs +++ b/crates/rds-client/src/local/mod.rs @@ -218,10 +218,15 @@ async fn run( result } +mod desktop; + +pub use desktop::{ManagedControl, ManagedDesktop, ManagedMessage, RelayedFrame}; + struct Output { reply: Reply, reservation: Option, tcp: Option<(RequestStreams, OwnedSemaphorePermit)>, + desktop: Option<(rds_desktop::client::DesktopSession, OwnedSemaphorePermit)>, } impl Output { @@ -230,6 +235,7 @@ impl Output { reply, reservation: None, tcp: None, + desktop: None, } } } @@ -280,6 +286,9 @@ 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?; + } } Ok(()) } @@ -371,6 +380,7 @@ async fn execute( reply: Reply::Connected(reservation.id), reservation: Some(reservation), tcp: None, + desktop: None, }) } Command::Renew { session, grant } => { @@ -393,6 +403,7 @@ async fn execute( reply: Reply::Done, reservation: Some(reservation), tcp: None, + desktop: None, }) } Command::Select { session } => { @@ -438,6 +449,37 @@ async fn execute( reply: Reply::Opened(session), reservation: None, tcp: Some((RequestStreams::new(pair), permit)), + desktop: None, + }) + } + Command::Desktop { session, hello } => { + let permit = streams + .try_acquire_owned() + .map_err(|_| ErrorCode::Capacity)?; + let (session, conn) = state::lock(shared)?.connection(session)?; + // Relay mode publishes encoded frames for IPC forwarding; the + // viewer decodes, so the manager never needs a codec. A fresh + // v2 session ID keeps the frame route isolated from any other + // desktop session on this connection. + let remote = rds_desktop::client::DesktopSession::connect_opts( + &conn, + *hello, + rds_desktop::client::SessionOpts { + session: Some(rand::random()), + relay_encoded: true, + ..Default::default() + }, + ) + .await + .map_err(|_| ErrorCode::Remote)?; + Ok(Output { + reply: Reply::DesktopOpened { + session, + caps: remote.caps().clone(), + }, + reservation: None, + tcp: None, + desktop: Some((remote, permit)), }) } } @@ -457,7 +499,7 @@ impl Client { } async fn exchange(&self, command: Command) -> Result<(Reply, UnixStream), Error> { - let body = matches!(command, Command::OpenTcp { .. }); + let body = matches!(command, Command::OpenTcp { .. } | Command::Desktop { .. }); let timeout = if matches!(command, Command::Sync { .. }) { SYNC_TIMEOUT + Duration::from_secs(10) } else { @@ -493,12 +535,35 @@ impl Client { } pub async fn request(&self, command: Command) -> Result { - if matches!(command, Command::OpenTcp { .. }) { + if matches!(command, Command::OpenTcp { .. } | Command::Desktop { .. }) { return Err(Error::Protocol); } self.exchange(command).await.map(|(reply, _)| reply) } + /// Open a managed desktop channel: the manager runs the remote session + /// and relays encoded frames; the caller decodes and sends controls. + /// The returned socket speaks `DesktopDown`/`DesktopUp`. + pub async fn desktop( + &self, + session: Option, + hello: rds_core::DesktopHello, + ) -> Result { + let display = hello.display; + let (reply, stream) = self + .exchange(Command::Desktop { + session, + hello: Box::new(hello), + }) + .await?; + match reply { + Reply::DesktopOpened { session, caps } => { + Ok(ManagedDesktop::new(stream, session, caps, display)) + } + _ => Err(Error::Protocol), + } + } + 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 460b9a9..30dddb7 100644 --- a/crates/rds-core/src/local.rs +++ b/crates/rds-core/src/local.rs @@ -1,9 +1,14 @@ //! Versioned local control protocol. This is not a remote service ALPN. use serde::{Deserialize, Serialize}; -pub const VERSION: u16 = 4; +/// Wire version 5 adds the managed desktop channel (`Desktop`, +/// `DesktopDown`/`DesktopUp`). +pub const VERSION: u16 = 5; pub const MAX_SESSIONS: usize = 32; pub const TCP_CHUNK: usize = 16 * 1024; +/// Bound on one encoded desktop frame payload on the local channel — the +/// same 32 MiB the remote frame reader enforces. +pub const MAX_DESKTOP_PAYLOAD: usize = 32 * 1024 * 1024; /// Local TCP bodies distinguish byte-direction FIN from caller departure. /// Never Debug: application bytes can contain secrets. @@ -13,6 +18,34 @@ pub enum TcpFrame { Finish, } +/// Manager→viewer messages on a [`Reply::DesktopOpened`] channel: encoded +/// frames exactly as the remote agent sent them (the viewer owns decoding), +/// control-plane events, and an explicit clean end marker. +/// +/// A `Frame` postcard message is followed on the wire by a big-endian u32 +/// length and that many encoded payload bytes (≤ `MAX_DESKTOP_PAYLOAD`) — +/// the same header-then-body shape the remote frame stream uses, since one +/// postcard message cannot hold a 32 MiB payload. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum DesktopDown { + /// One encoded frame past the remote sequence/stale-drop discipline. + Frame { + header: crate::FrameHeader, + }, + Event(crate::DesktopEvent), + /// The remote session ended; the socket closes after this frame. + Finished, +} + +/// Viewer→manager messages on a desktop channel. `Control` is forwarded +/// verbatim to the remote session's control stream; `Finished` ends the +/// relay without error. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum DesktopUp { + Control(crate::DesktopControl), + Finished, +} + /// Random, process-lifetime handle. Never aliases a session after agent restart. #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] #[serde(try_from = "String", into = "String")] @@ -106,6 +139,13 @@ pub enum Command { session: SessionId, operation: SyncOperation, }, + /// Agent-mediated desktop viewer session. A `DesktopOpened` reply is + /// followed by a bidirectional [`DesktopDown`]/[`DesktopUp`] channel on + /// the same socket until either side ends it. + Desktop { + session: Option, + hello: Box, + }, } /// Never Debug: paths may contain private information. @@ -164,6 +204,12 @@ pub enum Reply { info: crate::AgentInfo, }, Opened(SessionId), + /// Desktop relay established: negotiated caps, then the socket becomes a + /// [`DesktopDown`]/[`DesktopUp`] channel for the returned session. + DesktopOpened { + session: SessionId, + caps: crate::DesktopCaps, + }, Ticket(String), Synced { session: SessionId, @@ -248,12 +294,96 @@ mod tests { } } + #[test] + fn desktop_channel_round_trips() { + // Wire v5: managed desktop command, reply and body enums. + assert_eq!(VERSION, 5); + let id = SessionId([0x5a; 16]); + let request = Request { + version: VERSION, + command: Command::Desktop { + session: Some(id), + hello: Box::new(crate::DesktopHello { + display: 1, + max_fps: 30, + codec: crate::Codec::H264, + input_acks: true, + }), + }, + }; + let bytes = postcard::to_stdvec(&request).unwrap(); + match postcard::from_bytes::(&bytes).unwrap().command { + Command::Desktop { session, hello } => { + assert_eq!(session, Some(id)); + assert_eq!(hello.display, 1); + assert!(hello.input_acks); + } + _ => panic!("wrong variant"), + } + let response = Response { + version: VERSION, + result: Ok(Reply::DesktopOpened { + session: id, + caps: crate::DesktopCaps { + displays: vec![crate::DisplayInfo { + index: 0, + width: 1920, + height: 1080, + primary: true, + }], + codecs: vec![crate::Codec::H264], + }, + }), + }; + let bytes = postcard::to_stdvec(&response).unwrap(); + match postcard::from_bytes::(&bytes) + .unwrap() + .result + .unwrap() + { + Reply::DesktopOpened { session, caps } => { + assert_eq!(session, id); + assert_eq!(caps.displays.len(), 1); + } + _ => panic!("wrong variant"), + } + let header = crate::FrameHeader { + seq: 7, + capture_ts_ms: 1, + encode_done_ts_ms: 2, + send_ts_ms: 3, + keyframe: true, + codec: crate::Codec::H264, + width: 640, + height: 480, + }; + for down in [ + DesktopDown::Frame { header }, + DesktopDown::Event(crate::DesktopEvent::Heartbeat { seq: 9, ts_ms: 4 }), + DesktopDown::Finished, + ] { + let bytes = postcard::to_stdvec(&down).unwrap(); + assert!(bytes.len() < crate::MAX_MESSAGE_LEN as usize); + postcard::from_bytes::(&bytes).unwrap(); + } + for up in [ + DesktopUp::Control(crate::DesktopControl::RequestIdr), + DesktopUp::Control(crate::DesktopControl::Heartbeat { seq: 3, ts_ms: 11 }), + DesktopUp::Finished, + ] { + let bytes = postcard::to_stdvec(&up).unwrap(); + 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 ".*") { let _ = postcard::from_bytes::(&bytes); let _ = postcard::from_bytes::(&bytes); let _ = postcard::from_bytes::(&bytes); + let _ = postcard::from_bytes::(&bytes); + let _ = postcard::from_bytes::(&bytes); let _ = text.parse::(); } } diff --git a/crates/rds-desktop/src/client.rs b/crates/rds-desktop/src/client.rs index fff80b8..5cf6aed 100644 --- a/crates/rds-desktop/src/client.rs +++ b/crates/rds-desktop/src/client.rs @@ -168,6 +168,8 @@ pub struct DesktopSession { pub frame_headers: mailbox::Receiver, /// Server→client control events (input acks, heartbeat echoes). pub events: mailbox::Receiver, + /// Encoded wire frames in relay mode; `None` on a direct session. + pub encoded: Option>, /// Send input or encoder control to the serving side. ctrl_tx: mpsc::Sender, /// Next expected frame sequence — the lowest seq still accepted. @@ -234,6 +236,21 @@ pub struct SessionOpts { /// an ended session can never reach this session's inbox. `None` /// keeps the legacy shared `Desktop` route for old peers. pub session: Option<[u8; 16]>, + /// Relay mode (the local session manager): publish each encoded + /// payload to [`DesktopSession::encoded`] instead of decoding it — + /// no decoder is created and `frames` stays empty. Transport-level + /// resync (IDR on sequence gaps) still runs; decode-chain discipline + /// is the downstream viewer's job. + pub relay_encoded: bool, +} + +/// One encoded frame exactly as it arrived on the wire, published in relay +/// mode for consumers that forward rather than decode. +#[derive(Debug)] +pub struct EncodedDelivery { + pub header: FrameHeader, + /// Annex-B encoded payload, already bounded by the receive budget. + pub payload: bytes::Bytes, } impl DesktopSession { @@ -307,6 +324,13 @@ impl DesktopSession { let (header_tx, frame_headers) = mailbox::channel(64); let (ctrl_tx, mut ctrl_rx) = mpsc::channel::(64); let (events_tx, events) = mailbox::channel::(128); + let (encoded_tx, encoded) = match opts.relay_encoded { + true => { + let (tx, rx) = mailbox::channel::(4); + (Some(tx), Some(rx)) + } + false => (None, None), + }; let next_seq = Arc::new(AtomicU64::new(0)); // `u64::MAX` = "not measured": a loopback heartbeat can // legitimately round-trip in 0 ms, so 0 cannot be the sentinel. @@ -352,6 +376,7 @@ impl DesktopSession { ReceiveContext { frame_tx, header_tx, + encoded_tx, next_seq: next_seq.clone(), ctrl: ctrl_tx.clone(), clock: clock.clone(), @@ -369,6 +394,7 @@ impl DesktopSession { frames, frame_headers, events, + encoded, ctrl_tx, next_seq, input_seq: AtomicU64::new(0), @@ -456,6 +482,88 @@ impl DesktopSession { .await .map_err(|_| DesktopError::Input("control channel closed".into())) } + + /// Queue one fully-formed control message verbatim — callers that + /// manage their own sequencing use this; everyone else prefers the + /// typed helpers. + pub async fn send_control(&self, control: DesktopControl) -> Result<(), DesktopError> { + self.ctrl_tx + .send(control) + .await + .map_err(|_| DesktopError::Input("control channel closed".into())) + } + + /// Cloneable control-queue handle — the managed relay forwards a + /// viewer's verbatim `DesktopControl` messages through it while the + /// session's receivers are consumed elsewhere. Direct callers use the + /// typed helpers (`send_input`, `heartbeat`, `request_idr`, + /// `set_bitrate`) that fill sequence metadata in. + pub fn control_sender(&self) -> mpsc::Sender { + self.ctrl_tx.clone() + } +} + +/// What [`RelayDecoder::push`] made of one relayed encoded frame. +pub enum RelayOutcome { + /// A frame came out of the decoder (decoder-enabled builds only). + #[cfg(feature = "x11")] + Frame(RawFrame), + /// Consumed without producing a frame — the codec is still buffering + /// inputs, or this build has no decoder at all. + Pending, + /// The reference chain is broken; the viewer should send + /// `DesktopControl::RequestIdr` upstream. + NeedIdr, +} + +/// Viewer-side decode chain for a relayed frame channel +/// (`local::DesktopDown::Frame`): applies the same wait-for-keyframe and +/// broken-chain discipline a direct session does, while the resync request +/// stays with the caller's own control path. `NeedIdr` reports are rate +/// limited like the in-session request path, so one report covers a run +/// of broken frames. +pub struct RelayDecoder { + delivery: Delivery, + /// Wall clock of the last `NeedIdr` report. + last_idr_report: Option, +} + +impl RelayDecoder { + pub fn new() -> Self { + Self { + delivery: Delivery::new(), + last_idr_report: None, + } + } + + /// Feed one relayed encoded frame. Input must already have passed the + /// relay's sequence checks — this owns only decode-chain state. + pub fn push(&mut self, header: &FrameHeader, payload: Vec) -> RelayOutcome { + match self.delivery.decode(header, payload) { + #[cfg(feature = "x11")] + DecodeOutcome::Decoded(raw) => RelayOutcome::Frame(raw), + DecodeOutcome::Buffered => RelayOutcome::Pending, + DecodeOutcome::Failed => { + self.delivery.invalidate(); + let now = std::time::Instant::now(); + let due = self.last_idr_report.is_none_or(|last| { + now.duration_since(last).as_millis() as u64 >= IDR_MIN_INTERVAL_MS + }); + if due { + self.last_idr_report = Some(now); + RelayOutcome::NeedIdr + } else { + RelayOutcome::Pending + } + } + } + } +} + +impl Default for RelayDecoder { + fn default() -> Self { + Self::new() + } } fn next_control_seq(sequence: &AtomicU64) -> Result { @@ -469,6 +577,8 @@ fn next_control_seq(sequence: &AtomicU64) -> Result { struct ReceiveContext { frame_tx: mailbox::Sender, header_tx: mailbox::Sender, + /// Relay-mode tap: encoded payloads publish here instead of decoding. + encoded_tx: Option>, next_seq: Arc, ctrl: mpsc::Sender, clock: SessionClock, @@ -494,6 +604,15 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) { // read_one rejects u64::MAX before publishing anything. ctx.next_seq.store(header.seq + 1, Ordering::Relaxed); ctx.header_tx.send(header.clone()); + if let Some(tx) = &ctx.encoded_tx { + // Relay mode: forward the payload still encoded; the + // downstream viewer owns decode-chain discipline. + tx.send(EncodedDelivery { + header, + payload: body.into(), + }); + continue; + } let Ok(slot) = DECODE_SLOTS.acquire().await else { break; }; // No control/queue/connection handles escape into native work. // An already running call may finish after cancellation, but diff --git a/crates/rds-desktop/tests/session_v2.rs b/crates/rds-desktop/tests/session_v2.rs index fcce5c9..a04c5a3 100644 --- a/crates/rds-desktop/tests/session_v2.rs +++ b/crates/rds-desktop/tests/session_v2.rs @@ -224,6 +224,7 @@ async fn harness_with_input( SessionOpts { clock: Some(clock.clone()), session: Some(next_session_id()), + ..Default::default() }, ) .await @@ -654,6 +655,7 @@ async fn desktop_and_sync_share_one_connection() { SessionOpts { clock: Some(clock.clone()), session: Some(next_session_id()), + ..Default::default() }, ) .await @@ -779,6 +781,7 @@ async fn stale_and_foreign_frame_routes_never_reach_the_session_inbox() { SessionOpts { clock: Some(clock.clone()), session: Some(session_id), + ..Default::default() }, ) .await @@ -848,6 +851,7 @@ async fn legacy_shared_route_still_serves_v1_clients() { SessionOpts { clock: Some(clock.clone()), session: None, + ..Default::default() }, ) .await @@ -874,3 +878,76 @@ fn self_rss_kb() -> Option { fn self_rss_kb() -> Option { None } + +/// Relay mode (the local session manager): encoded payloads publish to +/// `session.encoded` verbatim and never touch the local decode chain; +/// sequence discipline and headers still run. A viewer-side +/// `RelayDecoder` then rebuilds the chain and reports `NeedIdr` on a +/// broken one — the caller forwards it over its own control path. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn relay_mode_publishes_encoded_frames_and_viewer_decodes() { + use rds_desktop::client::{RelayDecoder, RelayOutcome}; + + let fps = 30; + let frame_bytes = 32 * 1024; + let clock = SessionClock::default(); + let (server_ep, client_ep, _impair, target) = endpoints(None).await; + let server_task = spawn_serving( + server_ep, + fps, + frame_bytes, + 5, + clock.clone(), + TestInput::default(), + ) + .await; + let conn = client_ep.connect(target, rds_core::ALPN).await.unwrap(); + let mut session = DesktopSession::connect_opts( + &conn, + DesktopHello { + display: 0, + max_fps: fps, + codec: Codec::H264, + input_acks: false, + }, + SessionOpts { + clock: Some(clock.clone()), + session: Some(next_session_id()), + relay_encoded: true, + }, + ) + .await + .unwrap(); + let mut encoded = session + .encoded + .take() + .expect("relay mode publishes an encoded tap"); + let mut decoder = RelayDecoder::new(); + let mut last_seq = None; + let mut outcomes = 0u32; + for _ in 0..20 { + let delivery = tokio::time::timeout(Duration::from_secs(10), encoded.recv()) + .await + .unwrap() + .unwrap(); + // The relay still enforces in-order monotonic headers. + if let Some(prev) = last_seq { + assert!(delivery.header.seq > prev, "out-of-order relayed seq"); + } + last_seq = Some(delivery.header.seq); + assert!(!delivery.payload.is_empty()); + match decoder.push(&delivery.header, delivery.payload.to_vec()) { + RelayOutcome::Pending | RelayOutcome::NeedIdr => outcomes += 1, + #[cfg(feature = "x11")] + RelayOutcome::Frame(raw) => { + assert_eq!(raw.width, delivery.header.width); + outcomes += 1; + } + } + } + assert!(outcomes > 0); + // Relay mode never publishes decoded frames locally. + assert!(session.frames.try_recv().is_none()); + drop(session); + server_task.abort(); +} From 30fb78f8e171003e784048943ab01cd8b5fbeced Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 11:59:34 +0500 Subject: [PATCH 2/2] docs: managed desktop channel contract and W2.4 ledger local-sessions.md documents the v5 managed desktop body channel (header-then-payload shape, Finished markers, permit accounting, manager/viewer ownership split and --direct), the capability matrix gains service:desktop-managed, and the plan/progress ledgers record the implemented viewer API with installed-binary migration left open. --- README.md | 2 +- docs/capability-matrix.md | 7 ++++- docs/conventions.md | 2 +- docs/grant-leases.md | 2 +- docs/implementation-plan.md | 7 +++-- docs/local-sessions.md | 59 ++++++++++++++++++++++++++++++------ docs/remediation-plan.md | 12 +++++++- docs/remediation-progress.md | 44 +++++++++++++++++++++++++++ 8 files changed, 117 insertions(+), 18 deletions(-) diff --git a/README.md b/README.md index 3737973..8d9d3a9 100644 --- a/README.md +++ b/README.md @@ -121,7 +121,7 @@ agent and name CLI persist their policy revisions and freshness state. See the issuer or deploying managed access. [Grant v2 and renewable leases](docs/grant-leases.md) bind access to the controlled endpoint. `rds session renew --session --grant-file ` extends a live same-scope connection; automatic GDS issuance -is still pending. Old grants must be reissued, and CLI/agent IPC v4 upgraded together. +is still pending. Old grants must be reissued, and CLI/agent IPC v5 upgraded together. The agent's role, data-plane services, authority posture and deadlines can also live in one versioned JSON document (`--agent-config`), with diff --git a/docs/capability-matrix.md b/docs/capability-matrix.md index a8998dc..231498c 100644 --- a/docs/capability-matrix.md +++ b/docs/capability-matrix.md @@ -31,7 +31,8 @@ only what is actually served. `ping`/`info` are the always-on control plane. | service:ping | implemented | peer on `--allow` list | `rds-net` tests, CI | | service:info | implemented | peer on `--allow` list | `rds-cli` tests, CI | | service:tcp | implemented | enabled + `--allow` + grant scope for tcp | `rds-ssh` e2e, CI | -| service:desktop | experimental | enabled + `desktop` build flag; usable capture backend; grant scope | capture→encode→decode→stats contract tests; no viewer | +| service:desktop | experimental | enabled + `desktop` build flag; usable capture backend; grant scope | capture→encode→decode→stats contract tests | +| service:desktop-managed | implemented | manager-owned remote session relays encoded frames over local IPC v5; viewer decodes via `RelayDecoder` (`desktop` build flag for real decode) | relay e2e + managed-channel tests; `rds desktop` defaults to it, `--direct` kept | | service:sync | implemented | enabled + `--sync-dir` configured | `rds-sync` tests, resumable transfer tests | | service:audio | stub | no codec; wire shape reserved in v2 | `ServiceKind::Audio` variant only | | service:sync-read | implemented | grant scope on sync root; optional grant `sync_paths` subtree | grant/scope tests | @@ -106,6 +107,10 @@ Notes: still fails closed when no capture backend probes usable. X11 init failure is an explicit `Err`, and the x11 CI lane fails when the feature cannot initialize. +- `service:desktop-managed` is the local-IPC path: the manager forwards + encoded frames and controls without decoding, so it works on headless + manager builds; the *remote* still needs a capture backend, and the + viewer build needs `rds-desktop/x11` for real decode. - Historical `transfer` reports predate receiver-verified timing; the verified lane is renamed `transfer-receiver-ack-v1` on `fix/verified-transfer-benchmark` — update this row when it merges diff --git a/docs/conventions.md b/docs/conventions.md index e44da1c..f492a0c 100644 --- a/docs/conventions.md +++ b/docs/conventions.md @@ -55,7 +55,7 @@ Rules every change follows. CI enforces what it can; the rest is review. `rds-ssh`; do not duplicate them as RDS control messages. - Transport protocol selection uses ALPN (`rds/0`, `rds-relay/0`). Signed objects and local IPC also carry independent explicit versions: [grant v3 - and IPC v4](grant-leases.md) require a coordinated upgrade without weakening + and IPC v5](grant-leases.md) require a coordinated upgrade without weakening authorization — v2 grants still verify, but a deployment that pins the v3 tenant/policy claims refuses them. General remote capability negotiation remains W2.2. diff --git a/docs/grant-leases.md b/docs/grant-leases.md index 98f56db..0511447 100644 --- a/docs/grant-leases.md +++ b/docs/grant-leases.md @@ -165,7 +165,7 @@ reserved state; extra reply bytes are rejected. A lost response after successful publication remains an uncertain RPC outcome, inspectable through `session list`. Neither another device nor a remote command is selected/replayed automatically. -Local IPC is now **version 4** following [managed file transfers](local-sessions.md). Upgrade CLI and local agent together; incompatible +Local IPC is now **version 5** following the [managed desktop channel](local-sessions.md). Upgrade CLI and local agent together; incompatible versions fail before session mutation. Remote ALPN `rds/0` and existing service framing are unchanged. The signed grant is independently versioned and the new renewal variant is appended; old servers reject unsupported renewal without a diff --git a/docs/implementation-plan.md b/docs/implementation-plan.md index 79d5835..af47188 100644 --- a/docs/implementation-plan.md +++ b/docs/implementation-plan.md @@ -23,9 +23,10 @@ below was verified against vendored sources or upstream documents on ## 0. Doctrine The current product-facing increment is the W2.4 -[local session manager](local-sessions.md). Default CLI connectivity and local -key-inode runtime ownership are implemented. Remaining manager work covers -viewer APIs and installed-device qualification; single-file send/receive now +[local session manager](local-sessions.md). Default CLI connectivity, local +key-inode runtime ownership and the managed desktop viewer API (local wire +v5) are implemented. Remaining manager work is installed-device +qualification and coordinated binary migration; single-file send/receive use the manager. The [native SSH client](ssh.md) now covers standard shell/exec/PTY requests; host-key/account enrollment, broker/reattachment and real-network/platform qualification remain W5 work, diff --git a/docs/local-sessions.md b/docs/local-sessions.md index eb41f6f..74d4966 100644 --- a/docs/local-sessions.md +++ b/docs/local-sessions.md @@ -68,10 +68,13 @@ for this boundary; no package/version or external helper was added. Compatibility: `rds --direct --key-file ...` explicitly binds an independent endpoint and acquires the same exclusive key owner as the agent. An occupied key is an error. Direct mode requires a persisted key path; there is -no accidental ephemeral-identity fallback. `desktop` still requires this explicit -mode because its manager API remains unimplemented. Send/receive default to the manager. -The local wire version is now **4** (adds file transfers); upgrade CLI -and agent together. Old/new local versions fail without mutating session state. +no accidental ephemeral-identity fallback. +Send/receive default to the manager, and `desktop` now uses it too — the +manager owns the remote session and relays encoded frames while the CLI +decodes locally. `desktop --direct` remains for native/direct use. +The local wire version is now **5** (adds the managed desktop channel); +upgrade CLI and agent together. Old/new local versions fail without +mutating session state. Remote ALPN is unchanged, with an appended isolated sync greeting; signed grant v3 accepts v2 payloads but a pinned tenant/policy binding refuses them. See [renewal contract](grant-leases.md). @@ -153,6 +156,7 @@ broker. | IPC workers | 96; acceptance pauses at capacity | | Long-lived TCP streams | 64; leaves control worker space | | File transfers | 8 agent-wide; 1 per outgoing session; immediate capacity/busy refusal | +| Managed desktop channels | one shared stream-permit slot (of 64) per channel; encoded payloads ≤ 32 MiB each, forwarded in arrival order | | Request prelude / reply write | 5 seconds each | | Resolve + dial + Authz + operation | 45 seconds total | | Client exchange, including connect/response and control EOF | 55 seconds for short operations | @@ -163,7 +167,7 @@ broker. | CLI forwarding workers | positive 16-bit limit; default 64 | Wire types live in `rds-core::local`, with explicit request/response versions. -Local IPC is **version 4**; upgrade CLI and agent together. Remote ALPN stays +Local IPC is **version 5**; upgrade CLI and agent together. Remote ALPN stays `rds/0`, with the appended `SyncTransfer` greeting described below. There are no unbounded task/command queues; the OS backlog is separate. TCP bodies do not inherit the prelude deadline and use the existing reset/stop cancellation guard. Operations never automatically retry a @@ -223,14 +227,49 @@ end-to-end receipt after the CLI disappears. See the [managed-transfer receipt](reports/rds-managed-sync-20260926.md). +## Managed desktop channels + +`rds desktop ` (and `rds session desktop --session `) default to the +manager like send/recv. `Command::Desktop` opens a remote `DesktopV2` session +on the pinned connection in **relay mode** (`SessionOpts::relay_encoded`); the +reply is `Reply::DesktopOpened` with the negotiated session and caps, then the +socket switches to a bidirectional body channel: + +- Down, `DesktopDown::Frame { header }` + a big-endian u32 length + raw + encoded payload (≤ 32 MiB); `DesktopDown::Event(DesktopEvent)` for control + events; `DesktopDown::Finished` as the explicit clean end marker. +- Up, `DesktopUp::Control(DesktopControl)` verbatim into the remote session's + control sender; `DesktopUp::Finished` or caller EOF ends the channel. + +Encoded payloads travel beside postcard framing deliberately: a +`DesktopDown::Frame` postcard message stays under the 64 KiB control bound +while frame bodies reach 32 MiB. The manager never decodes or links a codec — +it republishes the remote session's encoded tap (sequence checks and +headers still run there), so a headless manager build can serve desktop. +The viewer owns decode: `rds_desktop::client::RelayDecoder` reapplies +wait-for-keyframe/broken-chain discipline on `RelayedFrame` payloads and +reports `NeedIdr`, rate-limited like the in-session path; the caller +forwards it with `request_idr` over its own control channel. Input +sequence numbers and heartbeat timestamps are stamped viewer-side +(`ManagedControl`, cloneable, serialized behind one writer half so +concurrent senders cannot interleave postcard bytes). + +The channel ends on viewer `Finished`/EOF, remote-session end (both taps +close) or a body error; the manager then drops the remote +`DesktopSession`, aborting its tasks and releasing the stream permit. +`Client::desktop` shares the connect/response deadline; the body itself +has no IPC timeout. `--direct` still opens a `DesktopSession` in-process +for native/direct desktop use, unchanged. + ## Remaining sequence and exit checks -1. **W2.4 migration:** viewer manager APIs and coordinated installed-binary - migration remain. CLI connectivity defaults and cooperative same-key-inode +1. **W2.4 migration:** the viewer manager API is implemented (managed + desktop channel, local wire v5); coordinated installed-binary migration + remains. CLI connectivity defaults and cooperative same-key-inode runtime ownership are implemented. Qualify native macOS credentials, - actual distinct-user rejection, relay-registration reuse, - FD/RSS budgets and manager service APIs for media. Single-file send/receive - are implemented; directory/two-way synchronization remains W8. + actual distinct-user rejection, relay-registration reuse and + FD/RSS budgets. Single-file send/receive are implemented; + directory/two-way synchronization remains W8. 2. **W5 SSH:** native client, explicit host pins, credential selection, OS PTY requests and terminal restoration are implemented. Complete GDS host/account provisioning, broker isolation and authorized reattachment; qualify native diff --git a/docs/remediation-plan.md b/docs/remediation-plan.md index 25f3ad5..40e6394 100644 --- a/docs/remediation-plan.md +++ b/docs/remediation-plan.md @@ -237,7 +237,17 @@ reuse the agent identity and pin session handles. New transfer-specific uni tags isolate delayed/canceled data, with bounded live routes and coordinated IPC v4 migration. [The receipt](reports/rds-managed-sync-20260926.md) covers both transport backends, preserved TCP on cancellation, grants and real CLI processes. -Viewer APIs, broader negotiation, native/installed qualification and directory + +W2.4 implementation update (2026-09-27): the [managed desktop channel](local-sessions.md#managed-desktop-channels) +lands the viewer-manager API on local wire v5 — `Command::Desktop`/`Reply::DesktopOpened` +open a relay-mode `DesktopV2` session on the pinned connection, then the socket +speaks `DesktopDown`/`DesktopUp`: postcard headers plus u32-length-prefixed raw +encoded payloads (≤ 32 MiB, beside the 64 KiB control bound). The manager +forwards without decoding (no codec linkage); the viewer owns decode via +`RelayDecoder` with in-session keyframe/broken-chain discipline and +rate-limited `NeedIdr`. `rds desktop` defaults to the managed path and +`session desktop` exists; `--direct` keeps native in-process sessions. +Broader negotiation, native/installed qualification and directory synchronization remain open; this increment does not close `r2-session-core`. W2.3 implementation update (2026-09-25): [grant v2 and explicit renewal](grant-leases.md) diff --git a/docs/remediation-progress.md b/docs/remediation-progress.md index 7f9af2c..b50339d 100644 --- a/docs/remediation-progress.md +++ b/docs/remediation-progress.md @@ -2004,3 +2004,47 @@ scope while sibling, traversal and write escapes refuse by name. Remaining W2.3: account-level scopes (OS identity), automatic GDS issuance/renewal and policy reconciliation — cross-repository, deferred. + +## Managed desktop viewer API — local wire v5 (2026-09-27) + +W2.4's viewer-manager API piece. The local wire moves to **version 5**: +`Command::Desktop { session, hello }` opens a desktop channel on the +pinned managed session; `Reply::DesktopOpened { session, caps }` answers +it, then the authenticated socket becomes a bidirectional body channel +(`DesktopDown`/`DesktopUp`). Frame bodies travel beside postcard framing — +a `Frame { header }` message followed by a big-endian u32 length and the +raw encoded payload (≤ `MAX_DESKTOP_PAYLOAD` = 32 MiB) — because a +postcard control message is bound to 64 KiB while encoded frames are not. + +The manager runs the remote `DesktopSession` in the new **relay mode** +(`SessionOpts::relay_encoded`): sequence checks and header publication +still run session-side, while encoded payloads publish to a bounded tap +instead of decoding — the manager never links a codec and stays buildable +headless. `DesktopSession::control_sender` exposes a cloneable control +queue and `send_control` forwards verbatim. The pump ends on viewer +`Finished`/EOF, remote-session end, or body error; dropping the session +aborts its remote legs and releases the stream permit it shares with TCP +bodies. + +The viewer side (`rds_client::local::ManagedDesktop`) owns the +authenticated socket: `recv` yields `ManagedMessage::{Frame,Event}` until +`Finished`/EOF (`None`), and a cloneable `ManagedControl` serializes +verbatim controls plus typed `send_input`/`heartbeat`/`request_idr`/ +`set_bitrate` behind one write half — concurrent senders cannot interleave +postcard bytes. `rds_desktop::client::RelayDecoder` reapplies +wait-for-keyframe and broken-chain discipline viewer-side, returning +`Frame`/`Pending`/`NeedIdr` with the same 500 ms resync rate limit the +in-session path uses. `rds desktop` defaults to the managed channel and +`rds session desktop` exists; `desktop --direct` keeps native in-process +sessions unchanged. + +Tests: local wire v5 round-trips and proptest decode-safety in rds-core; +framing bounds, truncation, Finished/EOF and concurrent-sender +serialization unit tests plus three real-loopback `serve` e2e tests +(frames+heartbeat echo+clean finish, remote drop, caller EOF) in +rds-client; a relay-mode transport e2e in `session_v2`; and a real-agent +managed-desktop test asserting clean `Rejected(Remote)` refusal with no +stream-permit leak on peers that cannot serve desktop. + +Remaining W2.4: coordinated installed-binary migration and native/installed +qualification — cross-platform, deferred.