From 9b6fae00acffbc2495176cb9f0932054caf1174d Mon Sep 17 00:00:00 2001 From: yashb404 Date: Sun, 13 Sep 2026 18:58:58 +0530 Subject: [PATCH] orbiscreen | fix: drop stale frames and eliminate stream latency accumulation (#77) --- CHANGELOG.md | 10 ++++ crates/orbiscreen-capture/src/kwin_virtual.rs | 2 +- crates/orbiscreen-capture/src/wayland.rs | 4 +- crates/orbiscreen-encode/src/lib.rs | 5 +- crates/orbiscreen-transport/src/aoa.rs | 2 +- crates/orbiscreen-transport/src/lib.rs | 49 ++++++++++++++----- 6 files changed, 55 insertions(+), 17 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 93188f7f..6808b958 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,16 @@ All notable changes to this project will be documented in this file. +## [Unreleased] + +### Bug Fixes +- **Drop Stale Frames & Eliminate Stream Latency Accumulation (#77)**: + - Configure HTTP MPEG-TS pipeline with `is-live=true do-timestamp=true` and `appsink drop=true sync=false max-buffers=1`, ensuring stale frames are dropped immediately under backpressure instead of accumulating seconds of video and input latency. + - Fix PTS timeline desync in `stream_handler` by clamping PTS delta across frame gaps (>250ms) to nominal frame time, preventing client decoder buffer inflation after laptop suspend/resume or heavy stalls. + - Reduce AOA USB accessory `sync_channel` capacity from 64 to 8 chunks to eliminate host-side transport queueing. + - Multi-thread software color conversion with `n-threads=4` on `videoconvert` in capture and encode pipelines to prevent CPU bottlenecks during color space conversion at higher resolutions. + - Constrain `appsrc` max-bytes to 1 uncompressed frame and `appsink` max-buffers to 1 across capture pipelines to prevent buffer bloat. + ## [v0.28.0] - 2026-09-12 USB AOA frame drop and stutter elimination, secondary display black screen resolution, relative mouse cross-screen traversal, and 1:1 touch coordinate alignment: eliminate USB video frame drops on Screen 1 by expanding socket buffer to 256KB and tuning ExoPlayer buffer pacing, fix secondary display black screen by implementing dynamic wl_output binding and configure retries in damage pump, restore standard relative mouse motion across all displays without boundary confinement, and align Android touch and stylus coordinates 1:1 to host display pixels. diff --git a/crates/orbiscreen-capture/src/kwin_virtual.rs b/crates/orbiscreen-capture/src/kwin_virtual.rs index e292a021..61fec68c 100644 --- a/crates/orbiscreen-capture/src/kwin_virtual.rs +++ b/crates/orbiscreen-capture/src/kwin_virtual.rs @@ -678,7 +678,7 @@ impl KwinVirtualCapture { let pipeline_str = format!( "pipewiresrc path={node_id} do-timestamp=true \ ! video/x-raw \ - ! videoconvert \ + ! videoconvert n-threads=4 \ ! videoscale \ ! video/x-raw,format=BGRA,width={},height={} \ ! appsink name=sink drop=true sync=false max-buffers=1 emit-signals=false", diff --git a/crates/orbiscreen-capture/src/wayland.rs b/crates/orbiscreen-capture/src/wayland.rs index abc3bec5..cfcc3030 100644 --- a/crates/orbiscreen-capture/src/wayland.rs +++ b/crates/orbiscreen-capture/src/wayland.rs @@ -169,10 +169,10 @@ impl WaylandCapture { let pipeline_str = format!( "pipewiresrc fd={} path={} do-timestamp=true \ ! video/x-raw \ - ! videoconvert \ + ! videoconvert n-threads=4 \ ! videoscale \ ! video/x-raw,format=BGRA,width={},height={} \ - ! appsink name=sink drop=false sync=false max-buffers=2 emit-signals=false", + ! appsink name=sink drop=true sync=false max-buffers=1 emit-signals=false", raw_fd, node_id, spec.width, spec.height ); let pipeline = gstreamer::parse::launch(&pipeline_str)? diff --git a/crates/orbiscreen-encode/src/lib.rs b/crates/orbiscreen-encode/src/lib.rs index 46e42f48..5d1d0d37 100644 --- a/crates/orbiscreen-encode/src/lib.rs +++ b/crates/orbiscreen-encode/src/lib.rs @@ -225,6 +225,7 @@ impl Encoder { .map_err(|_| EncodeError::Pipeline("appsrc downcast".into()))?; let videoconvert = make_element("videoconvert")?; + set_str_if_present(&videoconvert, "n-threads", "4"); let appsink = ElementFactory::make("appsink") .build() @@ -233,7 +234,7 @@ impl Encoder { .map_err(|_| EncodeError::Pipeline("appsink downcast".into()))?; appsink.set_sync(false); appsink.set_drop(true); - appsink.set_max_buffers(2); + appsink.set_max_buffers(1); appsink.set_caps(Some( &gstreamer::Caps::builder("video/x-h264") .field("stream-format", "byte-stream") @@ -264,7 +265,7 @@ impl Encoder { appsrc.set_format(gstreamer::Format::Time); appsrc.set_is_live(true); appsrc.set_do_timestamp(false); - appsrc.set_max_bytes((params.width as u64) * params.height as u64 * 4 * 4); + appsrc.set_max_bytes((params.width as u64) * params.height as u64 * 4); if encoder.find_property("bitrate").is_some() { encoder.set_property_from_str("bitrate", ¶ms.bitrate_kbps.to_string()); diff --git a/crates/orbiscreen-transport/src/aoa.rs b/crates/orbiscreen-transport/src/aoa.rs index fa980ca5..6e6742db 100644 --- a/crates/orbiscreen-transport/src/aoa.rs +++ b/crates/orbiscreen-transport/src/aoa.rs @@ -374,7 +374,7 @@ pub fn run_accessory_bridge( ); let (prio_tx, prio_rx) = std::sync::mpsc::channel::>(); - let (video_tx, video_rx) = std::sync::mpsc::sync_channel::>(64); + let (video_tx, video_rx) = std::sync::mpsc::sync_channel::>(8); let running_writer = running.clone(); let fd_writer = fd; let writer_handle = std::thread::spawn(move || { diff --git a/crates/orbiscreen-transport/src/lib.rs b/crates/orbiscreen-transport/src/lib.rs index 2e03ffcb..4c5da9d5 100644 --- a/crates/orbiscreen-transport/src/lib.rs +++ b/crates/orbiscreen-transport/src/lib.rs @@ -244,7 +244,7 @@ impl Transport { idr_tx: Option>, ) -> Result<(), TransportError> { let input_tx = self.input_tx; - let (video_tx, _video_rx) = tokio::sync::broadcast::channel::(32); + let (video_tx, _video_rx) = tokio::sync::broadcast::channel::(8); let state = AppState { config: self.cfg.clone(), input_tx, @@ -775,11 +775,11 @@ fn build_video_pipeline() -> Result< > { use gstreamer::prelude::*; use gstreamer_app::{AppSink, AppSrc}; - let pipeline_str = "appsrc name=src format=time is-live=false block=false \ + let pipeline_str = "appsrc name=src format=time is-live=true block=false do-timestamp=true \ ! video/x-h264,stream-format=byte-stream,alignment=au \ ! h264parse config-interval=1 \ ! mpegtsmux alignment=7 \ - ! appsink name=sink drop=false sync=false max-buffers=16 emit-signals=false"; + ! appsink name=sink drop=true sync=false max-buffers=1 emit-signals=false"; let p = gstreamer::parse::launch(pipeline_str).map_err(|_| ())?; let pipeline = p.downcast::().map_err(|_| ())?; let appsrc = pipeline @@ -808,7 +808,7 @@ async fn stream_handler( gstreamer::init().ok(); - let (tx, rx) = tokio::sync::mpsc::channel::>(32); + let (tx, rx) = tokio::sync::mpsc::channel::>(8); let tx_alive = tx.clone(); let setup_pipeline = |pipeline: &gstreamer::Pipeline, @@ -821,9 +821,13 @@ async fn stream_handler( .build(); appsrc.set_caps(Some(&caps)); appsrc.set_format(gstreamer::Format::Time); - appsrc.set_max_bytes(256 * 1024); + appsrc.set_max_bytes(128 * 1024); appsrc.set_block(false); + appsink.set_sync(false); + appsink.set_drop(true); + appsink.set_max_buffers(1); + appsink.set_callbacks( AppSinkCallbacks::builder() .new_sample(move |sink| match sink.pull_sample() { @@ -895,6 +899,8 @@ async fn stream_handler( } let pipeline_for_task = pipeline.clone(); + let refresh_hz = state.refresh_hz.max(1); + let nominal_frame_ns = 1_000_000_000u64 / u64::from(refresh_hz); let mut video_rx = state.video_tx.subscribe(); tokio::spawn(async move { @@ -902,7 +908,8 @@ async fn stream_handler( let _guard = ClientGuard(stats); let mut wait_keyframe = true; - let mut pts_base: Option = None; + let mut stream_pts_ns: u64 = 0; + let mut last_pkt_pts_ns: Option = None; loop { if tx_alive.is_closed() { debug!("stream client disconnected"); @@ -925,20 +932,40 @@ async fn stream_handler( continue; } wait_keyframe = false; - if pts_base.is_none() { - pts_base = Some(pkt.pts_ns); - } } - let base = pts_base.unwrap_or(0); + let delta_ns = match last_pkt_pts_ns { + Some(last) => { + let diff = pkt.pts_ns.saturating_sub(last); + // If the gap between packets exceeds 250ms (e.g. system suspend/resume, + // GPU sleep, or extreme lag spike), do not pass the massive time jump into + // mpegtsmux and downstream players. Clamp the progression to nominal frame duration. + if diff > 250_000_000 { + debug!( + diff_ms = diff / 1_000_000, + "large timestamp gap detected; clamping stream PTS delta to prevent player desync" + ); + nominal_frame_ns + } else { + diff + } + } + None => 0, + }; + last_pkt_pts_ns = Some(pkt.pts_ns); + stream_pts_ns = stream_pts_ns.saturating_add(delta_ns); + let mut normalized = pkt; - normalized.pts_ns = normalized.pts_ns.saturating_sub(base); + normalized.pts_ns = stream_pts_ns; if let Err(e) = push_h264_packet(&appsrc_clone, &normalized) { match e { gstreamer::FlowError::Flushing | gstreamer::FlowError::Eos => break, _ => { wait_keyframe = true; + if let Some(tx) = &idr_tx { + let _ = tx.try_send(()); + } } } }