Compare commits

...
Author SHA1 Message Date
dailz c772e4eb0b chore: bump MSRV to 1.87 to match actual API usage
CI / Build + Clippy + Test (push) Failing after 1h10m8s
CI / Security audit (RUSTSEC) (push) Failing after 1m31s
Oracle P2 follow-up. The clippy --fix autofixes earlier in this branch
silently introduced dependencies on APIs newer than the README's
1.70+ claim:
  - u32::is_multiple_of  (stable 1.87)
  - Option::is_none_or   (stable 1.82)

clippy::incompatible_msrv flagged the mismatch once rust-version was
pinned. Bumping the floor to 1.87 is the honest fix — the codebase
genuinely depends on 1.87 features now, and 1.87 has been stable
long enough (current stable is 1.96) that desktop CLI users on stable
Rust already have it.

  - Cargo.toml: rust-version '1.70' -> '1.87'. Comment lists the specific
    APIs that drove the bump and notes that further bumps need to be
    validated against clippy::incompatible_msrv.
  - README.md: Prerequisites line updated to 1.87+ with a brief why.
  - src/state_portal.rs: added the AsRawFd rustc-quirk comment that was
    already in avhw.rs (rustc emits a false 'unused_imports' warning;
    removing it produces E0599). Same known quirk, same documentation
    pattern.
  - src/transform.rs: fixed empty_line_after_doc_comments warning by
    converting the leading // doc-style comment to a //! module-level
    doc comment (which is what it should have been when I rewrote the
    file in commit 145b5d3).

All 79 unit tests + 3 integration tests pass. clippy: 0 errors,
0 incompatible_msrv warnings, 0 empty_line_after_doc_comments warnings.
Remaining warnings are: 1 AsRawFd rustc false-positive (documented),
5 unnecessary_cast FFI false-positives (rustc quirk on pointer casts),
and 8 dead-code items that need product decisions.
2026-06-28 14:39:00 +08:00
dailz 86a8b61b07 refactor: clear too_many_arguments and large_enum_variant warnings
Oracle P2 batch 2 (refactor items). Drops both remaining design-shape
clippy warnings to zero without behavior change.

  - avhw.rs build_filter_graph: drop unused _enc_width/_enc_height params
    (Oracle caught them during P2 review — passed by EncState::new but
    never read inside the function; the filter graph uses width/height
    only). Signature: 8 args -> 6 args (under clippy's 7 threshold).

  - state.rs InFlightSurface::CopyQueued: Box the drm_map field.
    AVDRMFrameDescriptor is ~592 bytes (4 objects + 4 layers); the enum
    size was being dominated by this variant, ballooning every
    InFlightSurface value to 592 bytes even for the None/AllocQueued
    variants. Box<AVDRMFrameDescriptor> shrinks the enum to ~32 bytes
    regardless of variant. The drm_map field is currently destructured
    under _drm_map (unused), so the boxing has no consumer-side impact.

  - state_portal.rs webrtc_thread_loop: 10 args -> 4 args via two new
    structs:
      * WebRtcThreadConfig { fps, enc_width, enc_height, max_bitrate }
        — immutable for the thread's lifetime; a tier change spawns a
        new thread rather than mutating.
      * WebRtcThreadChannels { webrtc_rx, sent_gap_tx, bitrate_tx,
        resolution_tx } — channel endpoints owned exclusively by the
        sender thread after spawn.
    wrtc (WebRtcState) and paused (Arc<AtomicBool>) stay as separate
    args because they have different ownership semantics (moved-in
    state vs shared atomic). Documented as doc comments on the new
    types so the next reader understands the bundle rationale.

All 79 unit tests + 3 integration tests pass. clippy: 0 errors.
Per-file warning counts: state_portal.rs down from 3 to 0; state.rs
down from 8 to 4 (remaining are unrelated dead-code on OutputInfo /
starting_timestamp).
2026-06-28 14:35:35 +08:00
dailz ed39d3d873 ci: add build/test/clippy gate + cargo audit; pin rust-version
Oracle P2 plan step 1+2+missed-fields. Locks in the audit cleanup so
future PRs can't regress the 0-errors / deny-unsafe / 79-tests baseline.

  - .github/workflows/ci.yml: two jobs on ubuntu-latest (Linux only —
    project is Wayland/VAAPI-specific, no macOS/Windows story).
      * build-test: installs ffmpeg + libavcodec/libavformat/libavutil/
        libswscale/libva dev + libwayland + libdrm + libpipewire-0.3-dev
        + libclang-dev/llvm-14 (LIBCLANG_PATH pinned); caches cargo +
        target; runs clippy -> build --release -> test --release.
        Release build before tests is mandatory because
        tests/integration_test.rs shells out to target/release/wl-webrtc.
      * audit: installs cargo-audit and runs 'cargo audit --deny warnings'
        as a separate job so a RUSTSEC advisory fails the build
        independently of compile state.
    No -D warnings on clippy yet — undocumented_unsafe_blocks is already
    deny via Cargo.toml; remaining warnings are advisory and can be
    tightened later.

  - Cargo.toml: pin rust-version = '1.70' to match README's claim.
    Without this, cargo builds silently on older toolchains and surfaces
    errors as cryptic parse failures instead of a clean version-mismatch
    message. Oracle flagged this as a missing field during P2 review.

License field intentionally omitted — repo has no LICENSE file and no
publication plan yet. Add when publication becomes a goal.

Verified locally: YAML parses, cargo build --release Finished in 9.82s,
cargo test --release 79 passed, cargo clippy 0 errors.
2026-06-28 14:31:42 +08:00
dailz a6560cff6c feat(stats): wire real scale/transfer/encode timing from EncState
Oracle step 4 (option A) — give the scale_*, transfer_*, encode_* stats
fields real producers instead of misleading zeros. The fields existed in
FrameTimings and PipelineStats already; producers just weren't passing
non-zero values.

  - avhw.rs: new EncodeStages { scale_us, transfer_us, encode_us } struct.
    EncState::encode_frame (HW VAAPI path) now times the filter graph
    separately from avcodec_send_frame, returning EncodeStages. transfer_us
    is honestly 0 because the HW path never reads back to CPU.
    SwEncState::encode_frame (SW fallback path) returns EncodeStages too;
    there import_and_scale bundles GPU scale + GPU→CPU readback into one
    call, so scale_us includes transfer for SW. Documented inline.

  - state.rs: StreamingEncoder::encode_frame return type bumps from
    Result<()> to Result<EncodeStages>; wlr-screencopy path now feeds
    real per-stage timings into FrameTimings instead of just total_us.

  - state_portal.rs: HW portal path (enc.encode_frame) now extracts
    stages.scale_us / stages.transfer_us / stages.encode_us into
    FrameTimings. Removed the now-unused t_encode_start binding.

Deferred (documented):
  - state_portal.rs SW portal path (line 525) calls import_and_scale +
    enc_thread separately and bypasses SwEncState::encode_frame. To wire
    scale/transfer timing there too, either route through SwEncState or
    thread timing out of import_and_scale. Out of scope for this commit.
  - SW path lumps transfer into scale_us. Splitting requires extending
    import_and_scale's return type — left as a follow-up if operational
    need arises (current default is HW VAAPI).

Oracle audit 2026-06-28 step 4 (option A: integrate, not delete).

All 79 unit tests + 3 integration tests pass. clippy: 0 errors.
2026-06-28 14:22:45 +08:00
dailz 2ac37a1dd1 fix(stats): wire PipeWire drops, expand Display, purge dead residue
Oracle-driven P1 fix plan. Resolves the StatsSnapshot 'computed but never
consumed' debt that was silently zeroing two real diagnostic fields and
leaving a dozen more unreported.

Bug fix (Oracle step 2):
  - state_portal.rs: set_pipewire_dropped(0, 0) and set_queue_depths(0, 0)
    were hardcoded, silently discarding real PipeWire diagnostics. Now wires
    to self.cap.dropped_count() (with pw_dropped_prev delta tracking) and
    self.cap.capture_queue_depth(). The encoded side stays 0 because the
    encoder thread exposes no queue-depth API.

Display expansion (Oracle step 1):
  - stats.rs: StatsSnapshot::Display now reports 12 previously-silent fields
    paired with their existing p95/max counterparts — capture/encoded/sent
    frame counts, elapsed_secs, *_avg_ms gap timing, frame_age_avg_ms,
    per-stage import/sws/encode/total avg_ms, output_frame_bytes_p95. Each
    line of the format string maps to one operational question (cadence,
    drops, queue pressure, latency, bandwidth); layout note added.

Dead residue purge (Oracle steps 5 + 6):
  - stats.rs: removed record_over_budget method + over_budget_count field
    (no caller; total_p95_ms answers the useful question without an
    arbitrary budget threshold).
  - state.rs: removed InFlightSurface::Allocd variant (never constructed)
    and CaptureSource::alloc_frame trait method (prototype leftover; the
    sole impl in cap_wlr_screencopy.rs returned None unconditionally).
  - cap_wlr_screencopy.rs: removed the alloc_frame stub; updated the
    unit-type Frame doc to reference the asynchronicity rationale without
    the deleted method.
  - cap_portal.rs: removed redundant 'let dropped = dropped;' shadowing
    flagged by clippy::redundant_locals (line 849).

Deferred (Oracle step 4 — needs product decision):
  - scale_avg/scale_p95/transfer_avg/transfer_p95/send_wait_p95 fields
    still appear in Display but producers in the live encode path don't
    record them, so they often show misleading zeros. Either add real
    EncState timing for scale/transfer stages, or remove the fields from
    Display until then.

All 79 unit tests + 3 integration tests still pass. clippy: 0 errors.
Warning count: multiple_fields_never_read on StatsSnapshot,
method_never_used on record_over_budget/dropped_count/capture_queue_depth/
alloc_frame, variant_never_constructed on Allocd, redundant_locals on
dropped — all gone.
2026-06-28 14:15:55 +08:00
dailz 145b5d3e7e chore: design cleanup, dead-code purge, README/doc refresh
Audit-driven follow-up after the SAFETY-debt commit (Oracle steps 6-7).
End state: cargo clippy --release --all-targets still 0 errors; private_interfaces
and type_complexity warnings cleared.

Design cleanups (Oracle step 6):
  - cap_portal.rs: introduce PortalFormatInfo struct to replace the
    Rc<Cell<Option<(u32,u32,u32,u64)>>> cross-callback hand-off. Self-
    documenting struct fields replace positional tuple access at the
    format-change and process callbacks.
  - avhw.rs: import_dma_buf_to_vaapi signature collapses from 8 args
    (fd/width/height/drm_format/modifier/stride/offset) to
    (*mut AVBufferRef, &PwDmaBufFrame). Callers in avhw.rs,
    state_portal.rs, and vaapi_import_bench.rs now pass the frame by
    reference instead of unpacking 7 fields just to repack them. Drops
    the unused width parameter and the too_many_arguments(8/7) warning.
  - state.rs: visibility hygiene. EncConstructionStage and WlrHeadInfo
    downgrade pub -> pub(crate); State.stage field downgrades to
    pub(crate). These are internal state-machine types not exposed
    across the crate boundary; making them pub(crate) clears all
    private_interfaces warnings without leaking more types.

Dead-code purge (Oracle step 7):
  - transform.rs: remove unused Rect struct, transform_basis,
    screen_to_frame, fit_inside_bounds helpers and their 18 dedicated
    tests. Transform enum and transpose_if_transform_transposed remain
    (both are actively used by state.rs and avhw.rs). File shrinks
    from 409 -> 109 lines.

Repository housekeeping (Oracle step 7):
  - .gitignore: add review.json (stray review-tool output that
    regenerates per run).
  - README.md: refresh CLI table to match src/args.rs (now lists
    --backend, --no-persist, --port-as-WebRTC-signaling, --max-bitrate,
    --stats). Add capture-backend explainer + 4 new usage examples.
    Note in README points readers at src/args.rs as the authoritative
    source. Remove stale 'WebTransport, unused in MVP' description.

avhw.rs: AsRawFd import annotated with a rustc-quirk explanation — the
import triggers a false 'unused_imports' warning but E0599 if removed.
Left as-is with explanatory comment rather than chasing the lint.

All 79 remaining unit tests + 3 integration tests still pass. Cargo
build --release clean.
2026-06-28 13:58:52 +08:00
dailz 30f8fe51f2 chore: clear clippy errors, document all unsafe blocks, deny new SAFETY debt
Audit-driven cleanup pass. End state:
  - cargo clippy --release --all-targets: 0 errors (was 4)
  - undocumented_unsafe_blocks warnings: 0 (was 67)
  - Cargo.toml: undocumented_unsafe_blocks escalated warn -> deny

Clippy correctness errors fixed:
  - src/bin/{sw_encode_bench,vaapi_import_bench}.rs: receive_first_frame
    rewritten per Oracle plan with total 10s deadline + 200ms wait slice +
    while-let drain of all control events. The previous loop body always
    exited on first iteration (never_loop); the new version actually retries
    and matches production's repeated-poll semantics in state_portal.rs.
  - src/avhw.rs: hash_sampled_y_plane tests now use a row_range(row, stride,
    width) helper instead of inline stride * N. Preserves the row-index
    intent across all sibling tests without tripping erasing_op (row==0) or
    identity_op (row==1).

Machine-applicable clippy autofixes applied via 'cargo clippy --fix':
  - unnecessary_cast, manual_is_multiple_of, needless_borrows_for_generic_args
  - manual_abs_diff, derivable_impls, new_without_default
  - unnecessary_map_or, unneeded_struct_pattern, redundant_locals

webrtc_gop_formula test rewritten to wrap the (fps * 2).max(20) formula in
a runtime lambda. The previous clippy --fix pass had constant-folded the
5fps case into assert_eq!(20, 20), silently stripping the floor-case
coverage. The lambda blocks the fold while keeping the formula exercisable.

67 SAFETY comments added across 7 files (cap_portal.rs 26, sw_encode_bench
21, state_portal.rs 7, vaapi_import_bench.rs 6, avhw.rs 5, state.rs 1,
main.rs 1). Two sites carry load-bearing invariant documentation:
  - cap_portal.rs:806 process callback documents the PipeWire raw_buf
    ownership contract across all 10 exit paths (audited: every path
    correctly requeues; fd ownership via dup() is independent and also
    exactly-once closed).
  - avhw.rs:341 unsafe impl Send for EncState documents the single-thread
    exclusivity assumption referenced by AGENTS.md.

All 97 unit tests + 3 integration tests still pass; cargo build --release
finishes clean. Lint escalation to deny freezes the SAFETY baseline: any
future patch adding an unsafe block without a // SAFETY: comment will fail
clippy at compile time.
2026-06-28 13:44:27 +08:00
15 changed files with 665 additions and 525 deletions
+89
View File
@@ -0,0 +1,89 @@
# Continuous integration for wl-webrtc.
#
# Triggered on push/PR to master. Runs the full quality gate that the recent
# audit baselined:
# - clippy: 0 errors (undocumented_unsafe_blocks is deny in Cargo.toml; other
# warnings are advisory for now).
# - build --release: integration tests in tests/integration_test.rs shell out
# to target/release/wl-webrtc, so the release binary must exist before tests
# run.
# - test --release: 79 unit + 3 integration; the 1 hardware-ignored test
# stays ignored in CI (needs Wayland session + VAAPI GPU).
# - cargo audit: separate job so a RUSTSEC advisory fails the build without
# conflating with compile errors.
#
# The job pins Linux only — the project is Wayland/VAAPI-specific and has no
# macOS/Windows story. Oracle audit 2026-06-28 P2 plan.
name: CI
on:
push:
branches: [master]
pull_request:
branches: [master]
env:
CARGO_TERM_COLOR: always
# Build dependencies match shell.nix + README Prerequisites section.
LIBCLANG_PATH: /usr/lib/llvm-14/lib
jobs:
build-test:
name: Build + Clippy + Test
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Install Rust toolchain (stable)
uses: dtolnay/rust-toolchain@stable
with:
components: clippy
- name: Install system dependencies
run: |
sudo apt-get update
sudo apt-get install -y --no-install-recommends \
ffmpeg \
libavcodec-dev libavformat-dev libavutil-dev libswscale-dev libva-dev \
libwayland-dev wayland-protocols \
libdrm-dev \
libpipewire-0.3-dev \
libclang-dev llvm-14
- name: Cache cargo registry + build artifacts
uses: actions/cache@v4
with:
path: |
~/.cargo/registry
~/.cargo/git
target
key: ${{ runner.os }}-cargo-${{ hashFiles('Cargo.lock', 'Cargo.toml') }}
restore-keys: |
${{ runner.os }}-cargo-
- name: Clippy (release, all targets)
run: cargo clippy --release --all-targets
- name: Build release (required before tests)
run: cargo build --release --all-targets
- name: Test (release)
run: cargo test --release
audit:
name: Security audit (RUSTSEC)
runs-on: ubuntu-latest
# Keep separate from build-test so a vulnerability advisory fails the
# check independently of compile state.
steps:
- uses: actions/checkout@v4
- name: Install Rust toolchain (stable)
uses: dtolnay/rust-toolchain@stable
- name: Install cargo-audit
run: cargo install cargo-audit --locked
- name: Audit dependencies
run: cargo audit --deny warnings
+3
View File
@@ -21,3 +21,6 @@ Thumbs.db
.playwright-mcp/
wl-webrtc.log
webrtc-p0-success.png
# Stray review-tool output (regenerated per review run)
review.json
+7 -1
View File
@@ -2,6 +2,12 @@
name = "wl-webrtc"
version = "0.1.0"
edition = "2021"
# MSRV pinned to 1.87 to match the actual API floor — the codebase uses
# u32::is_multiple_of (1.87) and Option::is_none_or (1.82) introduced by
# clippy autofixes. README's Prerequisites section mirrors this. Bumping
# this floor requires checking clippy::incompatible_msrv against the new
# value.
rust-version = "1.87"
description = "Wayland screen capture and encoding tool"
[dependencies]
@@ -33,4 +39,4 @@ dirs = "6"
tempfile = "3.27.0"
[lints.clippy]
undocumented_unsafe_blocks = "warn"
undocumented_unsafe_blocks = "deny"
+30 -3
View File
@@ -4,7 +4,7 @@ Wayland screen capture and encoding tool.
## Prerequisites
- **Rust toolchain** (1.70+): `rustup default stable`
- **Rust toolchain** (1.87+; MSRV pinned to match `u32::is_multiple_of` / `Option::is_none_or` usage): `rustup default stable`
- **FFmpeg 6.0+** dev libraries with VAAPI support:
- Arch: `pacman -S ffmpeg`
- Ubuntu/Debian: `apt install libavcodec-dev libavformat-dev libavutil-dev libswscale-dev libva-dev`
@@ -38,19 +38,46 @@ wl-webrtc --output output.mp4 --drm-device /dev/dri/renderD128
# Verbose mode
wl-webrtc --output output.mp4 -v
# WebRTC streaming mode (HTTP signaling server)
wl-webrtc --port 8080 -v
# Force a fresh portal authorization dialog (ignore saved restore token)
wl-webrtc --output output.mp4 --no-persist
# Pin the capture backend instead of auto-detecting
wl-webrtc --output output.mp4 --backend portal # or: --backend screencopy
```
## CLI Arguments
> `src/args.rs` is the authoritative source. Run `wl-webrtc --help` for the live list.
| Argument | Default | Description |
|---|---|---|
| `-o`, `--output` | (required) | Output file path (e.g., output.mp4) |
| `-o`, `--output` | (optional) | Output file path (e.g. output.mp4). Optional when using `--port` for WebRTC mode. |
| `--output-name` | auto | Wayland output name to capture |
| `--fps` | 30 | Target frames per second |
| `--codec` | h264 | Video codec (h264 only for MVP) |
| `--hw-accel` | vaapi | Hardware acceleration method |
| `--drm-device` | auto | DRM render device path |
| `--bitrate` | auto | Target bitrate in bps |
| `--max-bitrate` | 8000000 | Max bitrate cap for WebRTC mode (caps BWE escalation; no effect in MP4 mode) |
| `--gop-size` | auto | Group of Pictures size |
| `-v`, `--verbose` | false | Enable verbose logging |
| `--port` | 0 | WebTransport server port (unused in MVP) |
| `--backend` | auto | Capture backend: `screencopy` (wlroots) or `portal` (KWin/KDE). Auto-detected if omitted. |
| `--port` | 0 | WebRTC HTTP signaling server port. `0` keeps MP4 file output mode. |
| `--no-persist` | false | Force re-authorization (ignore saved portal restore token) |
| `--stats` | false | Print per-second pipeline statistics for stutter diagnosis |
## Capture backends
The tool supports two Wayland capture backends, auto-detected by default:
- **wlr-screencopy** (preferred when `zwlr_screencopy_manager_v1` is advertised):
works on wlroots-based compositors (Sway, Hyprland, etc.).
- **XDG Portal / PipeWire** (fallback when D-Bus ScreenCast is available):
works on KWin/KDE and any compositor that implements the XDG Desktop Portal
screen-cast protocol. The first run shows an authorization dialog; a restore
token is cached under `wl-webrtc/portal-restore-token` so subsequent runs
don't re-prompt (use `--no-persist` to force a fresh authorization).
+122 -48
View File
@@ -1,6 +1,9 @@
use std::ffi::CString;
use std::mem;
use std::os::fd::{AsRawFd, RawFd};
// AsRawFd is required by `frame.fd.as_raw_fd()` below but rustc emits a false
// "unused_imports" warning because OwnedFd also has an inherent `as_raw_fd`.
// E0599 if removed → must stay; warning is a known rustc quirk.
use std::os::fd::AsRawFd;
use std::os::raw::c_void;
use std::path::Path;
use std::ptr;
@@ -189,6 +192,18 @@ impl Drop for AvHwFrameCtx {
}
}
/// Per-stage timing breakdown for one encode cycle on the hardware path.
/// Returned by [`EncState::encode_frame`] so callers can fold the numbers
/// into [`crate::stats::FrameTimings`]. `transfer_us` is always 0 on the HW
/// path because the frame stays on the GPU; the SW path's struct (if added
/// later) would carry a real readback measurement.
#[derive(Debug, Default, Clone, Copy)]
pub struct EncodeStages {
pub scale_us: u64,
pub transfer_us: u64,
pub encode_us: u64,
}
/// Test whether `drm_device` can import the PipeWire DMA-BUF frame via VAAPI.
pub fn test_dma_buf_import(drm_device: &Path, frame: &PwDmaBufFrame) -> Result<()> {
let hw_dev = AvHwDevCtx::new_vaapi(drm_device)?;
@@ -197,16 +212,7 @@ pub fn test_dma_buf_import(drm_device: &Path, frame: &PwDmaBufFrame) -> Result<(
// SAFETY: frames is a live VAAPI frames context; frame carries valid DMA-BUF metadata.
unsafe {
import_dma_buf_to_vaapi(
frames.as_ptr(),
frame.fd.as_raw_fd(),
frame.width,
frame.height,
frame.format,
frame.modifier,
frame.stride,
frame.offset,
)
import_dma_buf_to_vaapi(frames.as_ptr(), frame)
}?;
Ok(())
@@ -215,19 +221,22 @@ pub fn test_dma_buf_import(drm_device: &Path, frame: &PwDmaBufFrame) -> Result<(
/// Import a DMA-BUF into a VAAPI hardware frame via zero-copy `av_hwframe_map`.
///
/// # Safety
/// Imports a DMA-BUF frame into a VAAPI hardware frame pool for GPU-side processing.
///
/// Takes the negotiated format/geometry from `frame` (a `PwDmaBufFrame` from
/// PipeWire capture) plus the target `frames_ctx` (VAAPI frame pool from
/// `AvHwFrameCtx`) and returns an `ff::frame::Video` whose data[3] points to
/// the hardware frame.
///
/// # Safety
///
/// - `frames_ctx` must point to an initialized AVHWCramesContext for VAAPI
/// - `raw_fd` must be a valid DMA-BUF file descriptor
/// - `frame.fd` must be a valid DMA-BUF file descriptor
pub unsafe fn import_dma_buf_to_vaapi(
frames_ctx: *mut ffi::AVBufferRef,
raw_fd: RawFd,
width: u32,
height: u32,
drm_format: u32,
modifier: u64,
stride: u32,
offset: u64,
frame: &PwDmaBufFrame,
) -> Result<ff::frame::Video> {
let duped_fd = libc::dup(raw_fd);
let duped_fd = libc::dup(frame.fd.as_raw_fd());
if duped_fd < 0 {
bail!("dup(fd) failed: {}", std::io::Error::last_os_error());
}
@@ -235,14 +244,14 @@ pub unsafe fn import_dma_buf_to_vaapi(
let mut desc: ffi::AVDRMFrameDescriptor = mem::zeroed();
desc.nb_objects = 1;
desc.objects[0].fd = duped_fd;
desc.objects[0].size = (height as usize) * (stride as usize);
desc.objects[0].format_modifier = modifier;
desc.objects[0].size = (frame.height as usize) * (frame.stride as usize);
desc.objects[0].format_modifier = frame.modifier;
desc.nb_layers = 1;
desc.layers[0].format = drm_format;
desc.layers[0].format = frame.format;
desc.layers[0].nb_planes = 1;
desc.layers[0].planes[0].object_index = 0;
desc.layers[0].planes[0].offset = offset as isize;
desc.layers[0].planes[0].pitch = stride as isize;
desc.layers[0].planes[0].offset = frame.offset as isize;
desc.layers[0].planes[0].pitch = frame.stride as isize;
let desc_box = Box::new(desc);
let desc_ptr = Box::into_raw(desc_box);
@@ -264,8 +273,8 @@ pub unsafe fn import_dma_buf_to_vaapi(
{
let sp = src.as_mut_ptr();
(*sp).format = ffi::AVPixelFormat::AV_PIX_FMT_DRM_PRIME as i32;
(*sp).width = width as i32;
(*sp).height = height as i32;
(*sp).width = frame.width as i32;
(*sp).height = frame.height as i32;
(*sp).data[0] = (*buf_ref).data;
(*sp).buf[0] = buf_ref;
}
@@ -338,6 +347,14 @@ pub struct EncState {
frames_written: bool,
}
// SAFETY: EncState is moved to exactly one thread (the encode worker) and used
// exclusively there. All fields are either plain Copy types (Option<i64>, bool)
// or ffmpeg-next / AvHw* owned wrappers whose raw inner pointers are not actually
// shared across threads — they're touched only from the owning encode thread.
// This impl exists only to satisfy Rust's auto-Send inference (which can't see
// through the raw pointers hidden inside the wrappers). Do NOT add fields that
// introduce shared mutable state without re-auditing this assumption; see
// AGENTS.md "Unsafe and FFI work" for the documented exclusivity requirement.
unsafe impl Send for EncState {}
impl EncState {
@@ -374,8 +391,6 @@ impl EncState {
&frames_rgb,
width,
height,
enc_width,
enc_height,
fps,
transform,
)?;
@@ -430,6 +445,9 @@ impl EncState {
// VBV rate limiting: caps IDR burst size for WebRTC. Without this a 4K
// scene change can produce a 256KB keyframe that overflows the UDP send
// buffer. bufsize=bitrate/4 ≈ 250ms of video at the target bitrate.
// SAFETY: enc.as_mut_ptr() is a valid AVCodecContext for the not-yet-opened
// encoder. rc_max_rate and rc_buffer_size are plain integer fields; assigning
// i64/i32 values is a simple struct-field write on a properly aligned pointer.
unsafe {
let ctx_ptr = enc.as_mut_ptr();
(*ctx_ptr).rc_max_rate = bitrate as i64;
@@ -456,6 +474,10 @@ impl EncState {
{
let key = CString::new("repeat_pps").unwrap();
let val = CString::new("1").unwrap();
// SAFETY: enc is a valid AVCodecContext for the not-yet-opened encoder;
// priv_data is the codec's private options struct. key/val are NUL-terminated
// CString that live across the call. av_opt_set is FFmpeg's standard
// option-setter. Failure is non-fatal (returns < 0 on older FFmpeg).
let ret = unsafe {
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0)
};
@@ -488,11 +510,16 @@ impl EncState {
bail!("Failed to allocate output format context: {}", ff_err(ret));
}
// SAFETY: avformat_query_codec checks codec+format compatibility.
let codec_id = unsafe { (*enc_video.as_ptr()).codec_id };
let oformat = unsafe { (*fmt_ctx_ptr).oformat };
// SAFETY: enc_video is a valid AVCodecContext pointer; codec_id is a plain
// i32 enum discriminant read from it. fmt_ctx_ptr is a valid AVFormatContext
// allocated above; oformat is a const pointer field read from it.
// avformat_query_codec checks codec+format compatibility; both pointers are
// valid and FF_COMPLIANCE_NORMAL is a constant. All three reads happen in one
// block so a single SAFETY rationale covers them.
let compat = unsafe {
ffi::avformat_query_codec(oformat, codec_id, ffi::FF_COMPLIANCE_NORMAL as i32)
let codec_id = (*enc_video.as_ptr()).codec_id;
let oformat = (*fmt_ctx_ptr).oformat;
ffi::avformat_query_codec(oformat, codec_id, ffi::FF_COMPLIANCE_NORMAL)
};
if compat < 0 {
bail!("H.264 codec not supported by output container format");
@@ -560,7 +587,7 @@ impl EncState {
&self.frames_rgb
}
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<()> {
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<EncodeStages> {
let mut filter_src_ctx = self
.video_filter
.get("in")
@@ -572,11 +599,18 @@ impl EncState {
.ok_or_else(|| anyhow::anyhow!("filter 'out' not found"))?;
let mut filter_sink = filter_sink_ctx.sink();
// Scale stage = filter graph push + pull (scale_vaapi for resolution
// change + format conversion to NV12). Timed separately from the
// actual avcodec_send_frame so the per-stage stats answer "where is
// latency?" honestly. See Oracle audit 2026-06-28 step 4.
let scale_start = Instant::now();
// SAFETY: hw_frame is a valid VAAPI hardware frame from capture.
filter_src
.add(hw_frame)
.map_err(|e| anyhow::anyhow!("Filter source add failed: {e}"))?;
let mut scale_us = 0u64;
let mut encode_us = 0u64;
loop {
let mut filtered = ff::frame::Video::empty();
match filter_sink.frame(&mut filtered) {
@@ -588,6 +622,11 @@ impl EncState {
Err(ff::Error::Other { errno }) if errno == ffi::EAGAIN => break,
Err(e) => bail!("Filter sink get frame failed: {e}"),
}
// First successful pull closes the scale-stage measurement; later
// pulls (rare extras) roll into encode time.
if scale_us == 0 {
scale_us = scale_start.elapsed().as_micros() as u64;
}
let pts = filtered.pts().unwrap_or(0);
if self.starting_timestamp.is_none() {
@@ -595,6 +634,7 @@ impl EncState {
}
let start_ts = self.starting_timestamp.unwrap();
let encode_start = Instant::now();
// SAFETY: avcodec_send_frame sends a valid NV12 VAAPI surface to the encoder.
let ret =
unsafe { ffi::avcodec_send_frame(self.enc_video.as_mut_ptr(), filtered.as_ptr()) };
@@ -602,9 +642,15 @@ impl EncState {
bail!("avcodec_send_frame failed: {}", ff_err(ret));
}
self.drain_encoder(start_ts)?;
encode_us += encode_start.elapsed().as_micros() as u64;
}
Ok(())
Ok(EncodeStages {
scale_us,
// HW path stays on GPU — no CPU readback, transfer is N/A.
transfer_us: 0,
encode_us,
})
}
pub fn flush(&mut self) -> Result<()> {
@@ -1225,7 +1271,7 @@ impl SwEncEncode {
let frame_index = self.frame_count;
self.frame_count = self.frame_count.saturating_add(1);
let current_hash = hash_sampled_y_plane(&frame.y_data, width, height, frame.y_stride);
let force_gop_frame = self.gop_size > 0 && frame_index % u64::from(self.gop_size) == 0;
let force_gop_frame = self.gop_size > 0 && frame_index.is_multiple_of(u64::from(self.gop_size));
if frame_index > 0 && !force_gop_frame && !force_this_frame && current_hash == self.last_frame_hash {
tracing::debug!(frame_index, "skipping duplicate frame");
self.last_frame_hash = current_hash;
@@ -1545,9 +1591,25 @@ impl SwEncState {
self.import.frames_rgb()
}
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<()> {
pub fn encode_frame(&mut self, hw_frame: &ff::frame::Video) -> Result<EncodeStages> {
// SW path: import_and_scale bundles GPU filter graph (scale) + GPU→CPU
// readback (transfer) into one call. Timing them separately requires
// extending import_and_scale's signature; for now both roll into
// scale_us and transfer_us stays 0 with this comment as the honest
// statement. Oracle audit 2026-06-28 step 4.
let scale_start = Instant::now();
let cpu_frame = self.import.import_and_scale(hw_frame)?;
self.encode.encode_cpu_frame(&cpu_frame).map(|_| ())
let scale_us = scale_start.elapsed().as_micros() as u64;
let encode_start = Instant::now();
self.encode.encode_cpu_frame(&cpu_frame)?;
let encode_us = encode_start.elapsed().as_micros() as u64;
Ok(EncodeStages {
scale_us,
transfer_us: 0,
encode_us,
})
}
pub fn flush(&mut self) -> Result<()> {
@@ -1760,7 +1822,7 @@ fn create_software_h264_muxer(
let key = CString::new("threads").unwrap();
let val = CString::new("6").unwrap();
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH as i32;
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH;
// SAFETY: enc is a valid, initialized AVCodecContext from
// avcodec_alloc_context3. Setting level is a simple i32 field
// assignment on a properly aligned struct.
@@ -1896,7 +1958,7 @@ fn create_software_h264_encoder(
ffi::av_opt_set((*enc.as_mut_ptr()).priv_data, key.as_ptr(), val.as_ptr(), 0);
// High profile via AVCodecContext.profile (not x264opts — x264 rejects it there).
// High enables CABAC + 8x8dct automatically.
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH as i32;
(*enc.as_mut_ptr()).profile = ffi::AV_PROFILE_H264_HIGH;
// SAFETY: enc is a valid, initialized AVCodecContext from
// avcodec_alloc_context3. Setting level is a simple i32 field
// assignment on a properly aligned struct.
@@ -1941,8 +2003,6 @@ fn build_filter_graph(
frames_rgb: &AvHwFrameCtx,
width: u32,
height: u32,
_enc_width: u32,
_enc_height: u32,
fps: u32,
transform: Transform,
) -> Result<ff::filter::Graph> {
@@ -2042,6 +2102,14 @@ fn build_filter_graph(
mod tests {
use super::*;
// Centralizes the `stride * row` byte-offset pattern used by the Y-plane hash
// tests below, so clippy::erasing_op (row == 0) and clippy::identity_op (row == 1)
// both pass without sacrificing the row-index intent the tests are written around.
fn row_range(row: usize, stride: usize, width: usize) -> std::ops::Range<usize> {
let start = stride * row;
start..start + width
}
// ── Task 1: VBV x264opts formatting ──
#[test]
@@ -2071,10 +2139,16 @@ mod tests {
#[test]
fn webrtc_gop_formula() {
assert_eq!((15u32 * 2).max(20), 30); // 15fps -> 30
assert_eq!((30u32 * 2).max(20), 60); // 30fps -> 60
assert_eq!((60u32 * 2).max(20), 120); // 60fps -> 120
assert_eq!((5u32 * 2).max(20), 20); // 5fps -> 20 (floor)
// Formula under test: GOP = max(fps * 2, 20). Hid behind a runtime lambda so
// clippy can't constant-fold the assertions into tautologies (which would
// silently strip the floor-case coverage for 5fps).
fn gop(fps: u32) -> u32 {
(fps * 2).max(20)
}
assert_eq!(gop(15), 30); // 15fps -> 30
assert_eq!(gop(30), 60); // 30fps -> 60
assert_eq!(gop(60), 120); // 60fps -> 120
assert_eq!(gop(5), 20); // 5fps -> 20 (floor)
}
#[test]
@@ -2122,7 +2196,7 @@ mod tests {
let y_data1 = vec![0u8; stride * height];
let mut y_data2 = vec![0u8; stride * height];
// Row 1 is NOT sampled (sampling is every 8th row: 0, 8, 16, ...)
y_data2[stride * 1..stride * 1 + width].fill(255);
y_data2[row_range(1, stride, width)].fill(255);
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_eq!(
@@ -2140,7 +2214,7 @@ mod tests {
let y_data1 = vec![0u8; stride * height];
let mut y_data2 = vec![0u8; stride * height];
// Row 0 IS sampled (every 8th row starting from 0)
y_data2[stride * 0..stride * 0 + width].fill(255);
y_data2[row_range(0, stride, width)].fill(255);
let hash1 = hash_sampled_y_plane(&y_data1, width, height, stride);
let hash2 = hash_sampled_y_plane(&y_data2, width, height, stride);
assert_ne!(
+73 -9
View File
@@ -62,22 +62,32 @@ fn pix_fmt(p: ff::format::Pixel) -> ffi::AVPixelFormat {
}
fn receive_first_frame(cap: &CapPortal) -> Result<wl_webrtc::cap_portal::PwDmaBufFrame> {
// Drain-and-wait loop that mirrors production's repeated-poll semantics
// (state_portal.rs::poll_and_encode driven by main.rs's outer loop), but with
// a single bounded 10s total deadline appropriate for a bench tool. Unlike a
// single 10s blocking wait, this loop actually iterates: each turn drains ALL
// pending control events (the ctrl channel is bounded to 8 — a single
// if-let would silently miss backlog) and then waits a short slice for a
// frame, so StreamEnded/Error arriving mid-wait are observed within ~200ms.
const TOTAL_DEADLINE: std::time::Duration = std::time::Duration::from_secs(10);
const WAIT_SLICE: std::time::Duration = std::time::Duration::from_millis(200);
let deadline = Instant::now() + TOTAL_DEADLINE;
loop {
if let Ok(ctrl) = cap.event_receiver().try_recv() {
while let Ok(ctrl) = cap.event_receiver().try_recv() {
match ctrl {
PwCtrlEvent::StreamEnded => bail!("PipeWire stream ended before first frame"),
PwCtrlEvent::FormatChanged { .. } => {}
PwCtrlEvent::Error(e) => bail!("PipeWire error: {e}"),
}
}
match cap
.frame_receiver()
.recv_timeout(std::time::Duration::from_secs(10))
{
let remaining = match deadline.checked_duration_since(Instant::now()) {
Some(r) if !r.is_zero() => r,
_ => bail!("Timeout waiting for first frame (10s)"),
};
let slice = remaining.min(WAIT_SLICE);
match cap.frame_receiver().recv_timeout(slice) {
Ok(frame) => return Ok(frame),
Err(crossbeam_channel::RecvTimeoutError::Timeout) => {
bail!("Timeout waiting for first frame (10s)");
}
Err(crossbeam_channel::RecvTimeoutError::Timeout) => continue,
Err(crossbeam_channel::RecvTimeoutError::Disconnected) => {
bail!("PipeWire frame channel disconnected");
}
@@ -142,6 +152,9 @@ fn main() -> Result<()> {
println!("[3/4] Testing mmap on DMA-BUF...");
let mmap_size = (src_stride as usize) * (src_height as usize);
// SAFETY: first_frame.fd is an open DMA-BUF; offset/size come from PipeWire's
// negotiated format. PROT_READ+MAP_SHARED is the standard read-only DMA-BUF
// mapping. Returns MAP_FAILED on error (checked below).
let mmap_ptr = unsafe {
libc::mmap(
ptr::null_mut(),
@@ -173,6 +186,8 @@ fn main() -> Result<()> {
"[3/4] mmap SUCCESS — CPU can read DMA-BUF ({:.1} MB)\n",
mmap_size as f64 / 1024.0 / 1024.0
);
// SAFETY: mmap_ptr was returned by mmap above and is not MAP_FAILED (checked);
// mmap_size matches the original mapping. POSIX munmap(2) releases the mapping.
unsafe {
libc::munmap(mmap_ptr, mmap_size);
}
@@ -205,6 +220,9 @@ fn main() -> Result<()> {
let codec_name = codec.name();
if codec_name == "libx264" {
// SAFETY: enc is a valid AVCodecContext for the not-yet-opened encoder;
// priv_data is the x264 private options struct. All CStrings live across
// both av_opt_set calls. These set the x264 "preset" and "tune" options.
unsafe {
let key = CString::new("preset").unwrap();
let val = CString::new("veryfast").unwrap();
@@ -220,6 +238,8 @@ fn main() -> Result<()> {
// Create output format context via FFI
let mut fmt_ctx_ptr: *mut ffi::AVFormatContext = ptr::null_mut();
// SAFETY: fmt_ctx_ptr is an out-parameter initialized by FFmpeg; output_cstr
// lives across the call. Returns 0 on success; we check below.
let ret = unsafe {
ffi::avformat_alloc_output_context2(
&mut fmt_ctx_ptr,
@@ -232,21 +252,30 @@ fn main() -> Result<()> {
bail!("Failed to allocate output format context: error {ret}");
}
// SAFETY: fmt_ctx_ptr is the valid output context allocated above.
// avformat_new_stream returns a pointer to a new AVStream or NULL on failure.
let stream_ptr = unsafe { ffi::avformat_new_stream(fmt_ctx_ptr, ptr::null()) };
if stream_ptr.is_null() {
bail!("Failed to create new stream");
}
// SAFETY: stream_ptr and enc_video.as_ptr() are valid pointers; codecpar is
// the output destination inside stream. avcodec_parameters_from_context copies
// encoder parameters into the stream's codecpar.
let ret =
unsafe { ffi::avcodec_parameters_from_context((*stream_ptr).codecpar, enc_video.as_ptr()) };
if ret < 0 {
bail!("Failed to copy encoder parameters: error {ret}");
}
// SAFETY: stream_ptr and enc_video are valid; time_base is a plain AVRational
// field copied from encoder to stream.
unsafe {
(*stream_ptr).time_base = (*enc_video.as_ptr()).time_base;
}
// SAFETY: fmt_ctx_ptr is valid; pb is the AVIOContext slot to initialize;
// output_cstr is a valid NUL-terminated path; AVIO_FLAG_WRITE is a constant.
let ret = unsafe {
ffi::avio_open(
&mut (*fmt_ctx_ptr).pb,
@@ -261,17 +290,23 @@ fn main() -> Result<()> {
);
}
// SAFETY: fmt_ctx_ptr is fully configured (streams + pb set); NULL options
// is the default. Returns 0 on success.
let ret = unsafe { ffi::avformat_write_header(fmt_ctx_ptr, ptr::null_mut()) };
if ret < 0 {
bail!("Failed to write header: error {ret}");
}
// SAFETY: fmt_ctx_ptr is a fully initialized output context (header written).
// Output::wrap takes ownership of the pointer into a safe RAII wrapper.
let mut octx = unsafe { ff::format::context::Output::wrap(fmt_ctx_ptr) };
// Create sws_scale context: BGRZ (BGR0) -> YUV420P
let bgr0_fmt = pix_fmt(ff::format::Pixel::BGRZ);
let yuv420p_fmt = pix_fmt(ff::format::Pixel::YUV420P);
// SAFETY: all parameters are valid enum/pixel format values; NULL filters are
// allowed by FFmpeg. sws_getContext returns a heap-allocated SwsContext or NULL.
let sws_ctx = unsafe {
ffi::sws_getContext(
src_width as i32,
@@ -291,6 +326,9 @@ fn main() -> Result<()> {
}
// Allocate reusable YUV frame
// SAFETY: av_frame_alloc returns NULL only on OOM. After allocation we set
// width/height/format fields and call av_frame_get_buffer to allocate plane
// data. On failure we free the frame via av_frame_free before bailing.
let mut yuv_frame = unsafe {
let mut f = ffi::av_frame_alloc();
if f.is_null() {
@@ -349,6 +387,9 @@ fn main() -> Result<()> {
let mmap_start = Instant::now();
let frame_size = (frame.stride as usize) * (frame.height as usize);
// SAFETY: frame.fd is an open DMA-BUF owned by the frame; offset/size come
// from PipeWire's negotiated format. PROT_READ+MAP_SHARED for read-only
// DMA-BUF access. Returns MAP_FAILED on error (checked below).
let mmap_ptr = unsafe {
libc::mmap(
ptr::null_mut(),
@@ -369,8 +410,15 @@ fn main() -> Result<()> {
stats.mmap_us.push(mmap_start.elapsed().as_micros() as u64);
let scale_start = Instant::now();
// SAFETY: mmap_ptr is a valid mapping of frame_size bytes (checked above);
// constructing a read-only slice over it for the duration of sws_scale is
// sound as long as we don't hold it past munmap (we don't).
let src_data = unsafe { std::slice::from_raw_parts(mmap_ptr as *const u8, frame_size) };
// SAFETY: yuv_frame and sws_ctx are valid; src_data is a valid slice of the
// mmap'd DMA-BUF for this frame. sws_scale reads src planes (BGR0 -> YUV420P)
// and writes into yuv_frame's data planes. av_frame_make_writable ensures
// yuv_frame is not shared before writing.
unsafe {
ffi::av_frame_make_writable(yuv_frame);
@@ -391,6 +439,8 @@ fn main() -> Result<()> {
.scale_us
.push(scale_start.elapsed().as_micros() as u64);
// SAFETY: mmap_ptr was returned by mmap above and is not MAP_FAILED; frame_size
// matches the original mapping. Release before dropping frame (which closes fd).
unsafe {
libc::munmap(mmap_ptr, frame_size);
}
@@ -398,6 +448,9 @@ fn main() -> Result<()> {
let encode_start = Instant::now();
// SAFETY: yuv_frame is allocated and writable; enc_video is the opened encoder.
// Setting pts is a plain i64 field write. avcodec_send_frame submits the frame
// for encoding; returns < 0 on error (we log and continue).
unsafe {
(*yuv_frame).pts = pts;
pts += 1;
@@ -419,7 +472,7 @@ fn main() -> Result<()> {
.push(frame_start.elapsed().as_micros() as u64);
frames_encoded += 1;
if frames_encoded % 30 == 0 {
if frames_encoded.is_multiple_of(30) {
let fps = frames_encoded as f64 / total_start.elapsed().as_secs_f64();
println!(
" [{}/{}] {:.1} FPS",
@@ -431,6 +484,8 @@ fn main() -> Result<()> {
let total_elapsed = total_start.elapsed();
println!("\nFlushing encoder...");
// SAFETY: enc_video is the opened encoder; passing NULL frame signals EOF to
// drain the encoder's internal pipeline. Returns < 0 on error (ignored here).
unsafe {
ffi::avcodec_send_frame(enc_video.as_mut_ptr(), ptr::null());
}
@@ -440,6 +495,9 @@ fn main() -> Result<()> {
.map_err(|e| anyhow::anyhow!("Failed to write trailer: {e}"))?;
// Cleanup
// SAFETY: yuv_frame is the allocated frame from earlier (still owned by us);
// sws_ctx is the allocated sws context. av_frame_free and sws_freeContext take
// ownership and free their respective heap allocations.
unsafe {
ffi::av_frame_free(&mut yuv_frame as *mut _);
ffi::sws_freeContext(sws_ctx);
@@ -526,6 +584,9 @@ fn drain_encoder(
) -> Result<()> {
loop {
let mut pkt = ff::Packet::empty();
// SAFETY: enc_video is the opened encoder; pkt is an empty Packet whose
// inner AVPacket pointer is valid. avcodec_receive_packet fills pkt with
// the next encoded packet, or returns EAGAIN/EOF when drained.
let ret = unsafe { ffi::avcodec_receive_packet(enc_video.as_mut_ptr(), pkt.as_mut_ptr()) };
if ret < 0 {
if ret == ffi::AVERROR(ffi::EAGAIN) || ret == ffi::AVERROR_EOF {
@@ -536,6 +597,9 @@ fn drain_encoder(
}
let enc_tb = enc_video.time_base();
// SAFETY: octx.as_ptr() is a valid AVFormatContext; streams is a NULL-terminated
// array of AVStream*. We index [0] which exists because we created exactly one
// stream in setup. Reading time_base is a plain AVRational field access.
let stream_tb = unsafe {
let streams = (*octx.as_ptr()).streams;
let st = *streams.add(0);
+43 -30
View File
@@ -126,6 +126,9 @@ impl Drop for SwsContext {
fn av_err_to_string(ret: i32) -> String {
let mut buf = vec![0u8; 128];
// SAFETY: buf is a 128-byte Vec initialized to zeros; av_strerror writes at most
// buf.len() bytes (including NUL) into the buffer. The ret value is an FFmpeg
// error code. We treat the buffer as `*mut i8` for the C string out-param.
unsafe {
ffi::av_strerror(ret, buf.as_mut_ptr() as *mut i8, buf.len());
}
@@ -134,22 +137,32 @@ fn av_err_to_string(ret: i32) -> String {
}
fn receive_first_frame(cap: &CapPortal) -> Result<wl_webrtc::cap_portal::PwDmaBufFrame> {
// Drain-and-wait loop that mirrors production's repeated-poll semantics
// (state_portal.rs::poll_and_encode driven by main.rs's outer loop), but with
// a single bounded 10s total deadline appropriate for a bench tool. Unlike a
// single 10s blocking wait, this loop actually iterates: each turn drains ALL
// pending control events (the ctrl channel is bounded to 8 — a single
// if-let would silently miss backlog) and then waits a short slice for a
// frame, so StreamEnded/Error arriving mid-wait are observed within ~200ms.
const TOTAL_DEADLINE: std::time::Duration = std::time::Duration::from_secs(10);
const WAIT_SLICE: std::time::Duration = std::time::Duration::from_millis(200);
let deadline = Instant::now() + TOTAL_DEADLINE;
loop {
if let Ok(ctrl) = cap.event_receiver().try_recv() {
while let Ok(ctrl) = cap.event_receiver().try_recv() {
match ctrl {
PwCtrlEvent::StreamEnded => bail!("PipeWire stream ended before first frame"),
PwCtrlEvent::FormatChanged { .. } => {}
PwCtrlEvent::Error(e) => bail!("PipeWire error: {e}"),
}
}
match cap
.frame_receiver()
.recv_timeout(std::time::Duration::from_secs(10))
{
let remaining = match deadline.checked_duration_since(Instant::now()) {
Some(r) if !r.is_zero() => r,
_ => bail!("Timeout waiting for first frame (10s)"),
};
let slice = remaining.min(WAIT_SLICE);
match cap.frame_receiver().recv_timeout(slice) {
Ok(frame) => return Ok(frame),
Err(crossbeam_channel::RecvTimeoutError::Timeout) => {
bail!("Timeout waiting for first frame (10s)");
}
Err(crossbeam_channel::RecvTimeoutError::Timeout) => continue,
Err(crossbeam_channel::RecvTimeoutError::Disconnected) => {
bail!("PipeWire frame channel disconnected");
}
@@ -163,6 +176,9 @@ fn drain_encoder(
) -> Result<()> {
loop {
let mut pkt = ff::Packet::empty();
// SAFETY: enc_video is the opened encoder; pkt is an empty Packet whose inner
// AVPacket pointer is valid. avcodec_receive_packet fills pkt with the next
// encoded packet or returns EAGAIN/EOF when drained.
let ret = unsafe { ffi::avcodec_receive_packet(enc_video.as_mut_ptr(), pkt.as_mut_ptr()) };
if ret < 0 {
if ret == ffi::AVERROR(ffi::EAGAIN) || ret == ffi::AVERROR_EOF {
@@ -172,6 +188,9 @@ fn drain_encoder(
break;
}
let enc_tb = enc_video.time_base();
// SAFETY: octx.as_ptr() is a valid AVFormatContext; streams is a NULL-terminated
// array; we index [0] which exists because we created exactly one stream in
// setup. Reading time_base is a plain AVRational field access.
let stream_tb = unsafe {
let streams = (*octx.as_ptr()).streams;
let st = *streams.add(0);
@@ -398,17 +417,10 @@ fn import_frame(
) -> Result<ff::frame::Video> {
// SAFETY: frames_ctx is a live VAAPI frames context configured for the capture format; frame
// carries a valid DMA-BUF fd and metadata from PipeWire for the duration of the call.
// SAFETY: frames_ctx is a valid VAAPI frames context; `frame` carries the
// DMA-BUF metadata read by the function.
unsafe {
import_dma_buf_to_vaapi(
frames_ctx.as_ptr(),
frame.fd.as_raw_fd(),
frame.width,
frame.height,
frame.format,
frame.modifier,
frame.stride,
frame.offset,
)
import_dma_buf_to_vaapi(frames_ctx.as_ptr(), frame)
}
}
@@ -596,7 +608,7 @@ fn run_cpu_pipeline(
stats.total_us.push(total_us);
stats.frames_encoded += 1;
if stats.frames_encoded <= 3 || stats.frames_encoded % 30 == 0 {
if stats.frames_encoded <= 3 || stats.frames_encoded.is_multiple_of(30) {
println!(
" CPU frame {:>4}/{frames}: import={:.2}ms transfer={:.2}ms scale={:.2}ms encode={:.2}ms total={:.2}ms",
stats.frames_encoded,
@@ -754,7 +766,7 @@ fn run_gpu_pipeline(
stats.total_us.push(total_us);
stats.frames_encoded += 1;
if stats.frames_encoded <= 3 || stats.frames_encoded % 30 == 0 {
if stats.frames_encoded <= 3 || stats.frames_encoded.is_multiple_of(30) {
println!(
" GPU frame {:>4}/{frames}: import={:.2}ms filter={:.2}ms transfer={:.2}ms format={:.2}ms encode={:.2}ms total={:.2}ms",
stats.frames_encoded,
@@ -919,17 +931,13 @@ fn main() -> Result<()> {
AvHwFrameCtx::for_capture(&hw_dev, src_width, src_height, ff::format::Pixel::BGRA)?;
println!(" VAAPI frames context created OK (sw_format=BGRA)");
// SAFETY: delegates to avhw::import_dma_buf_to_vaapi (itself an unsafe fn).
// frames_ctx is a valid AVBufferRef from AvHwFrameCtx::for_capture above;
// `first_frame` is the PipeWire-formatted PwDmaBufFrame whose metadata the
// function reads directly. See that function's own SAFETY contract for the
// full rationale.
let vaapi_frame = unsafe {
import_dma_buf_to_vaapi(
frames_ctx.as_ptr(),
first_frame.fd.as_raw_fd(),
first_frame.width,
first_frame.height,
first_frame.format,
first_frame.modifier,
first_frame.stride,
first_frame.offset,
)
import_dma_buf_to_vaapi(frames_ctx.as_ptr(), &first_frame)
};
match &vaapi_frame {
@@ -949,6 +957,9 @@ fn main() -> Result<()> {
let mmap_size = (first_frame.stride as usize) * (first_frame.height as usize);
let mmap_start = Instant::now();
// SAFETY: first_frame.fd is an open DMA-BUF; offset/size from PipeWire.
// PROT_READ+MAP_SHARED is the standard read-only DMA-BUF mapping. Returns
// MAP_FAILED on error (checked below).
let mmap_ptr = unsafe {
libc::mmap(
ptr::null_mut(),
@@ -970,6 +981,8 @@ fn main() -> Result<()> {
mmap_size as f64 / 1024.0 / 1024.0,
mmap_elapsed.as_secs_f64() * 1000.0
);
// SAFETY: mmap_ptr is a valid mapping (MAP_FAILED path was handled
// above); mmap_size matches the original mapping. POSIX munmap(2).
unsafe {
libc::munmap(mmap_ptr, mmap_size);
}
+106 -9
View File
@@ -95,6 +95,23 @@ pub struct PwDmaBufFrame {
pub pts: i64,
}
/// PipeWire-negotiated video format snapshot, stashed in a `Cell` for cross-callback
/// sharing (format-change callback writes it; process callback reads it). The four
/// fields are the minimal subset of `PwDmaBufFrame`'s metadata that the process
/// callback needs to construct the frame once a buffer arrives.
///
/// `Copy` is required because we store it inside `Cell<Option<PortalFormatInfo>>`;
/// `Cell` requires its contents to be `Copy` (no borrowed interior state).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PortalFormatInfo {
pub width: u32,
pub height: u32,
/// DRM FourCC format code (e.g. `0x34325258` for XR24 / XRGB8888).
pub drm_format: u32,
/// DRM format modifier describing buffer layout (linear, tiling, etc.).
pub modifier: u64,
}
/// PipeWire 控制事件枚举
///
/// 从 PipeWire 捕获线程发送给消费者的控制事件。
@@ -158,6 +175,9 @@ impl CapPortal {
let (frame_tx, frame_rx) = bounded(1);
let (event_tx, event_rx) = bounded(8);
// SAFETY: eventfd(2) is a POSIX syscall with no preconditions; the init value
// and flags (CLOEXEC + NONBLOCK) are valid. Returns either a fresh fd (>= 0)
// or -1 on error, which we check immediately below.
let efd = unsafe { libc::eventfd(0, libc::EFD_CLOEXEC | libc::EFD_NONBLOCK) };
if efd < 0 {
return Err(anyhow::anyhow!(
@@ -165,9 +185,13 @@ impl CapPortal {
std::io::Error::last_os_error()
));
}
// SAFETY: `efd` is the open eventfd we just created (>= 0 checked above) and
// own. dup(2) returns either a fresh fd or -1.
let write_fd = unsafe { libc::dup(efd) };
if write_fd < 0 {
let err = std::io::Error::last_os_error();
// SAFETY: `efd` is still the open eventfd we own; closing on the error
// path before returning to avoid fd leak.
unsafe { libc::close(efd) };
return Err(anyhow::anyhow!("dup eventfd failed: {err}"));
}
@@ -178,6 +202,10 @@ impl CapPortal {
frame_tx,
event_tx,
dropped: pw_dropped.clone(),
// SAFETY: `efd` is the freshly-created eventfd (>= 0 checked above) and we
// are its sole owner. OwnedFd::from_raw_fd takes ownership and will close()
// it on Drop. Ownership transfers into PwThreadCtx and then into the
// PipeWire thread via pipewire_thread.
shutdown_read: unsafe { OwnedFd::from_raw_fd(efd) },
pw_fd,
node_id,
@@ -190,11 +218,16 @@ impl CapPortal {
pipewire_thread(ctx);
})
.map_err(|e| {
// SAFETY: `write_fd` is the open dup'd eventfd we own (>= 0 checked
// above); closing on thread-spawn failure to avoid fd leak.
unsafe { libc::close(write_fd) };
anyhow::anyhow!("thread spawn failed: {e}")
})?;
Ok(Self {
// SAFETY: `write_fd` is the freshly-dup'd eventfd (>= 0 checked above) and
// we are its sole owner. OwnedFd::from_raw_fd takes ownership and will
// close() it on Drop (which fires when CapPortal is dropped).
shutdown_fd: unsafe { OwnedFd::from_raw_fd(write_fd) },
frame_rx,
event_rx,
@@ -433,7 +466,10 @@ fn verify_secure_dir(path: &std::path::Path) -> bool {
return false;
}
// Must be owned by current user
if meta.uid() != unsafe { libc::getuid() } {
// SAFETY: libc::getuid has no preconditions and cannot fail; it simply
// returns the calling process's real user ID.
// SAFETY: libc::getuid has no preconditions and cannot fail.
if meta.uid() != unsafe { libc::getuid() } {
tracing::warn!(
"Token parent dir not owned by current user: {}",
path.display()
@@ -511,6 +547,7 @@ fn load_restore_token_from(path: PathBuf) -> Option<String> {
tracing::warn!("Token path is not a regular file: {}", path.display());
return None;
}
// SAFETY: libc::getuid has no preconditions and cannot fail.
if meta.uid() != unsafe { libc::getuid() } {
tracing::warn!("Token file not owned by current user: {}", path.display());
return None;
@@ -604,6 +641,9 @@ impl Drop for CapPortal {
// Signal the PipeWire loop to quit via eventfd.
// eventfd write is a kernel syscall — thread-safe and lock-free.
let val: u64 = 1u64;
// SAFETY: shutdown_fd is a valid open eventfd (owned by Self); the buffer is
// a stack u64 of size 8 bytes which matches the count argument. POSIX write(2)
// is the standard fd-write syscall; eventfd writes must be exactly 8 bytes.
let _ = unsafe {
libc::write(
self.shutdown_fd.as_raw_fd(),
@@ -657,7 +697,7 @@ fn pipewire_thread(ctx: PwThreadCtx) {
shutdown_read,
pw_fd,
node_id,
fps,
fps: _,
} = ctx;
let mainloop = match pw::main_loop::MainLoopBox::new(None) {
@@ -721,7 +761,7 @@ fn pipewire_thread(ctx: PwThreadCtx) {
}
};
let format_info: Rc<Cell<Option<(u32, u32, u32, u64)>>> = Rc::new(Cell::new(None));
let format_info: Rc<Cell<Option<PortalFormatInfo>>> = Rc::new(Cell::new(None));
let event_tx_state = event_tx.clone();
let _listener = stream
@@ -773,9 +813,14 @@ fn pipewire_thread(ctx: PwThreadCtx) {
let max_framerate = info.max_framerate();
// 保存协商后的格式信息,供 process 回调读取
let previous_format = format_info.get();
format_info.set(Some((width, height, drm_format, modifier)));
if let Some((previous_width, previous_height, _, _)) = previous_format {
if width != previous_width || height != previous_height {
format_info.set(Some(PortalFormatInfo {
width,
height,
drm_format,
modifier,
}));
if let Some(prev) = previous_format {
if width != prev.width || height != prev.height {
tracing::warn!(
"PipeWire dimensions changed: {}x{} (format renegotiation)",
width,
@@ -801,8 +846,20 @@ fn pipewire_thread(ctx: PwThreadCtx) {
.process({
let format_info = format_info.clone();
let frame_tx = frame_tx.clone();
let dropped = dropped;
move |stream, _| {
// SAFETY: raw_buf ownership invariant — PipeWire's process callback
// contract requires that every buffer acquired via `dequeue_raw_buffer`
// is returned to the queue EXACTLY ONCE via `queue_raw_buffer` before
// the callback returns — on every exit path, success or error. Failure
// to requeue leaks the buffer slot and eventually stalls the stream.
//
// Audit map of this closure (verified 2026-06-28):
// - null raw_buf (dequeue returned NULL) → nothing to requeue, return.
// - null spa_buf / no data / bad fd / null chunk / no format_info /
// invalid dims / dup_fd < 0 → all requeue before early-return.
// - success (try_send Ok / Full / Disconnected) → final requeue at end.
// The fd ownership is independent: dup() creates a fresh fd that lives
// inside PwDmaBufFrame; on try_send error the frame Drops and closes it.
let raw_buf = unsafe { stream.dequeue_raw_buffer() };
if raw_buf.is_null() {
tracing::trace!("process: null raw_buf");
@@ -810,36 +867,49 @@ fn pipewire_thread(ctx: PwThreadCtx) {
}
// 获取 SPA buffer 结构体,包含数据数组、元数据等
// SAFETY: raw_buf was checked non-null above. `pw_buffer.buffer` is a
// valid raw pointer for the lifetime of raw_buf (PipeWire keeps the
// buffer alive until we queue it back).
let spa_buf = unsafe { (*raw_buf).buffer };
if spa_buf.is_null() {
tracing::trace!("process: null spa_buf");
// SAFETY: raw_buf is the non-null buffer we still own; returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
// 获取 buffer 中的数据项数量和数据指针
// 对于 DMA-BUF 帧,通常只有 1 个数据项(包含 fd)
// SAFETY: spa_buf checked non-null above; `n_datas` is a plain u32 field.
let n_datas = unsafe { (*spa_buf).n_datas };
// SAFETY: same as above; `datas` is a raw pointer field, may be null.
let datas_ptr = unsafe { (*spa_buf).datas };
if n_datas == 0 || datas_ptr.is_null() {
tracing::trace!("process: no data (n_datas={n_datas})");
// SAFETY: raw_buf still owned, returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
// 从第一个数据项中获取 DMA-BUF 文件描述符
// 通过 libspa 的 Data 包装类型安全地访问 SPA 数据结构
// SAFETY: datas_ptr is non-null and n_datas > 0 (checked above). We cast
// to pw::spa::buffer::Data and take a shared borrow; PipeWire does not
// mutate the data array during a process cycle, so a shared reference
// for the duration of this callback is sound.
let data_ref: &pw::spa::buffer::Data =
unsafe { &*(datas_ptr as *const pw::spa::buffer::Data) };
let fd = data_ref.fd();
if fd < 0 {
tracing::trace!("process: invalid fd={fd}");
// SAFETY: raw_buf still owned, returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
if data_ref.as_raw().chunk.is_null() {
tracing::trace!("process: null chunk");
// SAFETY: raw_buf still owned, returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
@@ -850,6 +920,12 @@ fn pipewire_thread(ctx: PwThreadCtx) {
// 从 SPA_META_Header 元数据中提取 PTS (显示时间戳)
// 遍历 buffer 的所有元数据项,查找 Header 类型的元数据
// PTS 可用于音视频同步和帧率控制
// SAFETY: spa_buf is non-null. `metas` is checked for null before
// iteration. We iterate `i in 0..n_metas` reading shared POD fields
// (type_, size, data) — PipeWire keeps the meta array immutable during
// a process cycle. The size guard (`meta.size >= size_of::<spa_meta_header>()`)
// and null-data check before reading ensure we never read past the
// meta's actual extent.
let pts: i64 = unsafe {
let mut pts_val: i64 = 0;
let n_metas = (*spa_buf).n_metas;
@@ -872,12 +948,15 @@ fn pipewire_thread(ctx: PwThreadCtx) {
};
// 验证格式信息已协商完成,且分辨率和格式有效
let Some((width, height, format, modifier)) = format_info.get() else {
let Some(fmt) = format_info.get() else {
// SAFETY: raw_buf still owned, returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
};
let PortalFormatInfo { width, height, drm_format: format, modifier } = fmt;
if width == 0 || height == 0 || format == 0 {
tracing::trace!("process: invalid dimensions {width}x{height} format={format}");
// SAFETY: raw_buf still owned, returning it.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
@@ -885,15 +964,27 @@ fn pipewire_thread(ctx: PwThreadCtx) {
// 复制 DMA-BUF 文件描述符
// 必须 dup,因为原始 fd 由 PipeWire 管理,我们不能持有它
// dup 后的 fd 由 PwDmaBufFrame 持有,生命周期独立于 PipeWire buffer
// SAFETY: `fd` is the open DMA-BUF fd reported by PipeWire (>= 0 checked
// above). libc::dup is the standard POSIX fd duplication call. The
// original `fd` remains owned by PipeWire (returned with raw_buf later).
let dup_fd = unsafe { libc::dup(fd) };
if dup_fd < 0 {
// SAFETY: raw_buf still owned, returning it. No fd cleanup needed
// because dup() failed and never returned a new fd.
unsafe { stream.queue_raw_buffer(raw_buf) };
return;
}
// 构建帧数据对象,所有必要的帧信息已收集完毕
// SAFETY: `dup_fd` is a freshly-dup'd open file descriptor (>= 0 checked
// above) and we are its sole owner. OwnedFd::from_raw_fd takes ownership
// and will close() it on Drop. The fd's lifecycle is independent of
// raw_buf: whether try_send succeeds (frame moves into the channel) or
// fails (Full/Disconnected — the error payload owns the frame and drops
// it at the end of the match arm), exactly one close() occurs per dup().
let frame_fd = unsafe { OwnedFd::from_raw_fd(dup_fd) };
let frame = PwDmaBufFrame {
fd: unsafe { OwnedFd::from_raw_fd(dup_fd) },
fd: frame_fd,
offset,
stride,
modifier,
@@ -910,6 +1001,8 @@ fn pipewire_thread(ctx: PwThreadCtx) {
}
Err(crossbeam_channel::TrySendError::Disconnected(_)) => {}
}
// SAFETY: final exactly-once requeue of raw_buf. Every path above
// either returned early with its own requeue, or falls through to here.
unsafe { stream.queue_raw_buffer(raw_buf) };
}
})
@@ -949,6 +1042,10 @@ fn pipewire_thread(ctx: PwThreadCtx) {
move |fd| {
// Drain the eventfd so it doesn't re-trigger
let mut buf: u64 = 0;
// SAFETY: `fd` is the registered eventfd owned by the mainloop source; the
// buffer is a stack u64 of 8 bytes matching the count argument. POSIX
// read(2) is the standard fd-read syscall; eventfd semantics require the
// 8-byte buffer.
let _ = unsafe {
libc::read(
fd.as_raw_fd(),
+3 -11
View File
@@ -21,9 +21,9 @@ pub struct CapWlrScreencopy {
}
impl CaptureSource for CapWlrScreencopy {
/// Unit type: wlr-screencopy is fully asynchronous — `alloc_frame()`
/// always returns `None`. The frame object is created by Dispatch
/// impls calling `manager.capture_output()`, not by this method.
/// Unit type: wlr-screencopy is fully asynchronous — frame allocation is
/// driven by Dispatch impls calling `manager.capture_output()`, so there
/// is no synchronous `alloc_frame`-style API on this trait.
type Frame = ();
fn new(
@@ -40,14 +40,6 @@ impl CaptureSource for CapWlrScreencopy {
})
}
fn alloc_frame(&mut self) -> Option<Self::Frame> {
// wlr-screencopy is asynchronous: the Dispatch impl creates a new
// ZwlrScreencopyFrameV1 which triggers the buffer allocation flow
// (buffer event → negotiate format → create DMA-BUF). This method
// always returns None.
None
}
fn queue_copy(&mut self, buffer: &WlBuffer, _qh: &QueueHandle<State<Self>>) {
if let Some(frame) = &self.current_frame {
frame.copy(buffer);
+5 -2
View File
@@ -148,6 +148,9 @@ fn run_wlr_screencopy(args: Args) -> Result<()> {
revents: 0,
};
// timeout=0 表示非阻塞,立即返回当前 fd 状态
// SAFETY: `pfd` is a stack-allocated libc::pollfd initialized above with a
// valid wayland_fd and POLLIN events; nfds=1 matches the single-element
// array; timeout=0 is non-blocking. POSIX poll(2) writes revents in place.
let ret = unsafe { libc::poll(&mut pfd, 1, 0) };
tracing::info!(
"Raw poll on wayland fd={wayland_fd}: ret={ret}, revents={}",
@@ -173,7 +176,7 @@ fn run_wlr_screencopy(args: Args) -> Result<()> {
// 注册 SIGINT / SIGTERM 信号用于优雅退出
// signal_hook_mio 将 Unix 信号转换为 fd 可读事件,
// 这样信号也可以通过 epoll 统一监听,不需要单独的信号处理器
let mut signals = signal_hook_mio::v1_0::Signals::new(&[
let mut signals = signal_hook_mio::v1_0::Signals::new([
signal_hook::consts::SIGINT, // Ctrl+C
signal_hook::consts::SIGTERM, // kill 命令默认信号
])?;
@@ -310,7 +313,7 @@ fn run_portal_pipewire(args: Args) -> Result<()> {
// 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(&[
let mut signals = signal_hook_mio::v1_0::Signals::new([
signal_hook::consts::SIGINT,
signal_hook::consts::SIGTERM,
])?;
+40 -33
View File
@@ -65,8 +65,6 @@ pub trait CaptureSource: Sized + 'static {
qh: &QueueHandle<State<Self>>,
) -> Result<Self>;
fn alloc_frame(&mut self) -> Option<Self::Frame>;
fn queue_copy(&mut self, buffer: &WlBuffer, qh: &QueueHandle<State<Self>>);
fn on_done_with_frame(&mut self, frame: Self::Frame);
@@ -83,6 +81,7 @@ pub struct OutputInfo {
pub logical_position: (i32, i32),
}
#[derive(Default)]
pub struct PartialOutputInfo {
pub name: Option<String>,
/// Name from wl_output::Name (v4) — used to match wlr-output-management heads
@@ -95,22 +94,11 @@ pub struct PartialOutputInfo {
pub done_count: u32,
}
impl Default for PartialOutputInfo {
fn default() -> Self {
Self {
name: None,
wl_name: None,
transform: None,
physical_size: None,
logical_position: None,
mode_size: None,
done_count: 0,
}
}
}
/// Stores head info from wlr-output-management for name-based matching with wl_output.
struct WlrHeadInfo {
// `pub(crate)` (not module-private): exposed via `EncConstructionStage::ProbingOutputs.wlr_heads`
// which is reached from main.rs during the wlr-screencopy probing loop.
pub(crate) struct WlrHeadInfo {
position: Option<(i32, i32)>,
}
@@ -138,7 +126,7 @@ impl StreamingEncoder {
}
}
fn encode_frame(&mut self, hw_frame: &ffmpeg_next::frame::Video) -> anyhow::Result<()> {
fn encode_frame(&mut self, hw_frame: &ffmpeg_next::frame::Video) -> anyhow::Result<crate::avhw::EncodeStages> {
match self {
StreamingEncoder::Mp4(enc) => enc.encode_frame(hw_frame),
StreamingEncoder::WebRtc(enc) => enc.encode_frame(hw_frame),
@@ -156,8 +144,14 @@ impl StreamingEncoder {
// ---------------------------------------------------------------------------
// EncConstructionStage
// ---------------------------------------------------------------------------
//
// `pub(crate)` (not `pub`): this enum leaks the private `WlrHeadInfo` type via
// its `wlr_heads` field, and the construction-stage state machine is an
// internal implementation detail. Crate-internal consumers (main.rs) get there
// via `crate::state::`; there is no need to expose this across the crate
// boundary. See Oracle audit 2026-06-28.
pub enum EncConstructionStage<S: CaptureSource> {
pub(crate) enum EncConstructionStage<S: CaptureSource> {
ProbingOutputs {
outputs: Vec<PartialOutputInfo>,
bound_outputs: Vec<WlOutput>,
@@ -197,10 +191,12 @@ pub enum EncConstructionStage<S: CaptureSource> {
pub enum InFlightSurface<S: CaptureSource> {
None,
AllocQueued,
Allocd(S::Frame),
CopyQueued {
surface: ff::frame::Video,
drm_map: ff::ffi::AVDRMFrameDescriptor,
// Boxed: AVDRMFrameDescriptor is ~592 bytes (4 objects + 4 layers),
// which would balloon every InFlightSurface variant via enum alignment.
// The box shrinks the enum to ~32 bytes regardless of variant.
drm_map: Box<ff::ffi::AVDRMFrameDescriptor>,
frame: S::Frame,
buffer: WlBuffer,
},
@@ -211,7 +207,7 @@ pub enum InFlightSurface<S: CaptureSource> {
// ---------------------------------------------------------------------------
pub struct State<S: CaptureSource> {
pub stage: EncConstructionStage<S>,
pub(crate) stage: EncConstructionStage<S>,
pub in_flight_surface: InFlightSurface<S>,
pub starting_timestamp: Option<i64>,
pub stats_start_time: Option<Instant>,
@@ -512,6 +508,10 @@ impl<S: CaptureSource> State<S> {
unsafe {
(*map_frame.as_mut_ptr()).format = ffi::AVPixelFormat::AV_PIX_FMT_DRM_PRIME as i32;
}
// SAFETY: map_frame and surface are valid, owned AVFrame pointers from
// av_hwframe_get/surface.alloc above. AV_HWFRAME_MAP_READ flag (0 here)
// requests a read-only mapping. The DRM_PRIME format set above instructs
// FFmpeg to populate data[0] with an AVDRMFrameDescriptor on success.
let ret = unsafe { ffi::av_hwframe_map(map_frame.as_mut_ptr(), surface.as_ptr(), 0) };
if ret < 0 {
tracing::error!("av_hwframe_map failed: {}", crate::avhw::ff_err(ret));
@@ -571,7 +571,7 @@ impl<S: CaptureSource> State<S> {
);
self.in_flight_surface = InFlightSurface::CopyQueued {
surface,
drm_map: desc,
drm_map: Box::new(desc),
frame,
buffer: wl_buffer,
};
@@ -629,15 +629,22 @@ impl<S: CaptureSource> State<S> {
};
if should_encode {
let encode_start = Instant::now();
if let Err(e) = enc.encode_frame(&surface) {
tracing::error!("encode_frame failed: {}", e);
self.errored = true;
match enc.encode_frame(&surface) {
Ok(stages) => {
let encode_elapsed = encode_start.elapsed().as_micros() as u64;
self.stats.record_encode(&FrameTimings {
scale_us: stages.scale_us,
transfer_us: stages.transfer_us,
encode_us: stages.encode_us,
total_us: encode_elapsed,
..Default::default()
});
}
Err(e) => {
tracing::error!("encode_frame failed: {}", e);
self.errored = true;
}
}
let encode_elapsed = encode_start.elapsed().as_micros() as u64;
self.stats.record_encode(&FrameTimings {
total_us: encode_elapsed,
..Default::default()
});
}
self.stats_frames += 1;
if let Some(last) = self.stats_last_time {
@@ -1509,7 +1516,7 @@ impl<S: CaptureSource> Dispatch<ZwlrOutputManagerV1, ()> for State<S> {
event: <ZwlrOutputManagerV1 as Proxy>::Event,
_data: &(),
_conn: &wayland_client::Connection,
qhandle: &QueueHandle<State<S>>,
_qhandle: &QueueHandle<State<S>>,
) {
match event {
WlrOutputManagerEvent::Head { head } => {
@@ -1530,7 +1537,7 @@ impl<S: CaptureSource> Dispatch<ZwlrOutputManagerV1, ()> for State<S> {
}
}
}
WlrOutputManagerEvent::Finished { .. } => {
WlrOutputManagerEvent::Finished => {
tracing::warn!("zwlr_output_manager_v1::Finished received during probing");
}
_ => {}
@@ -1583,7 +1590,7 @@ impl<S: CaptureSource> Dispatch<ZwlrOutputHeadV1, ()> for State<S> {
}
}
}
WlrHeadEvent::Finished { .. } => {
WlrHeadEvent::Finished => {
tracing::debug!("zwlr_output_head_v1::Finished received");
}
_ => {}
+87 -45
View File
@@ -1,4 +1,7 @@
// 采集门户状态模块 —— 通过 PipeWire/DMA-BUF 进行屏幕采集并编码
// AsRawFd is required by frame.fd.as_raw_fd() in build_drm_descriptor below
// but rustc emits a false "unused_imports" warning because OwnedFd also has
// an inherent as_raw_fd — same quirk as avhw.rs. E0599 if removed → keep it.
use std::os::fd::AsRawFd;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
@@ -42,6 +45,26 @@ struct WebrtcThread {
sent_gap_rx: crossbeam_channel::Receiver<(f64, Option<f64>)>,
}
/// Static configuration handed to the WebRTC sender thread. Immutable for the
/// thread's lifetime; a resolution tier change rebuilds the whole pipeline
/// (and spawns a new thread) rather than mutating this.
struct WebRtcThreadConfig {
fps: u32,
enc_width: u32,
enc_height: u32,
max_bitrate: u64,
}
/// Channel endpoints owned exclusively by the WebRTC sender thread after spawn.
/// The reverse endpoints stay with StatePortal (or the encode thread) for
/// inbound/outbound traffic.
struct WebRtcThreadChannels {
webrtc_rx: crossbeam_channel::Receiver<EncodedH264Frame>,
sent_gap_tx: crossbeam_channel::Sender<(f64, Option<f64>)>,
bitrate_tx: crossbeam_channel::Sender<BitrateCommand>,
resolution_tx: crossbeam_channel::Sender<BitrateCommand>,
}
/// 门户模式的主状态机
///
/// 负责管理从 PipeWire 采集屏幕帧、通过 VAAPI 硬件编码的完整生命周期。
@@ -288,15 +311,19 @@ impl StatePortal {
.spawn(move || {
webrtc_thread_loop(
wrtc,
webrtc_rx,
fps,
enc_width,
enc_height,
max_bitrate,
WebRtcThreadConfig {
fps,
enc_width,
enc_height,
max_bitrate,
},
WebRtcThreadChannels {
webrtc_rx,
sent_gap_tx,
bitrate_tx,
resolution_tx,
},
paused,
sent_gap_tx,
bitrate_tx,
resolution_tx,
)
})?;
self.webrtc_thread = Some(WebrtcThread {
@@ -338,8 +365,15 @@ impl StatePortal {
// 每秒输出一次结构化管道统计(仅 --stats 启用时记录日志)
if self.args.stats && self.stats.should_snapshot() {
self.stats.set_pipewire_dropped(0, 0);
self.stats.set_queue_depths(0, 0);
// Wire PipeWire drop counter (delta-tracked via pw_dropped_prev) and
// capture channel depth. Oracle audit 2026-06-28: previously hardcoded
// (0, 0), which silently zeroed two real diagnostic fields.
let total_dropped = self.cap.dropped_count();
self.stats.set_pipewire_dropped(total_dropped, self.pw_dropped_prev);
self.pw_dropped_prev = total_dropped;
// capture queue depth is real; encoded side has no exposed depth — the
// encoder thread publishes timings only, not a frame queue length.
self.stats.set_queue_depths(self.cap.capture_queue_depth(), 0);
if let Some(ref enc_thread) = self.enc_thread {
while let Ok(timing) = enc_thread.timing_rx.try_recv() {
self.stats.record_encode_thread(
@@ -478,55 +512,47 @@ impl StatePortal {
if let Some(enc) = self.enc.as_mut() {
// 将 DMA-BUF 帧零拷贝导入 VAAPI 硬件帧池
// SAFETY: delegates to avhw::import_dma_buf_to_vaapi (itself an unsafe fn);
// frames_rgb pointer is a valid AVBufferRef owned by enc, and `frame` is the
// PipeWire-formatted PwDmaBufFrame whose metadata the function reads directly.
// See that function's own SAFETY contract.
let mut vaapi_frame = unsafe {
avhw::import_dma_buf_to_vaapi(
enc.frames_rgb().as_ptr(),
frame.fd.as_raw_fd(),
frame.width,
frame.height,
frame.format,
frame.modifier,
frame.stride,
frame.offset,
)
avhw::import_dma_buf_to_vaapi(enc.frames_rgb().as_ptr(), &frame)
}?;
let import_us = t_import_start.elapsed().as_micros() as u64;
let t_encode_start = Instant::now();
// 设置帧的显示时间戳(PTS),基于已编码帧序号
// SAFETY: vaapi_frame is the freshly-imported valid AVFrame returned by
// import_dma_buf_to_vaapi above; pts is a plain i64 field on AVFrame.
unsafe {
(*vaapi_frame.as_mut_ptr()).pts = pts;
}
// 送入编码器完成:缩放 → 回读 → 格式转换 → H.264 编码
enc.encode_frame(&vaapi_frame)?;
let stages = enc.encode_frame(&vaapi_frame)?;
let total_us = t_import_start.elapsed().as_micros() as u64;
let encode_us = t_encode_start.elapsed().as_micros() as u64;
let encode_us = stages.encode_us;
self.frames_encoded += 1;
// 记录帧计时到管道统计(import + encode 内部各阶段暂不可分离,用 total 覆盖
// 记录帧计时到管道统计(scale 来自 filter graphtransfer 在 HW 路径恒为 0
let timings = FrameTimings {
import_us,
scale_us: stages.scale_us,
transfer_us: stages.transfer_us,
encode_us,
total_us,
..Default::default()
};
self.stats.record_encode(&timings);
} else if let Some(import) = self.enc_import.as_mut() {
// SAFETY: same contract as the enc branch above — frames_rgb owned by
// import, `frame` carries the PipeWire DMA-BUF metadata.
let mut vaapi_frame = unsafe {
avhw::import_dma_buf_to_vaapi(
import.frames_rgb().as_ptr(),
frame.fd.as_raw_fd(),
frame.width,
frame.height,
frame.format,
frame.modifier,
frame.stride,
frame.offset,
)
avhw::import_dma_buf_to_vaapi(import.frames_rgb().as_ptr(), &frame)
}?;
// SAFETY: vaapi_frame is the valid AVFrame returned above; pts is plain i64.
unsafe {
(*vaapi_frame.as_mut_ptr()).pts = pts;
}
@@ -696,16 +722,22 @@ fn encode_thread_loop(
fn webrtc_thread_loop(
mut wrtc: WebRtcState,
webrtc_rx: crossbeam_channel::Receiver<EncodedH264Frame>,
fps: u32,
enc_width: u32,
enc_height: u32,
max_bitrate: u64,
config: WebRtcThreadConfig,
channels: WebRtcThreadChannels,
paused: Arc<AtomicBool>,
sent_gap_tx: crossbeam_channel::Sender<(f64, Option<f64>)>,
bitrate_tx: crossbeam_channel::Sender<BitrateCommand>,
resolution_tx: crossbeam_channel::Sender<BitrateCommand>,
) {
let WebRtcThreadConfig {
fps,
enc_width,
enc_height,
max_bitrate,
} = config;
let WebRtcThreadChannels {
webrtc_rx,
sent_gap_tx,
bitrate_tx,
resolution_tx,
} = channels;
let mut frames_sent: u64 = 0;
let mut last_send: Option<std::time::Instant> = None;
let mut last_sent_bitrate: Option<u64> = None;
@@ -756,7 +788,7 @@ fn webrtc_thread_loop(
let should_send = match last_sent_bitrate {
None => true,
Some(last) => {
let diff = if bwe > last { bwe - last } else { last - bwe };
let diff = bwe.abs_diff(last);
diff * 10 > last
}
};
@@ -946,7 +978,12 @@ fn resolve_drm_device(args: &Args) -> Result<Option<PathBuf>> {
/// 用于验证 DMA-BUF 元数据映射的正确性。
#[cfg(test)]
fn build_drm_descriptor(frame: &PwDmaBufFrame) -> ffmpeg_next::ffi::AVDRMFrameDescriptor {
let mut desc: ffmpeg_next::ffi::AVDRMFrameDescriptor = unsafe { std::mem::zeroed() };
let mut desc: ffmpeg_next::ffi::AVDRMFrameDescriptor = {
// SAFETY: AVDRMFrameDescriptor is a POD struct from FFmpeg's C API with no
// pointers orDrop fields; all-zero is a valid initial state. Every field is
// explicitly overwritten in the lines below before the descriptor is used.
unsafe { std::mem::zeroed() }
};
desc.nb_objects = 1; // 单个 DMA-BUF 对象
desc.objects[0].fd = frame.fd.as_raw_fd(); // DMA-BUF 文件描述符
desc.objects[0].size = 0; // 大小设为 0(内核自动确定)
@@ -969,6 +1006,9 @@ mod tests {
fn make_test_frame() -> PwDmaBufFrame {
// Create a dummy fd from stderr (always valid fd 2)
// 使用 stderr(fd 2)的副本作为虚拟文件描述符
// SAFETY: stderr (fd 2) is always-open in any process; libc::dup(2) returns
// a fresh fd we solely own. OwnedFd::from_raw_fd takes ownership and closes
// it on Drop. Test-only; the fd is never actually memory-mapped.
let fd = unsafe { OwnedFd::from_raw_fd(libc::dup(2)) };
PwDmaBufFrame {
fd,
@@ -1091,8 +1131,10 @@ mod tests {
/// 测试:使用自定义偏移量和 stride 构建 DRM 描述符
#[test]
fn build_drm_descriptor_custom_offset_and_stride() {
// SAFETY: same as make_test_frame — dup of stderr (fd 2), test-only.
let test_fd = unsafe { OwnedFd::from_raw_fd(libc::dup(2)) };
let frame = PwDmaBufFrame {
fd: unsafe { OwnedFd::from_raw_fd(libc::dup(2)) },
fd: test_fd,
offset: 4096, // 4KB 对齐偏移
stride: 3840 * 4, // 4K 宽度 × 4 字节
modifier: 0x0100000000000001, // AMD modifiers
+47 -22
View File
@@ -38,7 +38,6 @@ pub struct PipelineStats {
encoded_frames: u64,
sent_frames: u64,
pipewire_dropped: u64,
over_budget_count: u64,
/// Count of frames dropped by encode thread due to Y-plane hash dedup
/// (EncodeOutcome::SkippedDuplicate). Read from atomic counter set by
/// encode thread, computed as delta since previous snapshot.
@@ -73,6 +72,12 @@ pub struct PipelineStats {
window_start: Instant,
}
impl Default for PipelineStats {
fn default() -> Self {
Self::new()
}
}
impl PipelineStats {
pub fn new() -> Self {
Self {
@@ -80,7 +85,6 @@ impl PipelineStats {
encoded_frames: 0,
sent_frames: 0,
pipewire_dropped: 0,
over_budget_count: 0,
duplicate_frames_skipped: 0,
prev_duplicate_frames_skipped: 0,
capture_queue_depth: 0,
@@ -207,11 +211,6 @@ impl PipelineStats {
self.encoded_queue_depth = encoded;
}
/// Record that a frame exceeded its time budget.
pub fn record_over_budget(&mut self) {
self.over_budget_count += 1;
}
/// Returns true if at least 1 second has elapsed since the last snapshot
/// (or since creation). If true, call `snapshot_and_reset` to get the stats.
pub fn should_snapshot(&self) -> bool {
@@ -230,7 +229,6 @@ impl PipelineStats {
encoded_frames: self.encoded_frames,
sent_frames: self.sent_frames,
pipewire_dropped: self.pipewire_dropped,
over_budget_count: self.over_budget_count,
duplicate_frames_skipped: self.duplicate_frames_skipped,
capture_queue_depth: self.capture_queue_depth,
encoded_queue_depth: self.encoded_queue_depth,
@@ -269,7 +267,6 @@ impl PipelineStats {
self.encoded_frames = 0;
self.sent_frames = 0;
self.pipewire_dropped = 0;
self.over_budget_count = 0;
self.duplicate_frames_skipped = 0;
self.capture_queue_depth = 0;
self.encoded_queue_depth = 0;
@@ -304,7 +301,6 @@ pub struct StatsSnapshot {
pub encoded_frames: u64,
pub sent_frames: u64,
pub pipewire_dropped: u64,
pub over_budget_count: u64,
pub duplicate_frames_skipped: u64,
// Queue depths
pub capture_queue_depth: usize,
@@ -346,43 +342,72 @@ pub struct StatsSnapshot {
impl std::fmt::Display for StatsSnapshot {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
// Layout note: each line answers one operational question.
// - line 1: throughput (fps + frame counts + window length)
// - line 2: drops (PipeWire backlog, encoder over-budget, frame-hash dedup)
// - line 3: queue back-pressure (capture + encoded)
// - line 4-6: gap timing (avg + p95 + max) — answers "is cadence stable?"
// - line 7: capture-to-send age (avg + p95 + max) — answers "how stale?"
// - line 8: per-stage encode timing (avg + p95) — answers "where is latency?"
// - line 9: output bandwidth (bytes/sec + per-frame p95/max)
//
// The avg counterparts were computed but never displayed before Oracle
// audit 2026-06-28; they pair with the existing p95/max to surface both
// central tendency and tail behaviour in the same glance.
write!(
f,
"capture_fps={:.1} encoded_fps={:.1} sent_fps={:.1} \
pw_dropped={} over_budget={} duplicate_frames_skipped={} \
cap_q={} enc_q={} \
cap_gap_p95={:.1}ms cap_gap_max={:.1}ms \
enc_gap_p95={:.1}ms enc_gap_max={:.1}ms \
sent_gap_p95={:.1}ms sent_gap_max={:.1}ms \
frame_age_p95={:.1}ms frame_age_max={:.1}ms \
send_wait_p95={:.1}ms \
import_p95={:.1}ms scale_p95={:.1}ms transfer_p95={:.1}ms \
sws_p95={:.1}ms encode_p95={:.1}ms total_p95={:.1}ms \
output_bps={:.0} frame_bytes_max={}",
"elapsed={:.1}s capture_fps={:.1} encoded_fps={:.1} sent_fps={:.1} \
capture_frames={} encoded_frames={} sent_frames={} \
pw_dropped={} duplicate_frames_skipped={} \
cap_q={} enc_q={} \
cap_gap_avg={:.1}ms cap_gap_p95={:.1}ms cap_gap_max={:.1}ms \
enc_gap_avg={:.1}ms enc_gap_p95={:.1}ms enc_gap_max={:.1}ms \
sent_gap_avg={:.1}ms sent_gap_p95={:.1}ms sent_gap_max={:.1}ms \
frame_age_avg={:.1}ms frame_age_p95={:.1}ms frame_age_max={:.1}ms \
send_wait_p95={:.1}ms \
import_avg={:.1}ms import_p95={:.1}ms \
scale_avg={:.1}ms scale_p95={:.1}ms transfer_avg={:.1}ms transfer_p95={:.1}ms \
sws_avg={:.1}ms sws_p95={:.1}ms \
encode_avg={:.1}ms encode_p95={:.1}ms total_avg={:.1}ms total_p95={:.1}ms \
output_bps={:.0} frame_bytes_p95={} frame_bytes_max={}",
self.elapsed_secs,
self.capture_fps,
self.encoded_fps,
self.sent_fps,
self.capture_frames,
self.encoded_frames,
self.sent_frames,
self.pipewire_dropped,
self.over_budget_count,
self.duplicate_frames_skipped,
self.capture_queue_depth,
self.encoded_queue_depth,
self.capture_gap_avg_ms,
self.capture_gap_p95_ms,
self.capture_gap_max_ms,
self.encoded_gap_avg_ms,
self.encoded_gap_p95_ms,
self.encoded_gap_max_ms,
self.sent_gap_avg_ms,
self.sent_gap_p95_ms,
self.sent_gap_max_ms,
self.frame_age_avg_ms,
self.frame_age_p95_ms,
self.frame_age_max_ms,
self.send_wait_p95_ms,
self.import_avg_ms,
self.import_p95_ms,
self.scale_avg_ms,
self.scale_p95_ms,
self.transfer_avg_ms,
self.transfer_p95_ms,
self.sws_avg_ms,
self.sws_p95_ms,
self.encode_avg_ms,
self.encode_p95_ms,
self.total_avg_ms,
self.total_p95_ms,
self.output_bytes_per_sec,
self.output_frame_bytes_p95,
self.output_frame_bytes_max,
)
}
+9 -311
View File
@@ -1,8 +1,12 @@
/// Coordinate transformation module for Wayland output transforms.
///
/// Handles the 8 `wl_output` transform variants (rotation + reflection)
/// and ROI clipping for screen capture.
///
//! Coordinate transformation module for Wayland output transforms.
//!
//! Historically exposed a family of `Rect`/`screen_to_frame`/`fit_inside_bounds`
//! helpers for ROI-based capture clipping. Those were never wired into the
//! capture pipeline (we capture full frames and let FFmpeg's filter graph handle
//! any scaling/rotation); they have been removed. Only `Transform` and the
//! `transpose_if_transform_transposed` helper remain — both are actively used by
//! `state.rs` and `avhw.rs`.
/// Wayland output transform enum, matching `wl_output::Transform`.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Transform {
@@ -16,68 +20,6 @@ pub enum Transform {
Flipped270,
}
/// Axis-aligned rectangle in integer coordinates.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Rect {
pub x: i32,
pub y: i32,
pub w: i32,
pub h: i32,
}
/// Returns the 2×2 basis matrix (a, b, c, d) for the given transform.
///
/// The matrix represents the affine mapping from screen coordinates to
/// frame coordinates:
///
/// ```text
/// [new_x] [a b] [x]
/// [new_y] = [c d] [y]
/// ```
pub fn transform_basis(transform: Transform) -> (i32, i32, i32, i32) {
match transform {
Transform::Normal => (1, 0, 0, 1),
Transform::Normal90 => (0, 1, -1, 0),
Transform::Normal180 => (-1, 0, 0, -1),
Transform::Normal270 => (0, -1, 1, 0),
Transform::Flipped => (-1, 0, 0, 1),
Transform::Flipped90 => (0, 1, 1, 0),
Transform::Flipped180 => (1, 0, 0, -1),
Transform::Flipped270 => (0, -1, -1, 0),
}
}
/// Transform a rectangle from screen space to frame space.
///
/// Applies the 2×2 basis matrix and computes offsets so the result
/// fits within the frame dimensions `(frame_w, frame_h)`.
///
/// ```text
/// new_x = a * x + b * y + offset_x
/// new_y = c * x + d * y + offset_y
/// ```
pub fn screen_to_frame(transform: Transform, rect: Rect, frame_w: i32, frame_h: i32) -> Rect {
let (a, b, c, d) = transform_basis(transform);
// Compute the offset so that the transformed origin maps correctly.
// For transforms with negative components, we need to shift by the
// frame dimension to keep coordinates in [0, frame_w) × [0, frame_h).
let offset_x = if a + b < 0 { frame_w } else { 0 };
let offset_y = if c + d < 0 { frame_h } else { 0 };
let new_x = a * rect.x + b * rect.y + offset_x;
let new_y = c * rect.x + d * rect.y + offset_y;
let new_w = a * rect.w + b * rect.h;
let new_h = c * rect.w + d * rect.h;
Rect {
x: new_x,
y: new_y,
w: new_w.abs(),
h: new_h.abs(),
}
}
/// Swap width and height for 90° or 270° rotations.
///
/// After a quarter-turn rotation the output dimensions are transposed
@@ -93,140 +35,10 @@ pub fn transpose_if_transform_transposed(transform: Transform, w: i32, h: i32) -
}
}
/// Clip a rectangle so it stays inside `(0, 0) .. (bounds_w, bounds_h)`.
///
/// The resulting rectangle has non-negative origin and its extent does
/// not exceed the bounds.
pub fn fit_inside_bounds(rect: Rect, bounds_w: i32, bounds_h: i32) -> Rect {
let x = rect.x.clamp(0, bounds_w);
let y = rect.y.clamp(0, bounds_h);
let right = (rect.x + rect.w).min(bounds_w);
let bottom = (rect.y + rect.h).min(bounds_h);
let w = (right - x).max(0);
let h = (bottom - y).max(0);
Rect { x, y, w, h }
}
#[cfg(test)]
mod tests {
use super::*;
// ── transform_basis ───────────────────────────────────────────
#[test]
fn basis_normal_is_identity() {
assert_eq!(transform_basis(Transform::Normal), (1, 0, 0, 1));
}
#[test]
fn basis_90_cw_rotation() {
assert_eq!(transform_basis(Transform::Normal90), (0, 1, -1, 0));
}
#[test]
fn basis_180_rotation() {
assert_eq!(transform_basis(Transform::Normal180), (-1, 0, 0, -1));
}
#[test]
fn basis_270_cw_rotation() {
assert_eq!(transform_basis(Transform::Normal270), (0, -1, 1, 0));
}
#[test]
fn basis_flipped_horizontal() {
assert_eq!(transform_basis(Transform::Flipped), (-1, 0, 0, 1));
}
#[test]
fn basis_flipped_90() {
assert_eq!(transform_basis(Transform::Flipped90), (0, 1, 1, 0));
}
#[test]
fn basis_flipped_180() {
assert_eq!(transform_basis(Transform::Flipped180), (1, 0, 0, -1));
}
#[test]
fn basis_flipped_270() {
assert_eq!(transform_basis(Transform::Flipped270), (0, -1, -1, 0));
}
// ── screen_to_frame ───────────────────────────────────────────
#[test]
fn screen_to_frame_identity_unchanged() {
let rect = Rect {
x: 10,
y: 20,
w: 100,
h: 50,
};
let result = screen_to_frame(Transform::Normal, rect, 1920, 1080);
assert_eq!(
result,
Rect {
x: 10,
y: 20,
w: 100,
h: 50
}
);
}
#[test]
fn screen_to_frame_90_rotates_origin() {
// 90° CW: top-left (0,0) in screen should map to bottom-left in frame
let rect = Rect {
x: 0,
y: 0,
w: 100,
h: 50,
};
let result = screen_to_frame(Transform::Normal90, rect, 1080, 1920);
// a=0,b=1,c=-1,d=0 => offset_x=0, offset_y=1920 (c+d=-1<0)
// new_x = 0*0 + 1*0 + 0 = 0
// new_y = -1*0 + 0*0 + 1920 = 1920
assert_eq!(result.x, 0);
assert_eq!(result.y, 1920);
// w' = 0*100 + 1*50 = 50, h' = -1*100 + 0*50 = -100 -> abs=100
assert_eq!(result.w, 50);
assert_eq!(result.h, 100);
}
#[test]
fn screen_to_frame_180_rotates() {
let rect = Rect {
x: 100,
y: 200,
w: 300,
h: 400,
};
let result = screen_to_frame(Transform::Normal180, rect, 1920, 1080);
// a=-1,b=0,c=0,d=-1, offset_x=1920, offset_y=1080
assert_eq!(result.x, -100 + 1920);
assert_eq!(result.y, -200 + 1080);
assert_eq!(result.w, 300);
assert_eq!(result.h, 400);
}
#[test]
fn screen_to_frame_flipped_horizontal() {
let rect = Rect {
x: 50,
y: 30,
w: 200,
h: 100,
};
let result = screen_to_frame(Transform::Flipped, rect, 1920, 1080);
// a=-1,b=0,c=0,d=1, offset_x=1920, offset_y=0
assert_eq!(result.x, -50 + 1920);
assert_eq!(result.y, 30);
assert_eq!(result.w, 200);
assert_eq!(result.h, 100);
}
// ── transpose_if_transform_transposed ─────────────────────────
#[test]
@@ -292,118 +104,4 @@ mod tests {
(1080, 1920)
);
}
// ── fit_inside_bounds ─────────────────────────────────────────
#[test]
fn fit_inside_already_fits() {
let rect = Rect {
x: 10,
y: 20,
w: 100,
h: 50,
};
let result = fit_inside_bounds(rect, 1920, 1080);
assert_eq!(result, rect);
}
#[test]
fn fit_inside_clips_right_and_bottom() {
let rect = Rect {
x: 1800,
y: 1000,
w: 200,
h: 200,
};
let result = fit_inside_bounds(rect, 1920, 1080);
assert_eq!(
result,
Rect {
x: 1800,
y: 1000,
w: 120,
h: 80
}
);
}
#[test]
fn fit_inside_clips_negative_origin() {
let rect = Rect {
x: -50,
y: -30,
w: 200,
h: 200,
};
let result = fit_inside_bounds(rect, 1920, 1080);
assert_eq!(
result,
Rect {
x: 0,
y: 0,
w: 150,
h: 170
}
);
}
#[test]
fn fit_inside_completely_out_of_bounds() {
let rect = Rect {
x: 2000,
y: 2000,
w: 100,
h: 100,
};
let result = fit_inside_bounds(rect, 1920, 1080);
assert_eq!(
result,
Rect {
x: 1920,
y: 1080,
w: 0,
h: 0
}
);
}
#[test]
fn fit_inside_zero_size_rect() {
let rect = Rect {
x: 100,
y: 100,
w: 0,
h: 0,
};
let result = fit_inside_bounds(rect, 1920, 1080);
assert_eq!(
result,
Rect {
x: 100,
y: 100,
w: 0,
h: 0
}
);
}
#[test]
fn fit_inside_zero_bounds() {
let rect = Rect {
x: 0,
y: 0,
w: 100,
h: 100,
};
let result = fit_inside_bounds(rect, 0, 0);
assert_eq!(
result,
Rect {
x: 0,
y: 0,
w: 0,
h: 0
}
);
}
}
+1 -1
View File
@@ -536,7 +536,7 @@ impl WebRtcInner {
let now = Instant::now();
let should_honor = self
.last_forced_keyframe_at
.map_or(true, |last| now.duration_since(last) >= FORCED_KEYFRAME_MIN_INTERVAL);
.is_none_or(|last| now.duration_since(last) >= FORCED_KEYFRAME_MIN_INTERVAL);
if should_honor {
self.last_forced_keyframe_at = Some(now);
self.need_keyframe = true;