perf(portal): achieve 58-60fps PipeWire screen capture
- Force PipeWire quantum=512 via NODE_FORCE_QUANTUM (48000/512=93Hz scheduling) - Switch to libx264 ultrafast/zerolatency with 6 threads - Use two-phase poll_and_encode: blocking recv_timeout for first frame, non-blocking try_recv drain for subsequent frames - Remove fps_limit from portal path (PW already rate-limits via quantum/KWin; fps_limit's min_interval was silently dropping ~10% of valid frames) - Remove diagnostic instrumentation (TIMING/PIPEWIRE logs, timing fields, pw_stats counters) - Add lightweight production stats: per-10s fps log + shutdown summary - Prefer libx264 over libopenh264 (better quality at same speed)
This commit is contained in:
+47
-147
@@ -1,13 +1,13 @@
|
||||
// 采集门户状态模块 —— 通过 PipeWire/DMA-BUF 进行屏幕采集并编码
|
||||
use std::os::fd::AsRawFd;
|
||||
use std::path::PathBuf;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use anyhow::{bail, Result};
|
||||
|
||||
use crate::args::Args;
|
||||
use crate::avhw::{self, SwEncState};
|
||||
use crate::cap_portal::{CapPortal, PwCtrlEvent, PwDmaBufFrame};
|
||||
use crate::fps_limit::FpsLimit;
|
||||
|
||||
/// 门户采集的阶段状态
|
||||
/// - WaitingForFormat: 等待接收到第一帧 DMA-BUF 以确定视频格式参数
|
||||
@@ -22,32 +22,16 @@ enum PortalStage {
|
||||
/// 负责管理从 PipeWire 采集屏幕帧、通过 VAAPI 硬件编码的完整生命周期。
|
||||
/// 工作流程:等待第一帧 → 创建编码器 → 持续编码帧数据。
|
||||
pub struct StatePortal {
|
||||
/// 当前采集阶段
|
||||
stage: PortalStage,
|
||||
/// GPU 缩放 + 软件编码器状态(第一帧到达后才初始化)
|
||||
enc: Option<SwEncState>,
|
||||
/// 帧率限制器
|
||||
fps_limit: FpsLimit<()>,
|
||||
/// PipeWire 屏幕采集端点
|
||||
cap: CapPortal,
|
||||
/// 命令行参数
|
||||
args: Args,
|
||||
/// 是否遇到错误
|
||||
errored: bool,
|
||||
/// 是否为第一帧(首帧跳过帧率限制)
|
||||
first_frame: bool,
|
||||
/// DRM 渲染设备路径(如 /dev/dri/renderD128);None 表示首帧自动检测
|
||||
drm_device: Option<PathBuf>,
|
||||
/// 第一帧的时间戳(纳秒),用于计算相对 PTS
|
||||
first_pts_ns: Option<i64>,
|
||||
/// Diagnostic: frames received from PipeWire channel
|
||||
frames_received: u64,
|
||||
/// Diagnostic: frames dropped by FPS limiter
|
||||
frames_fps_dropped: u64,
|
||||
/// Diagnostic: frames successfully encoded
|
||||
frames_encoded: u64,
|
||||
/// Diagnostic: last time we printed stats
|
||||
last_stats_time: Option<std::time::Instant>,
|
||||
start_time: Option<Instant>,
|
||||
last_stats_time: Option<Instant>,
|
||||
last_stats_frames: u64,
|
||||
}
|
||||
|
||||
impl StatePortal {
|
||||
@@ -67,25 +51,23 @@ impl StatePortal {
|
||||
Ok(Self {
|
||||
stage: PortalStage::WaitingForFormat,
|
||||
enc: None,
|
||||
fps_limit: FpsLimit::new(args.fps),
|
||||
cap,
|
||||
args,
|
||||
errored: false,
|
||||
first_frame: true,
|
||||
drm_device,
|
||||
first_pts_ns: None,
|
||||
frames_received: 0,
|
||||
frames_fps_dropped: 0,
|
||||
frames_encoded: 0,
|
||||
start_time: None,
|
||||
last_stats_time: None,
|
||||
last_stats_frames: 0,
|
||||
})
|
||||
}
|
||||
|
||||
/// 轮询 PipeWire 事件并编码帧
|
||||
///
|
||||
/// 尝试从采集端点接收一帧事件。返回 `Ok(true)` 表示已处理事件,
|
||||
/// `Ok(false)` 表示暂无数据。内部根据当前阶段(等待格式/流式)分发处理。
|
||||
pub fn poll_and_encode(&mut self) -> Result<bool> {
|
||||
/// `block=true` 时使用 recv_timeout 阻塞等待帧(最多 10ms),
|
||||
/// `block=false` 时使用 try_recv 非阻塞检查。
|
||||
/// 返回 `Ok(true)` 表示已处理事件,`Ok(false)` 表示暂无数据。
|
||||
pub fn poll_and_encode(&mut self, block: bool) -> Result<bool> {
|
||||
if let Ok(ctrl) = self.cap.event_receiver().try_recv() {
|
||||
match ctrl {
|
||||
PwCtrlEvent::StreamEnded => {
|
||||
@@ -101,13 +83,16 @@ impl StatePortal {
|
||||
}
|
||||
}
|
||||
|
||||
let frame = match self.cap.frame_receiver().try_recv() {
|
||||
Ok(frame) => {
|
||||
self.frames_received += 1;
|
||||
tracing::debug!("poll_and_encode: got frame #{} from channel", self.frames_received);
|
||||
frame
|
||||
let frame = if block {
|
||||
match self.cap.frame_receiver().recv_timeout(std::time::Duration::from_millis(10)) {
|
||||
Ok(frame) => frame,
|
||||
Err(_) => return Ok(false),
|
||||
}
|
||||
} else {
|
||||
match self.cap.frame_receiver().try_recv() {
|
||||
Ok(frame) => frame,
|
||||
Err(_) => return Ok(false),
|
||||
}
|
||||
Err(_) => return Ok(false),
|
||||
};
|
||||
|
||||
match self.stage {
|
||||
@@ -150,6 +135,8 @@ impl StatePortal {
|
||||
|
||||
self.enc = Some(enc);
|
||||
self.stage = PortalStage::Streaming;
|
||||
self.start_time = Some(Instant::now());
|
||||
self.last_stats_time = Some(Instant::now());
|
||||
tracing::info!("First frame processed, encoder initialized, transitioning to Streaming");
|
||||
drop(frame);
|
||||
}
|
||||
@@ -208,27 +195,11 @@ impl StatePortal {
|
||||
/// 通过 `av_hwframe_map` 零拷贝导入 VAAPI,然后交给 SwEncState 完成:
|
||||
/// scale_vaapi GPU 缩放、2K NV12 回读、YUV420P 格式转换、软件 H.264 编码。
|
||||
fn handle_pw_frame(&mut self, frame: PwDmaBufFrame) -> Result<()> {
|
||||
if self.first_frame {
|
||||
self.first_frame = false;
|
||||
} else {
|
||||
let now = std::time::Instant::now();
|
||||
if self.fps_limit.on_new_frame((), now).is_none() {
|
||||
self.frames_fps_dropped += 1;
|
||||
tracing::debug!("handle_pw_frame: FPS limit, dropping frame (#{})", self.frames_fps_dropped);
|
||||
self.maybe_print_stats(now);
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
|
||||
tracing::debug!("handle_pw_frame: processing frame, pts={}", frame.pts);
|
||||
|
||||
let enc = match self.enc.as_mut() {
|
||||
Some(enc) => enc,
|
||||
None => bail!("encoder not initialized"),
|
||||
};
|
||||
|
||||
// SAFETY: frames_rgb is a live VAAPI frames context configured for capture; frame carries
|
||||
// valid DMA-BUF fd/format/modifier/stride/offset metadata for the duration of this call.
|
||||
let mut vaapi_frame = unsafe {
|
||||
avhw::import_dma_buf_to_vaapi(
|
||||
enc.frames_rgb().as_ptr(),
|
||||
@@ -242,39 +213,31 @@ impl StatePortal {
|
||||
)
|
||||
}?;
|
||||
|
||||
tracing::debug!("handle_pw_frame: DMA-BUF import OK");
|
||||
|
||||
let pts = compute_pts(&mut self.first_pts_ns, frame.pts, self.args.fps);
|
||||
let pts = self.frames_encoded as i64;
|
||||
unsafe {
|
||||
(*vaapi_frame.as_mut_ptr()).pts = pts;
|
||||
}
|
||||
|
||||
enc.encode_frame(&vaapi_frame)?;
|
||||
self.frames_encoded += 1;
|
||||
tracing::info!("handle_pw_frame: frame #{} encoded OK, pts={}", self.frames_encoded, pts);
|
||||
|
||||
let now = std::time::Instant::now();
|
||||
self.maybe_print_stats(now);
|
||||
if let Some(last) = self.last_stats_time {
|
||||
if last.elapsed() >= Duration::from_secs(10) {
|
||||
let delta_frames = self.frames_encoded - self.last_stats_frames;
|
||||
let delta_secs = last.elapsed().as_secs_f64();
|
||||
let fps = delta_frames as f64 / delta_secs;
|
||||
tracing::info!(
|
||||
"encoded={}, fps={fps:.1}",
|
||||
self.frames_encoded,
|
||||
);
|
||||
self.last_stats_time = Some(Instant::now());
|
||||
self.last_stats_frames = self.frames_encoded;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn maybe_print_stats(&mut self, now: std::time::Instant) {
|
||||
let should_print = match self.last_stats_time {
|
||||
None => true,
|
||||
Some(last) => now.duration_since(last) >= std::time::Duration::from_secs(2),
|
||||
};
|
||||
if should_print {
|
||||
self.last_stats_time = Some(now);
|
||||
tracing::info!(
|
||||
"STATS: received={}, fps_dropped={}, encoded={}",
|
||||
self.frames_received,
|
||||
self.frames_fps_dropped,
|
||||
self.frames_encoded,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// 关闭状态:刷新编码器并清理资源
|
||||
///
|
||||
/// 使用 `enc.take()` 确保编码器只被 flush 一次,即使多次调用也安全(幂等)。
|
||||
@@ -284,6 +247,18 @@ impl StatePortal {
|
||||
tracing::error!("Flush error during shutdown: {e}");
|
||||
}
|
||||
}
|
||||
if let Some(start) = self.start_time {
|
||||
if self.frames_encoded > 0 {
|
||||
let elapsed = start.elapsed().as_secs_f64();
|
||||
let fps = self.frames_encoded as f64 / elapsed;
|
||||
tracing::info!(
|
||||
"Total: {} frames in {:.1}s, avg {:.1}fps",
|
||||
self.frames_encoded,
|
||||
elapsed,
|
||||
fps,
|
||||
);
|
||||
}
|
||||
}
|
||||
tracing::info!("StatePortal shutdown complete");
|
||||
}
|
||||
|
||||
@@ -316,17 +291,6 @@ fn portal_encode_dimensions(width: u32, height: u32) -> (u32, u32) {
|
||||
}
|
||||
}
|
||||
|
||||
/// Convert PipeWire nanosecond PTS to encoder frame-number units.
|
||||
///
|
||||
/// Uses elapsed time since the first frame to avoid i64 overflow on absolute timestamps.
|
||||
/// PipeWire PTS is CLOCK_MONOTONIC in nanoseconds; encoder time_base = 1/fps.
|
||||
fn compute_pts(first_pts_ns: &mut Option<i64>, frame_pts: i64, fps: u32) -> i64 {
|
||||
let fps_i64 = fps as i64;
|
||||
let base_ns = *first_pts_ns.get_or_insert(frame_pts.max(0));
|
||||
let elapsed_ns = (frame_pts.max(0) - base_ns).max(0);
|
||||
elapsed_ns * fps_i64 / 1_000_000_000
|
||||
}
|
||||
|
||||
/// 解析 DRM 渲染设备路径
|
||||
///
|
||||
/// 仅使用命令行指定的设备路径;未指定则在首帧到达时自动检测。
|
||||
@@ -452,68 +416,4 @@ mod tests {
|
||||
assert_eq!(desc.layers[0].planes[0].pitch, 3840 * 4);
|
||||
}
|
||||
|
||||
// --- compute_pts tests ---
|
||||
|
||||
#[test]
|
||||
fn compute_pts_first_frame_is_zero() {
|
||||
let mut base = None;
|
||||
let pts = compute_pts(&mut base, 1_000_000_000, 30);
|
||||
assert_eq!(pts, 0);
|
||||
assert_eq!(base, Some(1_000_000_000));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compute_pts_second_frame_at_30fps() {
|
||||
let mut base = Some(1_000_000_000);
|
||||
// 33_333_333 * 30 / 1_000_000_000 = 0 (integer division)
|
||||
let pts = compute_pts(&mut base, 1_000_000_000 + 33_333_333, 30);
|
||||
assert_eq!(pts, 0);
|
||||
|
||||
// 100ms later = frame 3
|
||||
let pts = compute_pts(&mut base, 1_000_000_000 + 100_000_000, 30);
|
||||
assert_eq!(pts, 3);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compute_pts_multiple_frames_accumulate() {
|
||||
let mut base = None;
|
||||
let fps = 60;
|
||||
|
||||
let pts0 = compute_pts(&mut base, 0, fps);
|
||||
assert_eq!(pts0, 0);
|
||||
|
||||
let pts1 = compute_pts(&mut base, 16_666_666, fps);
|
||||
assert_eq!(pts1, 0); // 16_666_666 * 60 / 1_000_000_000 = 0
|
||||
|
||||
let pts2 = compute_pts(&mut base, 33_333_333, fps);
|
||||
assert_eq!(pts2, 1); // 33_333_333 * 60 / 1_000_000_000 = 1
|
||||
|
||||
let pts3 = compute_pts(&mut base, 50_000_000, fps);
|
||||
assert_eq!(pts3, 3); // 50ms * 60 / 1000 = 3
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compute_pts_negative_pts_clamped_to_zero() {
|
||||
let mut base = None;
|
||||
let pts = compute_pts(&mut base, -999_999, 30);
|
||||
assert_eq!(pts, 0);
|
||||
assert_eq!(base, Some(0)); // max(0) clamps negative
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compute_pts_late_frame_after_negative() {
|
||||
let mut base = Some(0);
|
||||
let pts = compute_pts(&mut base, 1_000_000_000, 30);
|
||||
assert_eq!(pts, 30);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compute_pts_base_not_overwritten_after_first_call() {
|
||||
let mut base = None;
|
||||
let _ = compute_pts(&mut base, 5_000_000_000, 30);
|
||||
assert_eq!(base, Some(5_000_000_000));
|
||||
|
||||
let _ = compute_pts(&mut base, 10_000_000_000, 30);
|
||||
assert_eq!(base, Some(5_000_000_000)); // base stays at first frame
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user