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
5 changes: 5 additions & 0 deletions crates/rds-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,11 @@ pub const PROTOCOL_VERSION: u16 = 3;
/// Upper bound for a serialized greeting, guard against abusive peers.
pub const MAX_MESSAGE_LEN: u32 = 64 * 1024;

/// Desktop frame STOP_SENDING code: this predecessor is no longer needed
/// after a newer independent picture. It is neither fresh successful delivery
/// nor a broken current reference. Older senders may treat it as generic failure.
pub const DESKTOP_FRAME_OBSOLETE: u32 = 0x5244_5301;

/// First frame on every uni-directional stream (v3): routes the stream
/// to the service that owns it. The accepting side runs one
/// `accept_uni` demux per connection and hands each stream to the
Expand Down
26 changes: 19 additions & 7 deletions crates/rds-desktop/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -719,6 +719,7 @@ impl GapRepair {

enum FrameRead {
Complete(FrameHeader, Vec<u8>, FrameBudget),
Stale,
Rejected,
TimedOut,
}
Expand All @@ -736,6 +737,7 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
let mut admitted = 0u64;
let mut completed = 0u64;
let mut rejected = 0u64;
let mut stale = 0u64;
let mut timed_out = 0u64;
let mut gaps = 0u64;
let mut received_keys = 0u64;
Expand Down Expand Up @@ -815,13 +817,14 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) {
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");
tracing::info!(admitted,completed,rejected,stale,timed_out,gaps,received_keys,in_flight=ctx.receiving.load(Ordering::Relaxed),"desktop receiver health");
health=std::time::Instant::now();
}
},
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::Stale))=>{stale+=1;},
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());},
}
Expand Down Expand Up @@ -870,18 +873,27 @@ async fn read_one(
let frame_seq = header.seq;
let keyframe = header.keyframe;
let header_ms = started.elapsed().as_millis();
if header.seq == u64::MAX
|| crate::frame_bytes(header.width as usize, header.height as usize).is_none()
{
return FrameRead::Rejected;
}
let expected_seq = next_seq.load(Ordering::Relaxed);
if header.seq < expected_seq {
let _ = stream.stop(rds_core::DESKTOP_FRAME_OBSOLETE.into());
tracing::debug!(
frame_seq,
expected_seq,
"desktop obsolete predecessor released"
);
return FrameRead::Stale;
}
let deadline = if keyframe {
KEYFRAME_READ_TIMEOUT
} else {
FRAME_READ_TIMEOUT
};
let reading = async {
if header.seq == u64::MAX
|| header.seq < next_seq.load(Ordering::Relaxed)
|| crate::frame_bytes(header.width as usize, header.height as usize).is_none()
{
return None;
}
// Retain metadata-only progress when the future times out. read_to_end
// discards its partial buffer on cancellation and hid whether even a
// header or any media bytes arrived during observed freezes.
Expand Down
Loading
Loading