Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <id>
--grant-file <path>` 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
Expand Down
52 changes: 52 additions & 0 deletions crates/rds-agent/tests/local_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
}
101 changes: 96 additions & 5 deletions crates/rds-cli/src/managed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SessionId>,
#[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<()> {
Expand Down Expand Up @@ -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}"),
Expand Down Expand Up @@ -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?;
Expand Down Expand Up @@ -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(())
Expand Down Expand Up @@ -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")
}
1 change: 1 addition & 0 deletions crates/rds-client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading