//! # wl-webrtc 程序入口(main 函数所在文件) //! //! 本文件是 `wl-webrtc` 二进制 crate 的入口,等价于 Go 的 `func main()`。 //! 由于 Rust 的 `main()` 不允许返回错误(`Result`),本项目采用通用模式: //! 真正的业务逻辑写在 `fn run() -> Result<()>`,而 `main()` 直接 `run()` 完成所有工作。 //! //! 整体执行流程: //! 1. 通过 `clap` 解析命令行参数(`Args`,包含分辨率、编码格式、帧率等) //! 2. 初始化 `tracing` 日志系统(受 `RUST_LOG` 环境变量或 `-v` 参数控制) //! 3. MVP 阶段拒绝非 H.264 编码格式 //! 4. 要求至少提供 `--output`(输出到文件)或 `--port`(启动 WebRTC 信号服务器) //! 5. 调用 `backend_detect::detect_backend` 自动检测当前 Wayland 桌面支持的截屏后端 //! 6. 根据检测结果进入对应的事件循环: //! - 支持 `zwlr_screencopy_manager_v1` 的合成器(Sway/Hyprland)→ `run_wlr_screencopy` //! - 仅支持 XDG Portal ScreenCast 的桌面(GNOME/KDE)→ `run_portal_pipewire` //! //! 两个事件循环都基于 `mio`(一个手动驱动的事件循环库,类似 Go runtime netpoller 的手动版), //! 底层在 Linux 上使用 epoll。 // 获取 Unix 原始文件描述符所需的 trait // AsRawFd 提供了 as_raw_fd() 方法,用于从 std::io::Read/Write 等 Rust 抽象中 // 取出底层的 libc::c_int(POSIX 文件描述符),mio 注册 fd 监听时需要它 use std::os::unix::io::AsRawFd; // anyhow::Result 是一个简化的错误类型,等价于 Go 的 (T, error) // ? 操作符会将任何实现了 std::error::Error 的错误转换为 anyhow::Error use anyhow::Result; // clap::Parser 是一个 derive 宏,实现后 args.parse() 即可从 std::env::args() 解析 CLI 参数 // 类比 Go 的 flag.Parse(),但 clap 自动生成 --help 文本和错误处理 use clap::Parser; // mio::unix::SourceFd 是一个 bridge:将裸 fd 包装为实现 mio::Evented 的对象 // 这样 mio 的 epoll 可以监听任意 Unix fd,而不局限于 std::net::TcpStream 等标准类型 use mio::unix::SourceFd; // mio 是一个手动驱动的事件循环库(与 tokio 的异步运行时不同,mio 不调度 future) // - Poll:epoll/kqueue 的 Rust 封装,poll.poll() 会阻塞直到 fd 就绪 // - Interest:注册时的关注事件类型(READABLE / WRITABLE) // - Token:用户自定义的事件源标识(u64 包装),用于在 poll 返回时区分是哪个 fd 触发的 // - Events:poll 返回的事件集合(一个容量固定的 Vec) // 类比 Go runtime 的 netpoller,但 Go runtime 自动调度,mio 需要用户手动循环 use mio::{Events, Interest, Poll, Token}; // registry_queue_init 是 wayland-client 的便捷函数:连接到合成器并初始化全局注册表队列 // 它会在内部调用 Connection::connect_to_env() 并 roundtrip 一次拿到全局对象列表 use wayland_client::globals::registry_queue_init; // Connection 是与 Wayland 合成器的会话连接,封装了 Unix socket 的读写和协议解析 // 类比 Go 中的 net.Conn,但 Wayland 协议是有状态的消息流而非字节流 use wayland_client::Connection; // 各功能模块声明 mod args; // 命令行参数解析 mod avhw; // 音视频硬件加速 mod backend_detect; // 截屏后端自动检测(wlroots vs Portal/PipeWire) mod cap_portal; // XDG Portal 屏幕捕获 mod cap_wlr_screencopy; // wlroots wlr-screencopy 截屏协议 mod fps_limit; // 帧率限制器 mod state; // wlr-screencopy 后端的主状态机 mod state_portal; // Portal/PipeWire 后端的主状态机 mod stats; // 管道性能统计(卡顿诊断) mod transform; // 图像变换(旋转/翻转) mod webrtc; // WebRTC 传输(str0m Sans-IO) // 引入本 crate 内部模块,crate:: 前缀表示从 crate root 开始的绝对路径 // 类比 Go 中的 import "/args" 写法 use crate::args::Args; use crate::cap_wlr_screencopy::CapWlrScreencopy; use crate::state::EncConstructionStage; use crate::state::State; // mio 事件循环的 Token 标识:0 = Wayland 合成器事件,1 = 退出信号 const TOKEN_WAYLAND: Token = Token(0); const TOKEN_QUIT: Token = Token(1); /// 程序入口:解析参数 → 初始化日志 → 检测后端 → 启动对应的事件循环 /// /// 整体流程: /// 1. 解析命令行参数(分辨率、编码格式、帧率等) /// 2. 初始化日志系统(verbose 模式输出 DEBUG 级别,否则 INFO) /// 3. 检查编码格式(MVP 阶段仅支持 H.264) /// 4. 自动检测当前桌面环境支持的截屏后端 /// 5. 根据检测结果启动对应的事件循环: /// - wlroots 合成器(Sway/Hyprland)→ run_wlr_screencopy /// - GNOME/KDE 等 → run_portal_pipewire fn main() -> Result<()> { // 解析命令行参数 let args = Args::parse(); // 根据 verbose 模式或 RUST_LOG 环境变量设置日志级别 // 支持 RUST_LOG 粒度控制(如 RUST_LOG=wl_webrtc::webrtc=trace) // 详细解释: // - try_from_default_env() 返回 Result,读取 RUST_LOG 环境变量 // - unwrap_or_else(|_| {...}) 是 Result 的方法:成功则返回内部值,失败时调用闭包 // - |_| 是闭包参数语法:|参数| 表达式,单个 _ 表示忽略参数(这里是 Err 类型) // 类比 Go 的 if err != nil { fallback },但 Rust 用闭包传递 fallback 逻辑 let env_filter = tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| { if args.verbose { tracing_subscriber::EnvFilter::new("debug") } else { tracing_subscriber::EnvFilter::new("info") } }); // tracing_subscriber::fmt() 是 Builder 模式:链式调用配置,最后 .init() 消费 builder // 完成全局订阅注册。再次调用 .init() 会 panic,因此只能初始化一次。 // - with_env_filter: 设置过滤规则 // - with_writer: 设置日志输出目标(这里为 stderr,避免污染 stdout 用于视频流) // - init(): 消费 self,注册全局默认 subscriber,无返回值 tracing_subscriber::fmt() .with_env_filter(env_filter) .with_writer(std::io::stderr) .init(); tracing::info!("wl-webrtc starting"); tracing::debug!( "Args: output={:?} fps={} codec={} port={} verbose={}", args.output, args.fps, args.codec, args.port, args.verbose ); // MVP 阶段仅支持 H.264 编码,不支持 HEVC // anyhow::bail! 是一个宏(注意感叹号 !),立即返回 Err(anyhow::Error) // 类比 Go 的 fmt.Errorf("...") + return err,但是 Rust 用宏实现 if args.codec != "h264" { anyhow::bail!("HEVC not supported in MVP. Use --codec h264"); } if args.output.is_none() && args.port == 0 { anyhow::bail!("Either --output or --port is required"); } // 自动检测当前桌面环境可用的截屏后端 // 会尝试列举 Wayland 全局对象,判断合成器是否支持 wlr-screencopy 协议 // 行尾的 ? 是错误传播操作符:若 detect_backend 返回 Err,立即将该错误作为 fn main 的返回值 // 等价于 Go 的 if err != nil { return err },但 Rust 中 ? 适用于任何 Result/Option let backend = crate::backend_detect::detect_backend(&args)?; // 根据检测结果进入对应的事件循环 // match 是 Rust 的模式匹配表达式(类比 Go 的 switch 但更强大) // 每个 => 左侧是模式(这里是枚举变体),右侧是返回 Result<()> 的函数调用 // 由于 fn main 返回 Result<()>,这里直接把 match 表达式作为函数返回值(无分号 + 无 return) match backend { crate::backend_detect::CaptureBackend::WlrScreencopy => run_wlr_screencopy(args), crate::backend_detect::CaptureBackend::PortalPipeWire => run_portal_pipewire(args), } } /// 使用 wlroots wlr-screencopy 协议的事件循环(适用于 Sway、Hyprland 等 wlroots 合成器) /// /// 完整工作流程: /// 1. 连接 Wayland 合成器,获取全局注册表 /// 2. 初始化 wlr-screencopy 截屏状态机 /// 3. 获取 Wayland socket 的文件描述符 /// 4. 使用 mio 注册 fd 监听(Wayland 事件 + Unix 信号) /// 5. 进入主事件循环: /// - 监听合成器事件(帧就绪通知、输出信息等) /// - 定时请求截屏帧 /// - 将截取的帧编码为 H.264 并推流 /// 6. 收到退出信号或发生错误时,刷新编码器并退出 fn run_wlr_screencopy(args: Args) -> Result<()> { // Connect to Wayland compositor // 建立 Wayland 连接并初始化全局注册表 // 通过环境变量 $WAYLAND_DISPLAY 找到合成器的 Unix socket // 行尾的 ? 是 fn run_wlr_screencopy 内首次出现的错误传播操作符: // 若 connect_to_env 返回 Err,立即作为函数返回值向上抛出(类比 Go 的 return err) let conn = Connection::connect_to_env()?; // registry_queue_init 会绑定全局注册表回调, // 当合成器广播其全局对象(输出、截屏管理器等)时,State 会收到通知 // 返回值是元组 (GlobalManager, EventQueue),用 let 解构模式匹配赋值 // mut queue 表示 queue 在后续代码中会被修改(Rust 默认不可变,需 mut 显式声明) let (gm, mut queue) = registry_queue_init::>(&conn)?; let qhandle = queue.handle(); // State 是 wlr-screencopy 后端的核心状态机, // 内部管理输出探测、截屏请求、编码器构建、帧采集等阶段 let mut state = State::new(gm, args, qhandle)?; // Extract the Wayland fd and consume any immediately-available events. // prepare_read() flushes outgoing requests; read() pulls whatever the // compositor has already sent (may be EAGAIN if nothing yet). // 获取 Wayland socket 的文件描述符,并消费合成器已发送的事件 // 这个 fd 是后续 mio epoll 监听的对象,当合成器写入数据时变为可读 // 用 { ... } 块表达式将临时变量 guard 限制在作用域内,作用域结束自动 drop let wayland_fd = { let guard = queue .prepare_read() // ok_or_else 是 Option 的方法:None 时调用闭包生成 Err,得到 Result // || anyhow::anyhow!(...) 是无参数闭包语法(类比 JS 的 () => ...) // 行尾 ? 将 Result<_, Err> 解开为 Err 时立即从函数返回 .ok_or_else(|| anyhow::anyhow!("Failed to prepare Wayland read"))?; // 从 prepare_read 的 guard 中获取底层 socket 的原始文件描述符 let fd = guard.connection_fd().as_raw_fd(); // 尝试非阻塞读取合成器已发送但尚未消费的数据 // 如果没有数据会返回 EAGAIN,这里用 let _ 忽略 // let _ = expr 是显式忽略表达式返回值的惯用法,等价于 Go 的 _ = expr let _ = guard.read(); fd }; // 处理队列中的待处理事件 // 在进入主循环前,先处理注册表广播等初始化事件 // 此时状态机应进入 ProbingOutputs 阶段,正在探测可用的显示输出 queue.dispatch_pending(&mut state)?; tracing::info!( "Initial dispatch done, stage is ProbingOutputs: {}", matches!(state.stage, EncConstructionStage::ProbingOutputs { .. }) ); // 调试用:对 Wayland fd 做一次原始 poll,确认 fd 可读性 // 使用 libc 底层 poll 而非 mio,纯粹用于诊断初始化阶段的 fd 状态 { let mut pfd = libc::pollfd { fd: wayland_fd, events: libc::POLLIN, // 监听可读事件 revents: 0, }; // timeout=0 表示非阻塞,立即返回当前 fd 状态 // unsafe { ... } 是 Rust 的不安全块:内部调用 C 库 libc::poll,需要程序员 // 手动保证 &mut pfd 是有效的可变引用、fd 合法、不并发访问等不变量。 // unsafe 不关闭 Rust 借用检查,只是声明"我对外部 FFI 调用负责"。 let ret = unsafe { libc::poll(&mut pfd, 1, 0) }; tracing::info!( "Raw poll on wayland fd={wayland_fd}: ret={ret}, revents={}", pfd.revents ); } // Set up mio event loop // 使用 mio 创建事件循环,注册 Wayland fd 和 Unix 信号 // mio 底层在 Linux 上使用 epoll,macOS 上使用 kqueue let mut poll = Poll::new()?; // 事件缓冲区,容量 8 足够(实际只会同时处理 Wayland 事件和信号两种) let mut events = Events::with_capacity(8); // 将 Wayland socket fd 注册为可读监听 // 当合成器发送消息(如帧完成通知、配置变化)时,epoll 会唤醒 poll.registry().register( &mut SourceFd(&wayland_fd), TOKEN_WAYLAND, Interest::READABLE, )?; // 注册 SIGINT / SIGTERM 信号用于优雅退出 // signal_hook_mio 将 Unix 信号转换为 fd 可读事件, // 这样信号也可以通过 epoll 统一监听,不需要单独的信号处理器 let mut signals = signal_hook_mio::v1_0::Signals::new(&[ signal_hook::consts::SIGINT, // Ctrl+C signal_hook::consts::SIGTERM, // kill 命令默认信号 ])?; poll.registry() .register(&mut signals, TOKEN_QUIT, Interest::READABLE)?; tracing::info!("Event loop started"); // Flush outgoing before first poll iteration // 在首次 poll 前刷新所有待发送的 Wayland 请求 // 确保合成器能收到我们的初始化请求(如绑定全局对象、请求截屏等) conn.flush()?; // 主事件循环 // 这是 wlr-screencopy 后端的核心运行循环,负责: // - 接收合成器事件(截屏帧就绪、输出变化) // - 定时触发帧采集和编码 // - 响应退出信号 let mut running = true; while running { // 准备读取 Wayland 事件(非阻塞) // prepare_read() 会先刷出所有待发送的请求, // 然后进入"准备读取"状态,告诉合成器我们已准备好接收数据 let read_guard = queue.prepare_read(); // 如果无法 prepare_read(有待处理数据),先分发 // 返回 None 说明队列中已有待处理的合成器事件, // 需要先 dispatch 掉,否则新事件无法进入队列 if read_guard.is_none() { queue.dispatch_pending(&mut state)?; } // 阻塞等待事件,超时 100ms(用于帧率控制) // poll 会阻塞当前线程,直到以下任一条件满足: // 1. Wayland fd 可读(合成器发来了消息) // 2. 信号 fd 可读(收到了 SIGINT/SIGTERM) // 3. 超过 100ms 没有任何事件(超时返回,触发下一帧采集) // 100ms 超时 ≈ 10 FPS 的帧率上限 poll.poll(&mut events, Some(std::time::Duration::from_millis(100))) .unwrap_or_else(|e| { // EINTR 是信号中断,属于正常情况,继续循环 // 当进程收到信号时,阻塞中的 poll 会被中断并返回 EINTR, // 这不是错误,下一轮循环会继续正常 poll if e.kind() == std::io::ErrorKind::Interrupted { return; } tracing::error!("poll failed: {e}"); running = false; }); // 检查是否收到退出信号 // for x in &collection 是 Rust 的迭代语法,&events 表示借用 Events(不消费) // 类比 Go 的 for _, ev := range events {} for event in &events { if event.token() == TOKEN_QUIT { tracing::info!("Received quit signal"); running = false; } } // Wayland fd 可读时,读取并分发合成器事件 // 合成器可能发来多种事件:帧数据就绪、输出信息变化、协议错误等 // events.iter().any(|e| ...) 是迭代器方法,|e| 是单参数闭包 if events.iter().any(|e| e.token() == TOKEN_WAYLAND) { // if let Some(x) = opt 是 Option 的模式匹配简写(类比 Go 的 if v, ok := m[k]; ok) if let Some(guard) = read_guard { // match 是本函数内首次出现的多分支模式匹配 // Ok(_) 中下划线表示忽略成功值的具体内容(只关心成功/失败本身) match guard.read() { Ok(_) => { // 读取成功后,dispatch_pending 会将合成器事件 // 分发给 State 的对应回调方法处理 queue.dispatch_pending(&mut state)?; } Err(e) => { tracing::error!("Wayland read error: {e}"); running = false; } } } } // 请求状态机分配并编码一帧 // queue_alloc_frame 会根据当前状态机阶段执行不同操作: // - ProbingOutputs: 还在探测输出,跳过 // - AwaitingFrame: 已请求截屏,等待合成器回调 // - FrameReady: 有帧就绪,执行 DMA-BUF → H.264 编码 → 推流 // - Streaming: 正常采集中,请求下一帧 state.queue_alloc_frame(); state.poll_webrtc()?; // 状态机遇到致命错误时退出 if state.errored { tracing::error!("Fatal error in state machine (check preceding error logs), exiting"); running = false; } // 每轮循环结束前刷新 Wayland 发送缓冲区 // 将本轮回合中产生的所有 Wayland 请求(如截屏请求)发送给合成器 conn.flush()?; } // 关闭前刷新编码器,确保所有帧数据已写出 // 先通知帧率限制器停止,再刷新编码器缓冲区中残余的帧数据 tracing::info!("Shutting down, flushing encoder..."); state.fps_limit.flush(); // 仅在编码器已构建完成(Streaming 阶段)时才需要刷新 // if let 枚举变体模式匹配:Streaming { enc, .. } 解构出内部字段 enc,.. 忽略其他字段 // &mut state.stage 表示可变借用(类比 Go 的指针,但 Rust 编译期保证独占) if let crate::state::EncConstructionStage::Streaming { enc, .. } = &mut state.stage { // if let Err(e) = result 只关心失败分支,成功值用 _ 隐式忽略 if let Err(e) = enc.flush() { tracing::error!("Failed to flush encoder: {e}"); } } tracing::info!("Done"); Ok(()) } /// 使用 XDG Portal / PipeWire 后端的事件循环(适用于 KWin、KDE、GNOME 等桌面环境) /// /// 完整工作流程: /// 1. 初始化 Portal 状态机(内部通过 D-Bus 与 XDG Portal 通信) /// 2. 仅注册 Unix 信号监听(不需要监听 Wayland fd,帧数据由 PipeWire 投递) /// 3. 进入主事件循环: /// - 每 10ms 轮询一次 PipeWire 缓冲区 /// - 有帧数据时,取出并编码为 H.264 推流 /// - 收到退出信号时停止 /// 4. 退出时关闭 Portal 连接并释放 PipeWire 资源 fn run_portal_pipewire(args: Args) -> Result<()> { // 函数内 use 声明:将长路径名简化为局部短名,仅在该函数作用域内生效 // 类比 Go 函数内的局部 import 别名 use crate::state_portal::StatePortal; tracing::info!("Using Portal/PipeWire backend (KWin/KDE/GNOME)"); // StatePortal 初始化时会: // 1. 通过 D-Bus 连接到 XDG Portal 的 ScreenCast 接口 // 2. 请求用户授权屏幕录制权限 // 3. 建立 PipeWire 流连接,准备接收帧数据 // 行尾 ? 是本函数内首次出现的错误传播操作符:失败时立即从 fn run_portal_pipewire 返回 Err let mut state = StatePortal::new(args)?; // Set up signal handling only (no Wayland fd needed) // Portal 后端不需要监听 Wayland fd,只需处理 Unix 信号 // 因为帧数据是通过 PipeWire 独立投递的,不走 Wayland 协议 let mut signals = signal_hook_mio::v1_0::Signals::new(&[ signal_hook::consts::SIGINT, signal_hook::consts::SIGTERM, ])?; let mut poll = mio::Poll::new()?; let mut events = mio::Events::with_capacity(8); // 只注册信号 fd,没有 Wayland fd // 所以 poll.poll 在这里只负责检测 SIGINT/SIGTERM // 实际的帧采集完全依赖 poll_and_encode 的轮询 poll.registry() .register(&mut signals, mio::Token(1), mio::Interest::READABLE)?; // 主事件循环(非阻塞信号检测 + recv_timeout 等待帧) // poll 超时为 0ms(非阻塞),实际等待由 poll_and_encode 的 recv_timeout 实现 let mut running = true; while running { // poll 在此循环中只监听信号 fd(非阻塞): // - 收到 SIGINT/SIGTERM → 事件触发,设置 running=false // - 无事件 → 立即返回,继续执行 poll_and_encode(内部 recv_timeout 等待帧) poll.poll(&mut events, Some(std::time::Duration::from_millis(0))) .unwrap_or_else(|e| { if e.kind() == std::io::ErrorKind::Interrupted { return; } tracing::error!("poll failed: {e}"); running = false; }); for event in &events { if event.token() == mio::Token(1) { tracing::info!("Received quit signal"); running = false; } } // Process all available PipeWire frames // 处理所有可用的 PipeWire 帧数据 // poll_and_encode 会从 PipeWire 缓冲区取出帧, // 编码为 H.264 并推送。返回 true 表示还有更多帧待处理, // 返回 false 表示当前没有帧了,while 循环退出等待下一轮 poll // 外层 if 触发首次取帧(drain_first=true 表示允许阻塞等待), // 内层 while state.poll_and_encode(false)? {} 是空循环体语法: // 循环条件持续求值,只要返回 true 就重复,循环体 {} 不做额外事 if state.poll_and_encode(true)? { while state.poll_and_encode(false)? {} } // Portal 状态机遇到致命错误时退出 if state.is_errored() { tracing::error!( "Fatal error in portal state machine (check preceding error logs), exiting" ); running = false; } } tracing::info!("Shutting down..."); // 关闭 Portal 连接,释放 PipeWire 流和编码器资源 state.shutdown(); tracing::info!("Done"); Ok(()) }