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
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion crates/orbiscreen-capture/src/kwin_virtual.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
4 changes: 2 additions & 2 deletions crates/orbiscreen-capture/src/wayland.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?
Expand Down
5 changes: 3 additions & 2 deletions crates/orbiscreen-encode/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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")
Expand Down Expand Up @@ -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", &params.bitrate_kbps.to_string());
Expand Down
2 changes: 1 addition & 1 deletion crates/orbiscreen-transport/src/aoa.rs
Original file line number Diff line number Diff line change
Expand Up @@ -374,7 +374,7 @@ pub fn run_accessory_bridge(
);

let (prio_tx, prio_rx) = std::sync::mpsc::channel::<Vec<u8>>();
let (video_tx, video_rx) = std::sync::mpsc::sync_channel::<Vec<u8>>(64);
let (video_tx, video_rx) = std::sync::mpsc::sync_channel::<Vec<u8>>(8);
let running_writer = running.clone();
let fd_writer = fd;
let writer_handle = std::thread::spawn(move || {
Expand Down
49 changes: 38 additions & 11 deletions crates/orbiscreen-transport/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ impl Transport {
idr_tx: Option<mpsc::Sender<()>>,
) -> Result<(), TransportError> {
let input_tx = self.input_tx;
let (video_tx, _video_rx) = tokio::sync::broadcast::channel::<H264Packet>(32);
let (video_tx, _video_rx) = tokio::sync::broadcast::channel::<H264Packet>(8);
let state = AppState {
config: self.cfg.clone(),
input_tx,
Expand Down Expand Up @@ -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::<gstreamer::Pipeline>().map_err(|_| ())?;
let appsrc = pipeline
Expand Down Expand Up @@ -808,7 +808,7 @@ async fn stream_handler(

gstreamer::init().ok();

let (tx, rx) = tokio::sync::mpsc::channel::<Vec<u8>>(32);
let (tx, rx) = tokio::sync::mpsc::channel::<Vec<u8>>(8);
let tx_alive = tx.clone();

let setup_pipeline = |pipeline: &gstreamer::Pipeline,
Expand All @@ -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() {
Expand Down Expand Up @@ -895,14 +899,17 @@ 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 {
let _pipeline_guard = PipelineGuard(pipeline_for_task);
let _guard = ClientGuard(stats);

let mut wait_keyframe = true;
let mut pts_base: Option<u64> = None;
let mut stream_pts_ns: u64 = 0;
let mut last_pkt_pts_ns: Option<u64> = None;
loop {
if tx_alive.is_closed() {
debug!("stream client disconnected");
Expand All @@ -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(());
}
}
}
}
Expand Down
Loading