From fffa440e68a9734e1aad22bff3f785a6717294e9 Mon Sep 17 00:00:00 2001 From: dailz Date: Mon, 22 Jun 2026 17:18:00 +0800 Subject: [PATCH] =?UTF-8?q?docs(state=5Fportal):=20[1/2]=20=E4=B8=AD?= =?UTF-8?q?=E6=96=87=E6=B3=A8=E9=87=8A=20Portal=20=E7=8A=B6=E6=80=81?= =?UTF-8?q?=E6=9C=BA=E5=88=9D=E5=A7=8B=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/state_portal.rs | 101 ++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 101 insertions(+) diff --git a/src/state_portal.rs b/src/state_portal.rs index 15eb327..40d3a54 100644 --- a/src/state_portal.rs +++ b/src/state_portal.rs @@ -1,3 +1,32 @@ +//! Portal 后端的主状态机:通过 PipeWire + DMA-BUF 进行屏幕采集并软件编码。 +//! +//! ## 整体角色 +//! +//! `StatePortal` 与 `src/state.rs::State` 是平行的两条采集路径: +//! - `state.rs`(wlroots 路径):由外层 `mio` 事件循环驱动(手工版 epoll), +//! 通过 `zwlr_screencopy_manager_v1` 协议一帧一帧地拉取。 +//! - `state_portal.rs`(本文件,XDG Portal / PipeWire 路径):由 `CapPortal` +//! 通过 `crossbeam_channel::Receiver` 推帧;本状态机只负责"消费"。 +//! +//! ## 异步模型的真相 +//! +//! 本文件**不**使用 `mio` 或 `tokio`——`CapPortal` 内部在独立线程跑 PipeWire +//! asyncio loop,把 DMA-BUF 帧通过 crossbeam channel 投递出来;外层 `main.rs` +//! 只需在 `while !is_errored()` 循环里轮询 `poll_and_encode(block)`。编码线程与 +//! WebRTC 线程通过 `std::thread::spawn`(不是 `tokio::spawn`)启动,再借助 +//! crossbeam channel 与主线程通信——类比 Go 的 `go func()` + channel。 +//! +//! ## 阶段机 +//! +//! `PortalStage::WaitingForFormat`(等首帧以确定格式)→ `Streaming`(持续编码)。 +//! +//! ## 注意 +//! +//! - T9a(本块)覆盖文件头 + struct 定义 + `impl StatePortal`(至 `fn encode_thread_loop` 之前); +//! T9b 覆盖 `encode_thread_loop` / `webrtc_thread_loop` / `resolve_drm_device` 等自由函数。 +//! - 多处 `unsafe` 调用 FFmpeg/VAAPI FFI;现有英文 `// SAFETY:` 保留不动, +//! 本任务在每个 unsafe 块上方加普通 `//` 中文概述(不新增 `// SAFETY:`)。 + // 采集门户状态模块 —— 通过 PipeWire/DMA-BUF 进行屏幕采集并编码 use std::os::fd::AsRawFd; use std::path::PathBuf; @@ -24,12 +53,24 @@ enum PortalStage { Streaming, } +/// 编码线程单帧计时回执——由 `encode_thread_loop` 通过 `timing_tx` 发回主线程, +/// 用于在 `PipelineStats` 中窗口化统计 `sws_us`(libswscale 缩放开销)和 +/// `encode_us`(H.264 软件编码开销)。类比 Go 的 `type EncodeThreadTiming struct`。 struct EncodeThreadTiming { sws_us: u64, encode_us: u64, output_bytes: usize, } +/// 编码工作线程的句柄与通信端点。 +/// +/// 由主线程持有,负责把 NV12 帧 (`CpuNv12Frame`) 通过 `input_tx` 投递给 +/// `encode_thread_loop`;编码完成后通过 `timing_rx` 收回单帧计时;`duplicate_count` +/// 是跨线程共享的 `Arc`(类比 Go 的 `*uint64` protected by atomic), +/// 用于统计被去重跳过的帧数(影响 BWE 与丢弃策略)。 +/// +/// 字段全部用 `Option<...>`/`Sender`/`Receiver` 包装,是为了在 `shutdown` 时 +/// 能用 `Option::take()` 把所有权转移到本地变量、显式 drop `input_tx`、再 `join()`。 struct EncodeThread { handle: Option>, input_tx: crossbeam_channel::Sender, @@ -37,6 +78,11 @@ struct EncodeThread { duplicate_count: std::sync::Arc, } +/// WebRTC 工作线程的句柄与单向上行通道。 +/// +/// 主线程只能通过 `sent_gap_rx` **被动接收** WebRTC 线程上报的"已发送帧间隔/老化" +/// 指标(用于 `PipelineStats::record_send_from_thread`)。下行的码率 / 分辨率 / 暂停 +/// 控制走另一组 channel(`bitrate_tx` / `resolution_tx` / `webrtc_paused`),不在此处。 struct WebrtcThread { handle: Option>, sent_gap_rx: crossbeam_channel::Receiver<(f64, Option)>, @@ -71,6 +117,14 @@ pub struct StatePortal { last_pts_emitted: Option, } +// `impl StatePortal` 块集中了门户路径的所有主线程逻辑: +// - `new`:构造(DRM 设备探测 + CapPortal 初始化;编码器延后到首帧)。 +// - `poll_and_encode`:外层 main 循环每轮调用一次,处理 1 个 PipeWire 帧 / 控制事件。 +// - `shutdown`:幂等清理(编码线程 → WebRTC 线程 → MP4 flush)。 +// - 私有辅助:`record_capture_timeout` / `record_frame_arrival`(采集空闲日志节流)、 +// `resolve_drm_device_for_frame`(DMA-BUF 导入兼容性探测)、 +// `handle_pw_frame`(VAAPI 导入 + 软件编码)、`compute_capture_pts`(90kHz RTP PTS)。 +// 内部不使用任何锁——所有 `&mut self` 由外层 main 循环单线程串行化保证独占。 impl StatePortal { /// 创建门户状态实例 /// @@ -219,6 +273,10 @@ impl StatePortal { if self.webrtc.is_some() { let paused = self.webrtc_paused.as_ref() .ok_or_else(|| anyhow::anyhow!("internal invariant broken: webrtc_paused missing while WebRTC mode is active"))?; + // WebRTC 模式需要 6 路 crossbeam channel 协调主线程 ↔ 编码线程 ↔ WebRTC 线程。 + // `crossbeam_channel::bounded::(n)` 类比 Go 的 `make(chan T, n)`—— + // 容量满时 `send` 阻塞、空时 `recv` 阻塞;返回的 `(Sender, Receiver)` 各占一份 + //所有权,可 move 到不同线程(前提是元素类型 `T: Send`)。 let (resolution_tx, resolution_rx) = crossbeam_channel::bounded::(4); let (encoder_resolution_tx, encoder_resolution_rx) = @@ -252,7 +310,19 @@ impl StatePortal { let duplicate_count = std::sync::Arc::new( std::sync::atomic::AtomicU64::new(0), ); + // Arc 引用计数克隆(不是深拷贝)——`duplicate_count` 留在主线程, + // `duplicate_count_for_thread` move 进编码线程;两者指向同一原子。 + // 类比 Go 的 `*uint64` + atomic.Store,但 Rust 用类型系统保证线程安全。 let duplicate_count_for_thread = duplicate_count.clone(); + // `std::thread::Builder::new().name(...).spawn(move || {...})?`: + // - 类比 Go 的 `go func() {...}()`,但返回 `JoinHandle` 而非 fire-and-forget—— + // 主线程可在 shutdown 时 `handle.join()` 等待子线程退出。 + // - **不**用 `tokio::spawn`:编码是 CPU 密集 + 阻塞 FFmpeg 调用, + // 不需要 async/await;标准线程更直接。 + // - `move ||` 闭包:把 `encode` / `input_rx` / `timing_tx` / + // `duplicate_count_for_thread` 的所有权**转移**给子线程(类比 Go 里把变量 + // 显式传入 goroutine 闭包参数)。 + // - `?` 传播 `io::Error`——线程创建可能失败(资源限制)。 let handle = std::thread::Builder::new() .name("wl-webrtc-encode".into()) .spawn(move || { @@ -283,6 +353,10 @@ impl StatePortal { let max_bitrate = self.args.max_bitrate; let (sent_gap_tx, sent_gap_rx) = crossbeam_channel::bounded::<(f64, Option)>(64); + // WebRTC 工作线程:同上 `std::thread::spawn(move || ...)` 模式—— + // 内部跑 str0m 的 asyncio loop(`WebRtcState` 自己驱动), + // 通过 `webrtc_rx` 接收 H.264 帧、通过 `bitrate_tx` / `resolution_tx` + // 接收码率/分辨率指令、通过 `sent_gap_tx` 上报发送指标。 let webrtc_handle = std::thread::Builder::new() .name("wl-webrtc-webrtc".into()) .spawn(move || { @@ -366,6 +440,11 @@ impl StatePortal { Ok(true) } + /// 记录"采集超时"——本次轮询未取到帧(PipeWire 队列空)。 + /// + /// 因为 Wayland 是 damage-driven(只有画面变化才推帧),静态画面下长时间无帧 + /// 是**正常**行为,不是 compositor 卡死。所以本函数只做"5 秒阈值后的 DEBUG 一次性日志", + /// 用 `idle_log_start` 字段保证每次空闲区间只发一条日志(issue #15 / #18)。 fn record_capture_timeout(&mut self) { let Some(last_capture_arrival) = self.last_capture_arrival else { return; @@ -391,6 +470,11 @@ impl StatePortal { } } + /// 记录"采集到达"——本次轮询成功取到一帧。 + /// + /// 与 `record_capture_timeout` 互补:若之前处于空闲区间,则通过 `Option::take()` + /// 取出 `idle_log_start` 并发一条 "resumed after idle" DEBUG 日志;然后刷新 + /// `last_capture_arrival` 时间戳。两者共同实现"一次性空闲日志"语义。 fn record_frame_arrival(&mut self) { if let Some(idle_start) = self.idle_log_start.take() { tracing::debug!( @@ -459,6 +543,10 @@ impl StatePortal { // processing — DMA-BUF import, VAAPI scale, NV12 clone, channel send, and // encode thread wakeup. This eliminates ~60fps of pointless work during // the pre-connect idle window. MP4 mode (webrtc_paused == None) is unaffected. + // `Arc` 类比 Go 的 `*atomic.Bool`——`Arc` 提供跨线程共享所有权 + // (引用计数原子递增/递减),`AtomicBool` 提供无锁读/写。 + // `Ordering::Relaxed`:只保证单变量原子性,不建立与其他变量的 happens-before 关系—— + // 对"暂停标志"足够(不需要它做屏障同步)。 if let Some(paused) = &self.webrtc_paused { if paused.load(Ordering::Relaxed) { return Ok(()); @@ -478,6 +566,10 @@ impl StatePortal { if let Some(enc) = self.enc.as_mut() { // 将 DMA-BUF 帧零拷贝导入 VAAPI 硬件帧池 + // unsafe:FFI 调用 FFmpeg `av_hwframe_ctx_init` / `av_hwframe_map` 系列, + // 内部会读取 `enc.frames_rgb()` 指向的 `AVBufferRef`(硬件帧池), + // 并把 `frame.fd.as_raw_fd()`(DMA-BUF dmabuf fd)注册到 VAAPI。 + // 安全性前提:`enc` 在本线程独占(main 串行化保证)、`frame.fd` 未被 close。 let mut vaapi_frame = unsafe { avhw::import_dma_buf_to_vaapi( enc.frames_rgb().as_ptr(), @@ -515,6 +607,8 @@ impl StatePortal { }; self.stats.record_encode(&timings); } else if let Some(import) = self.enc_import.as_mut() { + // 同上 unsafe:DMA-BUF → VAAPI 导入;`import.frames_rgb()` 是与编码线程 + // **不共享**的独立硬件帧池(避免与 `import_and_scale` 的回读路径竞争)。 let mut vaapi_frame = unsafe { avhw::import_dma_buf_to_vaapi( import.frames_rgb().as_ptr(), @@ -540,6 +634,9 @@ impl StatePortal { "internal invariant broken: encode thread missing while async import is active" ) })?; + // `try_send` 类比 Go 的 `select { case ch <- v: default: }`—— + // 非阻塞投递;三种结果分别处理:成功递增、满了丢弃(DEBUG 日志)、 + // 对端关闭(致命,置 `errored=true` 让外层循环退出)。 match enc_thread.input_tx.try_send(cpu_nv12) { Ok(()) => { self.frames_encoded += 1; @@ -608,6 +705,10 @@ impl StatePortal { self.shutdown_started = true; // 1. Stop encode thread (drops webrtc_tx → signals WebRTC thread to exit) + // `Option::take()` 把 `EncodeThread` 的所有权从 `self.enc_thread` 转移到本地 `enc_thread`, + // 同时 `self.enc_thread` 变成 `None`——这是 Rust 里"消费字段但保留父结构体"的标准习语, + // 类比 Go 里把字段设为 nil 但保留外层 struct。接下来显式 `drop(input_tx)` 关闭 channel, + // 编码线程的 `input_rx.recv()` 会返回 `Err(Disconnected)` 从而退出循环。 if let Some(mut enc_thread) = self.enc_thread.take() { drop(enc_thread.input_tx); if let Some(handle) = enc_thread.handle.take() {