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
40 changes: 32 additions & 8 deletions crates/rds-cli/src/desktop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,9 @@ pub async fn read_grant(
#[cfg(feature = "desktop")]
mod control;

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

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

Expand Down Expand Up @@ -417,6 +420,7 @@ mod native {
view.status("Waiting for screen");
let control = channel.control_handle();
let (progress, last_frame) = tokio::sync::watch::channel(tokio::time::Instant::now());
let control_progress = liveness::ControlWatchdog::default();
// 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.
Expand Down Expand Up @@ -462,10 +466,16 @@ 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::Heartbeat { seq, ts_ms }) => {
if control_progress.echoed(seq, 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 { id, bytes }) => {
control_progress.clipboard_ready(id, bytes);
view.clipboard_ack(id, bytes);
}
None => {
Expand All @@ -481,9 +491,16 @@ mod native {
result = media => result,
}
};
let controls = control::pump(input, &control, &last_frame, started, |message| {
view.input_sent(message);
});
let controls = control::pump_with_liveness(
input,
&control,
&last_frame,
&control_progress,
started,
|message| {
view.input_sent(message);
},
);
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
Expand Down Expand Up @@ -529,6 +546,7 @@ mod native {
let ctrl = session.control_sender();
let mut last_frame = tokio::time::Instant::now();
let mut watchdog = control::VideoWatchdog::default();
let control_progress = liveness::ControlWatchdog::default();
let mut tick = tokio::time::interval(Duration::from_secs(1));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let result = loop {
Expand All @@ -537,11 +555,13 @@ mod native {
message = input.recv() => match message {
Some(ViewerInput::Control(message)) => {
view.input_sent(&message);
control_progress.sent(&message);
tokio::time::timeout(Duration::from_secs(2),ctrl.send(message)).await.map_err(|_|anyhow::anyhow!("direct desktop control stalled"))??;
},
Some(ViewerInput::Close)|None => break Ok(true),
},
_ = tick.tick() => {
control_progress.check()?;
match watchdog.observe(last_frame) {
control::VideoAction::Reconnect => anyhow::bail!("remote video stopped making progress"),
control::VideoAction::Repair => {
Expand All @@ -550,12 +570,16 @@ mod native {
},
control::VideoAction::Healthy => {},
}
tokio::time::timeout(Duration::from_secs(2),ctrl.send(rds_core::DesktopControl::Heartbeat { seq: 0,ts_ms: started.elapsed().as_millis() as u64 })).await.map_err(|_|anyhow::anyhow!("direct desktop heartbeat stalled"))??;
let heartbeat = control_progress.heartbeat(started.elapsed().as_millis() as u64)?;
control_progress.sent(&heartbeat);
tokio::time::timeout(Duration::from_secs(2),ctrl.send(heartbeat)).await.map_err(|_|anyhow::anyhow!("direct desktop heartbeat stalled"))??;
},
event = session.events.recv() => match event {
Some(rds_core::DesktopEvent::Heartbeat { ts_ms,.. }) => view.control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)),
Some(rds_core::DesktopEvent::Heartbeat { seq, ts_ms }) => {
if control_progress.echoed(seq, 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 { id, bytes }) => {view.clipboard_ack(id, bytes);},
Some(rds_core::DesktopEvent::ClipboardReady { id, bytes }) => {control_progress.clipboard_ready(id, bytes);view.clipboard_ack(id, bytes);},
None => break Ok(false),
},
frame = session.frames.recv() => match frame {
Expand Down
150 changes: 147 additions & 3 deletions crates/rds-cli/src/desktop/control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,10 +79,44 @@ async fn send_control(
Ok(())
}

pub(super) async fn pump(
#[cfg(test)]
async fn pump(
input: &mut impl Input,
sender: &impl Sender,
progress: &watch::Receiver<tokio::time::Instant>,
started: Instant,
sent: impl Fn(&DesktopControl),
) -> anyhow::Result<bool> {
pump_inner(input, sender, progress, None, started, sent).await
}

pub(super) async fn pump_with_liveness(
input: &mut impl Input,
sender: &impl Sender,
progress: &watch::Receiver<tokio::time::Instant>,
control_progress: &super::liveness::ControlWatchdog,
started: Instant,
sent: impl Fn(&DesktopControl),
) -> anyhow::Result<bool> {
pump_inner(
input,
sender,
progress,
Some(control_progress),
started,
|message| {
control_progress.sent(message);
sent(message);
},
)
.await
}

async fn pump_inner(
input: &mut impl Input,
sender: &impl Sender,
progress: &watch::Receiver<tokio::time::Instant>,
control_progress: Option<&super::liveness::ControlWatchdog>,
started: Instant,
sent: impl Fn(&DesktopControl),
) -> anyhow::Result<bool> {
Expand All @@ -93,6 +127,7 @@ pub(super) async fn pump(
let message = tokio::select! {
biased;
_ = tick.tick() => {
if let Some(monitor) = control_progress { monitor.check()?; }
let progress = *progress.borrow();
match watchdog.observe(progress) {
VideoAction::Reconnect => {
Expand All @@ -107,8 +142,10 @@ pub(super) async fn pump(
},
VideoAction::Healthy => {},
}
DesktopControl::Heartbeat {
seq: 0, ts_ms: started.elapsed().as_millis() as u64,
let ts_ms = started.elapsed().as_millis() as u64;
match control_progress {
Some(monitor) => monitor.heartbeat(ts_ms)?,
None => DesktopControl::Heartbeat { seq:0, ts_ms },
}
},
message = input.recv() => match message {
Expand Down Expand Up @@ -169,6 +206,113 @@ mod tests {
}
}

#[tokio::test(start_paused = true)]
async fn fresh_video_cannot_mask_unacknowledged_control_and_cancellation_drops_media() {
let (wire, mut remote) = tokio::io::duplex(256);
let (tx, mut input) = mpsc::channel(8);
let (progress, last_frame) = watch::channel(tokio::time::Instant::now());
let stop = CancellationToken::new();
let dropped = Arc::new(AtomicBool::new(false));
let held = dropped.clone();
let mut task = tokio::spawn(async move {
let monitor = super::super::liveness::ControlWatchdog::default();
let media = async {
let _drop = Dropped(held);
std::future::pending().await
};
run(
pump_with_liveness(
&mut input,
&Wire(Mutex::new(wire)),
&last_frame,
&monitor,
Instant::now(),
|_| {},
),
media,
&stop,
)
.await
});
tx.send(ViewerInput::Control(DesktopControl::Input(InputEvent {
seq: 37,
event_ts_ms: 0,
display_id: 0,
kind: InputKind::KeyDown { code: 56 },
})))
.await
.unwrap();
let started = tokio::time::Instant::now();
let mut inputs = 0;
let error = loop {
tokio::select! {
result = &mut task => break result.unwrap().unwrap_err(),
message = rds_net::read_frame::<_,DesktopUp>(&mut remote) => {
let message = match message {
Ok(DesktopUp::Control(message)) => message,
Err(error) => {
assert_eq!(error.kind(),std::io::ErrorKind::UnexpectedEof);
break (&mut task).await.unwrap().unwrap_err();
},
_ => panic!("unexpected framing"),
};
match message {
DesktopControl::Heartbeat {..} => { progress.send_replace(tokio::time::Instant::now()); },
DesktopControl::Input(event) => { assert_eq!(event.seq,37);inputs+=1; },
DesktopControl::RequestIdr => panic!("fresh video does not need repair"),
_ => panic!("unexpected control"),
}
}
}
};
assert!(
error
.to_string()
.contains("control stopped making progress")
);
assert_eq!(started.elapsed(), Duration::from_secs(8));
assert_eq!(inputs, 1, "no input replay while stalled");
assert!(dropped.load(Ordering::SeqCst));
}

#[tokio::test(start_paused = true)]
async fn matched_control_echoes_preserve_the_healthy_session_during_continuous_media() {
let (wire, mut remote) = tokio::io::duplex(256);
let (tx, mut input) = mpsc::channel(8);
let (progress, last_frame) = watch::channel(tokio::time::Instant::now());
let monitor = super::super::liveness::ControlWatchdog::default();
let sending = monitor.clone();
let task = tokio::spawn(async move {
pump_with_liveness(
&mut input,
&Wire(Mutex::new(wire)),
&last_frame,
&sending,
Instant::now(),
|_| {},
)
.await
});
let started = tokio::time::Instant::now();
let mut previous = None;
while started.elapsed() < Duration::from_secs(60) {
let DesktopUp::Control(DesktopControl::Heartbeat { seq, ts_ms }) =
rds_net::read_frame(&mut remote).await.unwrap()
else {
panic!("unexpected heartbeat framing");
};
if let Some(n) = previous {
assert_eq!(seq, n + 1);
}
previous = Some(seq);
assert!(monitor.echoed(seq, ts_ms));
progress.send_replace(tokio::time::Instant::now());
assert!(!task.is_finished());
}
tx.send(ViewerInput::Close).await.unwrap();
assert!(task.await.unwrap().unwrap());
}

#[tokio::test(start_paused = true)]
async fn video_watchdog_repairs_before_reconnect_and_rearms_only_on_frames() {
let first = tokio::time::Instant::now();
Expand Down
Loading
Loading