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
93 changes: 87 additions & 6 deletions crates/rds-desktop/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -123,14 +123,16 @@ impl Delivery {
self.waiting_keyframe = true;
}

fn request_idr(&mut self, ctrl: &mpsc::Sender<DesktopControl>, now_ms: u64) {
fn request_idr(&mut self, ctrl: &mpsc::Sender<DesktopControl>, now_ms: u64) -> bool {
if self
.last_idr_req_ms
.is_none_or(|last| now_ms.saturating_sub(last) >= IDR_MIN_INTERVAL_MS)
&& ctrl.try_send(DesktopControl::RequestIdr).is_ok()
{
self.last_idr_req_ms = Some(now_ms);
return true;
}
false
}

fn decode(&mut self, header: &FrameHeader, body: Vec<u8>) -> DecodeOutcome {
Expand Down Expand Up @@ -683,6 +685,38 @@ struct ReceiveContext {
control_rtt_ms: Arc<AtomicU64>,
}

/// Successors of one missing reference must not queue repeated large IDRs
/// while a reliable recovery request/key is still in transit. Actual reader
/// failure or a received key ends that episode; absent either, retry within
/// the existing key-reader deadline. Failed control admission never arms it.
#[derive(Default)]
struct GapRepair {
requested_at_ms: Option<u64>,
}
impl GapRepair {
fn request(
&mut self,
delivery: &mut Delivery,
ctrl: &mpsc::Sender<DesktopControl>,
now_ms: u64,
) -> bool {
if self
.requested_at_ms
.is_some_and(|at| now_ms.saturating_sub(at) < KEYFRAME_READ_TIMEOUT.as_millis() as u64)
{
return false;
}
if delivery.request_idr(ctrl, now_ms) {
self.requested_at_ms = Some(now_ms);
return true;
}
false
}
fn clear(&mut self) {
self.requested_at_ms = None;
}
}

enum FrameRead {
Complete(FrameHeader, Vec<u8>, FrameBudget),
Rejected,
Expand All @@ -698,6 +732,7 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
// Keep the decoded queue open for the session even in a headless build.
let _frames = &ctx.frame_tx;
let mut delivery = Delivery::new();
let mut gap_repair = GapRepair::default();
let mut admitted = 0u64;
let mut completed = 0u64;
let mut rejected = 0u64;
Expand All @@ -718,6 +753,7 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
ctx.header_tx.send(header.clone());
completed += 1;
if header.keyframe {
gap_repair.clear();
received_keys += 1;
}
if let Some(tx) = &ctx.encoded_tx {
Expand Down Expand Up @@ -773,9 +809,10 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
let wait = reorder_wait(ctx.control_rtt_ms.load(Ordering::Relaxed));
if readers.is_empty() && ordered.expire(wait) {
gaps+=1;
tracing::warn!(expected_seq=ctx.next_seq.load(Ordering::Relaxed),reorder_budget_ms=wait.as_millis(),"desktop reference gap expired");
delivery.invalidate();
delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());
let now_ms = ctx.clock.now_ms();
let request_sent = gap_repair.request(&mut delivery, &ctx.ctrl, now_ms);
tracing::warn!(expected_seq=ctx.next_seq.load(Ordering::Relaxed),reorder_budget_ms=wait.as_millis(),request_sent,pending_recovery_age_ms=gap_repair.requested_at_ms.map(|at| now_ms.saturating_sub(at)),"desktop reference gap expired");
}
if health.elapsed()>=std::time::Duration::from_secs(5) {
tracing::info!(admitted,completed,rejected,timed_out,gaps,received_keys,in_flight=ctx.receiving.load(Ordering::Relaxed),"desktop receiver health");
Expand All @@ -785,8 +822,8 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
frame = readers.join_next(), if !readers.is_empty() => {
match frame {
Some(Ok(FrameRead::Complete(header,body,budget)))=>ordered.push(header,(body,budget)),
Some(Ok(FrameRead::TimedOut))=>{timed_out+=1;delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());},
_=>{rejected+=1;delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());},
Some(Ok(FrameRead::TimedOut))=>{timed_out+=1;gap_repair.clear();delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());},
_=>{rejected+=1;gap_repair.clear();delivery.request_idr(&ctx.ctrl,ctx.clock.now_ms());},
}
}
// A compressed relay consumer can backpressure completed readers.
Expand Down Expand Up @@ -832,6 +869,7 @@ async fn read_one(
};
let frame_seq = header.seq;
let keyframe = header.keyframe;
let header_ms = started.elapsed().as_millis();
let deadline = if keyframe {
KEYFRAME_READ_TIMEOUT
} else {
Expand Down Expand Up @@ -859,7 +897,19 @@ async fn read_one(
Some((header, body, budget))
};
match tokio::time::timeout(deadline.saturating_sub(started.elapsed()), reading).await {
Ok(Some((header, body, budget))) => FrameRead::Complete(header, body, budget),
Ok(Some((header, body, budget))) => {
if started.elapsed() >= std::time::Duration::from_millis(250) {
tracing::warn!(
frame_seq,
keyframe,
body_bytes,
header_ms,
elapsed_ms = started.elapsed().as_millis(),
"desktop frame receive completed slowly"
);
}
FrameRead::Complete(header, body, budget)
}
Ok(None) => {
tracing::debug!(
?frame_seq,
Expand Down Expand Up @@ -972,6 +1022,37 @@ mod tests {
assert!(matches!(rx.try_recv(), Ok(DesktopControl::RequestIdr)));
}

#[test]
fn gap_repair_coalesces_until_a_key_failure_or_bounded_retry() {
let (tx, mut rx) = mpsc::channel(2);
let mut delivery = Delivery::new();
let mut repair = GapRepair::default();
assert!(repair.request(&mut delivery, &tx, 0));
assert!(matches!(rx.try_recv(), Ok(DesktopControl::RequestIdr)));
for now in [500, 1000, 7999] {
assert!(!repair.request(&mut delivery, &tx, now));
}
assert!(rx.try_recv().is_err());
assert!(repair.request(&mut delivery, &tx, 8000));
assert!(matches!(rx.try_recv(), Ok(DesktopControl::RequestIdr)));
repair.clear();
assert!(repair.request(&mut delivery, &tx, 8500));
assert!(matches!(rx.try_recv(), Ok(DesktopControl::RequestIdr)));
}

#[test]
fn a_full_control_queue_does_not_arm_gap_repair() {
let (tx, mut rx) = mpsc::channel(1);
tx.try_send(DesktopControl::RequestIdr).unwrap();
let mut delivery = Delivery::new();
let mut repair = GapRepair::default();
assert!(!repair.request(&mut delivery, &tx, 0));
assert!(repair.requested_at_ms.is_none());
rx.try_recv().unwrap();
assert!(repair.request(&mut delivery, &tx, 0));
assert!(matches!(rx.try_recv(), Ok(DesktopControl::RequestIdr)));
}

#[test]
fn decoded_dimensions_are_bounded_before_bgra_allocation() {
assert_eq!(crate::frame_bytes(7680, 4320), Some(7680 * 4320 * 4));
Expand Down
18 changes: 13 additions & 5 deletions crates/rds-desktop/src/render/viewer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -297,11 +297,19 @@ impl ViewerHandle {
}
pub fn input_sent(&self, control: &DesktopControl) {
if let DesktopControl::Input(event) = control {
lock(&self.state).input_latency.sent(
event.seq,
event.event_ts_ms,
self.started.elapsed().as_millis() as u64,
);
let sent_ms = self.started.elapsed().as_millis() as u64;
lock(&self.state)
.input_latency
.sent(event.seq, event.event_ts_ms, sent_ms);
let event_class = match event.kind {
InputKind::KeyDown { .. } => "key_down",
InputKind::KeyUp { .. } => "key_up",
InputKind::PointerMove { .. } | InputKind::PointerMotion { .. } => "pointer_move",
InputKind::PointerButton { pressed: true, .. } => "button_down",
InputKind::PointerButton { pressed: false, .. } => "button_up",
InputKind::Scroll { .. } => "scroll",
};
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 input_ack(&self, seq: u64) {
Expand Down
153 changes: 153 additions & 0 deletions crates/rds-desktop/tests/recovery_requests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
//! One missing reference must not queue another recovery while its key is pending.
use rds_core::{
Codec, DesktopCaps, DesktopControl, DesktopHello, FrameHeader, HelloAck, StreamHello, UniHello,
};
use rds_desktop::client::{DesktopSession, SessionOpts};
use rds_net::{Backend, Connection, EndpointConfig, read_frame, write_frame};
use std::time::Duration;

async fn frame(conn: &Connection, id: [u8; 16], seq: u64, keyframe: bool) {
let mut stream = conn.open_uni().await.unwrap();
write_frame(&mut stream, &UniHello::DesktopFrames { id })
.await
.unwrap();
write_frame(
&mut stream,
&FrameHeader {
seq,
keyframe,
capture_ts_ms: 0,
encode_done_ts_ms: 0,
send_ts_ms: 0,
codec: Codec::H264,
width: 32,
height: 32,
},
)
.await
.unwrap();
stream
.write_all(b"synthetic encoded relay payload")
.await
.unwrap();
stream.finish().unwrap();
}

#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn repeated_gap_waits_for_one_recovery_but_a_later_gap_requests_another() {
for backend in [Backend::Iroh, Backend::Noq] {
let config = EndpointConfig {
backend,
discovery: false,
bind_addrs: vec!["127.0.0.1:0".parse().unwrap()],
..Default::default()
};
let client_ep = rds_net::bind_endpoint(config.clone()).await.unwrap();
let server_ep = rds_net::bind_endpoint(config).await.unwrap();
let (a, b) = tokio::join!(client_ep.connect(server_ep.addr(), rds_core::ALPN), async {
server_ep.accept().await.unwrap().await
});
let (a, b) = (a.unwrap(), b.unwrap());
let (client, server) = tokio::join!(
DesktopSession::connect_opts(
&a,
DesktopHello {
display: 0,
max_fps: 30,
codec: Codec::H264,
input_acks: false
},
SessionOpts {
session: Some(rand::random()),
relay_encoded: true,
..Default::default()
}
),
async {
let (mut send, mut recv) = b.accept_bi().await.unwrap();
let StreamHello::DesktopV2 { session, .. } = read_frame(&mut recv).await.unwrap()
else {
panic!("expected tagged session")
};
write_frame(
&mut send,
&HelloAck::Desktop(DesktopCaps {
displays: vec![],
codecs: vec![Codec::H264],
}),
)
.await
.unwrap();
(send, recv, session)
}
);
let mut client = client.unwrap();
let (_control_send, mut control_recv, id) = server;
let (control_tx, mut controls) = tokio::sync::mpsc::channel(16);
let reader = tokio::spawn(async move {
while let Ok(message) = read_frame::<_, DesktopControl>(&mut control_recv).await {
if control_tx.send(message).await.is_err() {
break;
}
}
});
frame(&b, id, 0, true).await;
assert_eq!(
tokio::time::timeout(
Duration::from_secs(2),
client.encoded.as_mut().unwrap().recv()
)
.await
.unwrap()
.unwrap()
.header
.seq,
0
);
frame(&b, id, 2, false).await; // reference 1 never arrives
assert!(matches!(
tokio::time::timeout(Duration::from_secs(2), controls.recv())
.await
.unwrap(),
Some(DesktopControl::RequestIdr)
));
tokio::time::sleep(Duration::from_millis(600)).await;
frame(&b, id, 3, false).await; // another successor of the same missing reference
assert!(
tokio::time::timeout(Duration::from_millis(800), controls.recv())
.await
.is_err(),
"same recovery episode queued a second expensive keyframe"
);
frame(&b, id, 4, true).await;
assert_eq!(
tokio::time::timeout(
Duration::from_secs(2),
client.encoded.as_mut().unwrap().recv()
)
.await
.unwrap()
.unwrap()
.header
.seq,
4
);
frame(&b, id, 6, false).await; // a new gap after successful recovery
assert!(
matches!(
tokio::time::timeout(Duration::from_secs(2), controls.recv())
.await
.unwrap(),
Some(DesktopControl::RequestIdr)
),
"a new broken chain still needs recovery"
);
drop(client);
reader.abort();
let _ = reader.await;
a.close(0u32.into(), b"fixture complete");
b.close(0u32.into(), b"fixture complete");
client_ep.close().await;
server_ep.close().await;
}
}
20 changes: 20 additions & 0 deletions docs/native-viewer.md
Original file line number Diff line number Diff line change
Expand Up @@ -374,3 +374,23 @@ that work until it returns; cancellation does not permit an unbounded series
of replacement decoders. Controlled blocking-pool tests cover both cases.
Managed and direct decode use the same work boundary. These diagnostics do
not themselves prove native visible latency or stability under contention.

## Coalesced reference-gap recovery

Repeated successors of the same missing reference may expire while a requested
keyframe is still travelling over reliable QUIC. The receiver now sends one
recovery request for that episode. A received key clears it; an actual rejected
or timed-out reader also permits repair again. Without either outcome, another
request becomes eligible at the existing eight-second key-reader bound. A full
control queue does not arm the hold, and the existing 500 ms request rate limit
still applies. Reader count, memory budgets and read deadlines are unchanged.

Gap records distinguish a sent request from a coalesced observation and include
pending recovery age. Successful frame reads taking at least 250 ms now report
header time, total time, sequence, keyframe flag and byte count. A short scoped
`RUST_LOG=rds_desktop::input_timing=trace` capture can correlate local native
input dispatch with a read-only local image observer, without relying on the
automation tool's clock. It records only sequence, event class and local timing;
key/button codes, coordinates, text and clipboard contents are excluded.
These records locate delay; neither input dispatch nor a frame read proves
physical display response.
Loading
Loading