feat(stats): populate frame_age metric for WebRTC path (partial #20)

Add quantitative capture-to-send latency measurement so we can diagnose
remaining latency sources after #24/#25 PTS fixes.

Previously frame_age_p95 was always 0.0ms because the WebRTC code path
never propagated capture timestamps, even though stats.rs already
supported the metric. The infrastructure existed but was disconnected.

Changes:

- avhw.rs: Add capture_time: Instant field to CpuNv12Frame (set when
  PipeWire delivers frame) and EncodedH264Frame (propagated through
  encode thread via new last_capture_time side-channel on SwEncEncode).

- state_portal.rs: Change sent_gap channel type from Sender<f64> to
  Sender<(f64, Option<f64>)> so WebRTC thread can send pre-computed
  age_ms = capture_time.elapsed() at the exact send moment (not at
  stats drain time, which would inflate the measurement by ~1s).

- stats.rs: record_send_from_thread now accepts Option<f64> age_ms
  and pushes to frame_age_ms Vec when Some.

After this commit:
- stats: log lines show real frame_age_p95 / frame_age_max in ms
- Expected range: 5-30ms (import + scale + encode + channel send)
- If much higher: server pipeline has queueing issue
- If low but user still sees latency: confirms bottleneck is network
  or browser-side (jitter buffer, decode queue)

Scope notes:

- Only Portal/PipeWire path is instrumented. wlr-screencopy path uses
  different code path (EncState, not SwEncState) and will continue to
  report frame_age=0.0ms. Adding wlr instrumentation is separate scope.

- This is diagnostic only — does NOT change user-visible behavior.
  No encoding, sending, or stats output format changes.

Tests:
- cargo build --release: 0 new warnings (19 baseline preserved)
- cargo test: 96 lib + 3 integration, 0 failed
- SAFETY comments preserved verbatim
- 3 files changed, +39/-10 lines

Refs #20.
This commit is contained in:
dailz
2026-06-20 23:20:36 +08:00
parent ad28af6ff3
commit b4b9990efe
3 changed files with 39 additions and 10 deletions
+18
View File
@@ -725,6 +725,8 @@ pub struct EncodedH264Frame {
/// PTS in encoder time_base units (1/fps seconds), normalized so first frame = 0. /// PTS in encoder time_base units (1/fps seconds), normalized so first frame = 0.
/// Derived from real capture time, NOT frame counter. /// Derived from real capture time, NOT frame counter.
pub pts_ticks: i64, pub pts_ticks: i64,
/// Wall-clock capture time, propagated from CpuNv12Frame for frame_age stat.
pub capture_time: std::time::Instant,
} }
pub enum FrameOutput { pub enum FrameOutput {
@@ -740,6 +742,9 @@ pub struct CpuNv12Frame {
pub y_stride: usize, pub y_stride: usize,
pub uv_stride: usize, pub uv_stride: usize,
pub pts: i64, pub pts: i64,
/// Wall-clock time when this frame was captured (PipeWire delivery).
/// Used for frame_age stat: time from capture to WebRTC send.
pub capture_time: std::time::Instant,
} }
pub struct SwEncImport { pub struct SwEncImport {
@@ -987,6 +992,7 @@ impl SwEncImport {
y_stride, y_stride,
uv_stride, uv_stride,
pts, pts,
capture_time: std::time::Instant::now(),
} }
}; };
@@ -1020,6 +1026,10 @@ pub struct SwEncEncode {
/// every `encode_cpu_frame` call (even on early returns) so stale values /// every `encode_cpu_frame` call (even on early returns) so stale values
/// from a previous frame can never leak out. /// from a previous frame can never leak out.
last_timing: SwEncodeTiming, last_timing: SwEncodeTiming,
/// Capture time of the frame currently being encoded. Saved from the
/// input `CpuNv12Frame` so `drain_encoder` can propagate it into the
/// emitted `EncodedH264Frame` for the frame_age stat (issue #20).
last_capture_time: Option<Instant>,
} }
const FNV1A_OFFSET_BASIS: u64 = 0xcbf29ce484222325; const FNV1A_OFFSET_BASIS: u64 = 0xcbf29ce484222325;
@@ -1090,6 +1100,7 @@ impl SwEncEncode {
gop_size, gop_size,
force_keyframe_pending: false, force_keyframe_pending: false,
last_timing: SwEncodeTiming::default(), last_timing: SwEncodeTiming::default(),
last_capture_time: None,
}) })
} }
@@ -1130,6 +1141,7 @@ impl SwEncEncode {
gop_size, gop_size,
force_keyframe_pending: false, force_keyframe_pending: false,
last_timing: SwEncodeTiming::default(), last_timing: SwEncodeTiming::default(),
last_capture_time: None,
}) })
} }
@@ -1154,6 +1166,9 @@ impl SwEncEncode {
pub fn encode_cpu_frame(&mut self, frame: &CpuNv12Frame) -> Result<EncodeOutcome> { pub fn encode_cpu_frame(&mut self, frame: &CpuNv12Frame) -> Result<EncodeOutcome> {
self.last_timing = SwEncodeTiming::default(); self.last_timing = SwEncodeTiming::default();
// Save capture_time so drain_encoder can propagate it into the
// EncodedH264Frame emitted via the WebRTC channel (issue #20).
self.last_capture_time = Some(frame.capture_time);
if self.webrtc_disconnected { if self.webrtc_disconnected {
return Ok(EncodeOutcome::SkippedDisconnected); return Ok(EncodeOutcome::SkippedDisconnected);
@@ -1417,6 +1432,9 @@ impl SwEncEncode {
match tx.try_send(EncodedH264Frame { match tx.try_send(EncodedH264Frame {
data: data.to_vec(), data: data.to_vec(),
pts_ticks, pts_ticks,
capture_time: self
.last_capture_time
.unwrap_or_else(Instant::now),
}) { }) {
Ok(()) => {} Ok(()) => {}
Err(crossbeam_channel::TrySendError::Full(frame)) => { Err(crossbeam_channel::TrySendError::Full(frame)) => {
+14 -7
View File
@@ -38,7 +38,7 @@ struct EncodeThread {
struct WebrtcThread { struct WebrtcThread {
handle: Option<std::thread::JoinHandle<()>>, handle: Option<std::thread::JoinHandle<()>>,
sent_gap_rx: crossbeam_channel::Receiver<f64>, sent_gap_rx: crossbeam_channel::Receiver<(f64, Option<f64>)>,
} }
/// 门户模式的主状态机 /// 门户模式的主状态机
@@ -281,7 +281,8 @@ impl StatePortal {
.clone(); .clone();
let fps = self.args.fps; let fps = self.args.fps;
let max_bitrate = self.args.max_bitrate; let max_bitrate = self.args.max_bitrate;
let (sent_gap_tx, sent_gap_rx) = crossbeam_channel::bounded(64); let (sent_gap_tx, sent_gap_rx) =
crossbeam_channel::bounded::<(f64, Option<f64>)>(64);
let webrtc_handle = std::thread::Builder::new() let webrtc_handle = std::thread::Builder::new()
.name("wl-webrtc-webrtc".into()) .name("wl-webrtc-webrtc".into())
.spawn(move || { .spawn(move || {
@@ -349,8 +350,8 @@ impl StatePortal {
} }
} }
if let Some(ref webrtc_thread) = self.webrtc_thread { if let Some(ref webrtc_thread) = self.webrtc_thread {
while let Ok(gap_ms) = webrtc_thread.sent_gap_rx.try_recv() { while let Ok((gap_ms, age_ms)) = webrtc_thread.sent_gap_rx.try_recv() {
self.stats.record_send_from_thread(gap_ms); self.stats.record_send_from_thread(gap_ms, age_ms);
} }
} }
let snap = self.stats.snapshot_and_reset(); let snap = self.stats.snapshot_and_reset();
@@ -694,7 +695,7 @@ fn webrtc_thread_loop(
enc_height: u32, enc_height: u32,
max_bitrate: u64, max_bitrate: u64,
paused: Arc<AtomicBool>, paused: Arc<AtomicBool>,
sent_gap_tx: crossbeam_channel::Sender<f64>, sent_gap_tx: crossbeam_channel::Sender<(f64, Option<f64>)>,
bitrate_tx: crossbeam_channel::Sender<BitrateCommand>, bitrate_tx: crossbeam_channel::Sender<BitrateCommand>,
resolution_tx: crossbeam_channel::Sender<BitrateCommand>, resolution_tx: crossbeam_channel::Sender<BitrateCommand>,
) { ) {
@@ -799,8 +800,12 @@ fn webrtc_thread_loop(
let gap_ms = last_send let gap_ms = last_send
.map(|l| l.elapsed().as_secs_f64() * 1000.0) .map(|l| l.elapsed().as_secs_f64() * 1000.0)
.unwrap_or(0.0); .unwrap_or(0.0);
// Compute capture-to-send age on the sending thread so the
// frame_age stat stays accurate when batch-drained later.
let age_ms =
Some(enc_frame.capture_time.elapsed().as_secs_f64() * 1000.0);
last_send = Some(std::time::Instant::now()); last_send = Some(std::time::Instant::now());
let _ = sent_gap_tx.try_send(gap_ms); let _ = sent_gap_tx.try_send((gap_ms, age_ms));
} }
} else { } else {
while webrtc_rx.try_recv().is_ok() {} while webrtc_rx.try_recv().is_ok() {}
@@ -817,8 +822,10 @@ fn webrtc_thread_loop(
let gap_ms = last_send let gap_ms = last_send
.map(|l| l.elapsed().as_secs_f64() * 1000.0) .map(|l| l.elapsed().as_secs_f64() * 1000.0)
.unwrap_or(0.0); .unwrap_or(0.0);
let age_ms =
Some(enc_frame.capture_time.elapsed().as_secs_f64() * 1000.0);
last_send = Some(std::time::Instant::now()); last_send = Some(std::time::Instant::now());
let _ = sent_gap_tx.try_send(gap_ms); let _ = sent_gap_tx.try_send((gap_ms, age_ms));
} }
} }
Err(crossbeam_channel::RecvTimeoutError::Timeout) => {} Err(crossbeam_channel::RecvTimeoutError::Timeout) => {}
+7 -3
View File
@@ -168,13 +168,17 @@ impl PipelineStats {
/// Record a frame sent from a background WebRTC thread. /// Record a frame sent from a background WebRTC thread.
/// `gap_ms` is the pre-computed time since the previous send (0.0 = first frame). /// `gap_ms` is the pre-computed time since the previous send (0.0 = first frame).
/// Unlike `record_send`, this does not sample `Instant::now()`, so it remains /// `age_ms` is the pre-computed capture-to-send latency (None if unavailable).
/// accurate even when batch-drained at stats snapshot time. /// Both are pre-computed on the sending thread to remain accurate when
pub fn record_send_from_thread(&mut self, gap_ms: f64) { /// batch-drained at stats snapshot time on the main thread.
pub fn record_send_from_thread(&mut self, gap_ms: f64, age_ms: Option<f64>) {
if gap_ms > 0.0 { if gap_ms > 0.0 {
self.sent_gaps_ms.push(gap_ms); self.sent_gaps_ms.push(gap_ms);
} }
self.sent_frames += 1; self.sent_frames += 1;
if let Some(age) = age_ms {
self.frame_age_ms.push(age);
}
} }
/// Update PipeWire dropped counter (absolute value from AtomicU64). /// Update PipeWire dropped counter (absolute value from AtomicU64).