Skip to content
Closed
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
95 changes: 67 additions & 28 deletions crates/rds-cli/src/desktop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ pub async fn read_grant(
}
}

#[cfg(feature = "desktop")]
mod control;

#[cfg(feature = "desktop")]
mod native {
use super::*;
Expand Down Expand Up @@ -263,8 +266,22 @@ mod native {
failures = 0;
}
failures = failures.saturating_add(1);
let snapshot = view.snapshot();
tracing::warn!(
error = ?result.as_ref().err(), attempt = failures,
network_stage = %snapshot.network_stage,
render_stage = %snapshot.render_stage,
decoded_frame_age_ms = snapshot.decoded_frame_age_ms,
submission_age_ms = snapshot.submission_age_ms,
ui_event_age_ms = snapshot.ui_event_age_ms,
pending_input_acks = snapshot.report.pending_input_acks,
oldest_input_ack_age_ms = snapshot.report.oldest_input_ack_age_ms,
control_rtt_ms = snapshot.report.control_rtt_ms,
occluded = snapshot.occluded,
"desktop reconnecting"
);
view.status("Reconnecting");
tracing::warn!(error = ?result.as_ref().err(), attempt=failures,"desktop reconnecting");
view.stage("reconnecting");
let wait = Duration::from_millis((250u64 << failures.min(5)).min(8000));
if !retry_pause(wait, &stop, input).await {
return Ok(());
Expand Down Expand Up @@ -339,46 +356,68 @@ mod native {
extent(view, &channel.caps, options.display)?;
view.status("Waiting for screen");
let control = channel.control_handle();
let mut decoder = RelayDecoder::new();
let mut last_frame = Instant::now();
let mut tick = tokio::time::interval(Duration::from_secs(1));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let result = loop {
tokio::select! {
_ = stop.cancelled() => break Ok(true),
message = input.recv() => match message {
Some(ViewerInput::Control(message)) => {
view.input_sent(&message);
tokio::time::timeout(Duration::from_secs(2), control.control(message)).await.map_err(|_|anyhow::anyhow!("desktop control write stalled"))??;
},
Some(ViewerInput::Close)|None => break Ok(true),
},
_ = tick.tick() => {
anyhow::ensure!(last_frame.elapsed() < Duration::from_secs(15),"remote video stopped making progress");
tokio::time::timeout(Duration::from_secs(2), control.control(rds_core::DesktopControl::Heartbeat { seq: 0,ts_ms: started.elapsed().as_millis() as u64 })).await.map_err(|_|anyhow::anyhow!("desktop heartbeat write stalled"))??;
},
message = channel.recv() => match message? {
let (progress, last_frame) = tokio::sync::watch::channel(tokio::time::Instant::now());
// Keep the entire receive/decode future alive while controls progress.
// Selecting individual recv calls and awaiting decode in their handler
// prevents input, heartbeat and close from being polled during decode.
let media = async {
let mut decoder = RelayDecoder::new();
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(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::Frame(frame)) => {
view.stage("decoding");
let received = Instant::now();
view.media_timing(&frame.header);
let (next,outcome) = tokio::time::timeout(Duration::from_secs(5), decoder.push_bounded(frame.header,frame.payload)).await.map_err(|_|anyhow::anyhow!("desktop decode stalled"))??;
let (next, outcome) = tokio::time::timeout(
Duration::from_secs(5),
decoder.push_bounded(frame.header, frame.payload),
)
.await
.map_err(|_| anyhow::anyhow!("desktop decode stalled"))??;
decoder = next;
match outcome {
RelayOutcome::Frame(raw) => { last_frame = Instant::now(); view.frame(raw,received); },
RelayOutcome::NeedIdr => control.request_idr().await?,
RelayOutcome::Pending => {},
RelayOutcome::Frame(raw) => {
progress.send_replace(tokio::time::Instant::now());
view.frame(raw, received);
}
RelayOutcome::NeedIdr => {
tokio::time::timeout(Duration::from_secs(2), control.request_idr())
.await
.map_err(|_| {
anyhow::anyhow!("desktop repair write stalled")
})??
}
RelayOutcome::Pending => {}
}
view.stage("receiving");
}
}
}
};
let _ = tokio::time::timeout(Duration::from_secs(1), channel.finish()).await;
let controls = control::pump(input, &control, &last_frame, started, |message| {
view.input_sent(message);
});
let result = control::run(controls, media, 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.
drop(channel);
result
}

Expand Down
Loading
Loading