diff --git a/.gitea/workflows/cargo-deny.yml b/.gitea/workflows/cargo-deny.yml new file mode 100644 index 0000000..bfefc94 --- /dev/null +++ b/.gitea/workflows/cargo-deny.yml @@ -0,0 +1,34 @@ +name: cargo-deny + +# Enforce the supply-chain policy in deny.toml (advisories / bans / licenses / +# sources) on every push to main and every PR. Runs on a *locked* tree so the +# pinned, vetted versions in Cargo.lock are exactly what get audited — see the +# deny.toml header and VERSIONING.md. A new poisoned release of a dependency +# cannot reach CI until Cargo.lock is deliberately updated. + +on: + push: + branches: [main] + pull_request: + +jobs: + cargo-deny: + runs-on: ubuntu-latest + # rust:1 provides the cargo toolchain that cargo-deny shells out to for + # `cargo metadata`. Adjust the runner label if your act_runner uses a + # different one. + container: rust:1 + steps: + - uses: actions/checkout@v4 + + - name: Install cargo-deny (pinned prebuilt) + run: | + set -euo pipefail + version=0.19.9 + curl -sSfL \ + "https://github.com/EmbarkStudios/cargo-deny/releases/download/${version}/cargo-deny-${version}-x86_64-unknown-linux-musl.tar.gz" \ + | tar -xz -C /usr/local/bin --strip-components=1 --wildcards '*/cargo-deny' + cargo-deny --version + + - name: cargo deny check + run: cargo deny --locked check diff --git a/.gitea/workflows/windows-build.yml b/.gitea/workflows/windows-build.yml new file mode 100644 index 0000000..59dc481 --- /dev/null +++ b/.gitea/workflows/windows-build.yml @@ -0,0 +1,85 @@ +name: windows-build + +# Milestone M1 of the Windows port (see docs/handoff windows-migration-plan): +# prove the tree compiles for `x86_64-pc-windows-msvc` and the unit tests pass. +# The audio backend is the Phase 0 `CpalBackend` stub for now — this job guards +# the *compile* boundary (cfg gating, platform deps, the PlatformAudioBackend +# alias) so a Unix-only assumption can't sneak back in and break Windows. +# +# RUNNER REQUIREMENT: this needs a Windows act_runner registered with the +# `windows-latest` label (the Linux `cargo-deny` job's container approach does +# NOT apply here — Windows jobs run on the host, not a Linux container). If your +# runner advertises a different label, change `runs-on` below. Until a Windows +# runner exists this workflow is simply skipped/queued, not a failure of the +# Linux CI. +# +# BUILD-HOST REQUIREMENTS (validated by the opus spike, see +# peerspeak-windows-opus-spike.md): +# - MSVC C toolchain (Visual Studio Build Tools) — to compile vendored libopus. +# - CMake on PATH — `audiopus_sys` builds libopus from source via cmake. +# - CMAKE_POLICY_VERSION_MINIMUM=3.5 (set below) — the vendored libopus declares +# an ancient `cmake_minimum_required` that CMake >= 4.0 refuses without it. +# GitHub-hosted `windows-latest` images ship MSVC + CMake; a self-hosted runner +# must provide both. + +on: + push: + # `main` plus the in-progress port branches, so the Windows path is exercised + # before merge rather than only after. + branches: [main, "windows-port-**"] + pull_request: + # Allow manual runs from the Gitea Actions UI. + workflow_dispatch: + +permissions: + contents: read + +env: + CARGO_TERM_COLOR: always + # The vendored libopus (audiopus_sys -> cmake) uses cmake_minimum_required < 3.5, + # which CMake 4.x rejects unless this is set. See the opus spike report. + CMAKE_POLICY_VERSION_MINIMUM: "3.5" + +jobs: + windows-build: + runs-on: windows-latest + steps: + - uses: actions/checkout@v4 + + - name: Install Rust (MSVC, pinned to repo toolchain if present) + uses: dtolnay/rust-toolchain@stable + with: + targets: x86_64-pc-windows-msvc + components: clippy + + - name: Show toolchain + build prerequisites + shell: bash + run: | + set -euo pipefail + rustc --version + cargo --version + # libopus is built from source via cmake; fail early with a clear + # message if the runner lacks it rather than deep in the opus build. + if ! command -v cmake >/dev/null 2>&1; then + echo "::error::cmake not found on PATH. The opus crate builds libopus from source via cmake; install CMake on this runner." + exit 1 + fi + cmake --version + + # Build on a *locked* tree so the pinned, vetted Cargo.lock versions are what + # get compiled — same supply-chain stance as the cargo-deny job. + - name: Build (all targets, msvc) + run: cargo build --all-targets --locked --target x86_64-pc-windows-msvc + + # Unit (lib) tests only: the `transport_loopback` integration tests stand up + # real iroh/QUIC endpoints and need working loopback networking, which isn't + # guaranteed on a CI runner. Add `--tests` here once a networked Windows + # runner is confirmed. + - name: Unit tests (lib, msvc) + run: cargo test --lib --locked --target x86_64-pc-windows-msvc + + # Informational for now (not `-D warnings`): the Windows tree may surface + # platform-specific lints we haven't triaged. Tighten to deny-warnings once + # it's clean. + - name: Clippy (msvc) + run: cargo clippy --all-targets --locked --target x86_64-pc-windows-msvc diff --git a/Cargo.lock b/Cargo.lock index a177314..c7e880b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -105,6 +105,28 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" +[[package]] +name = "alsa" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed7572b7ba83a31e20d1b48970ee402d2e3e0537dcfe0a3ff4d6eb7508617d43" +dependencies = [ + "alsa-sys", + "bitflags 2.11.1", + "cfg-if", + "libc", +] + +[[package]] +name = "alsa-sys" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "db8fee663d06c4e303404ef5f40488a53e062f89ba8bfed81f42325aafad1527" +dependencies = [ + "libc", + "pkg-config", +] + [[package]] name = "android-activity" version = "0.6.1" @@ -114,12 +136,12 @@ dependencies = [ "android-properties", "bitflags 2.11.1", "cc", - "jni", + "jni 0.22.4", "libc", "log", - "ndk", + "ndk 0.9.0", "ndk-context", - "ndk-sys", + "ndk-sys 0.6.0+11769913", "num_enum", "thiserror 2.0.18", ] @@ -712,6 +734,12 @@ dependencies = [ "shlex", ] +[[package]] +name = "cesu8" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6d43a04d8753f35258c91f8ec639f792891f748a1edbd759cf1dcea3382ad83c" + [[package]] name = "cexpr" version = "0.6.0" @@ -1002,6 +1030,26 @@ dependencies = [ "libm", ] +[[package]] +name = "coreaudio-rs" +version = "0.11.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "321077172d79c662f64f5071a03120748d5bb652f5231570141be24cfcd2bace" +dependencies = [ + "bitflags 1.3.2", + "core-foundation-sys", + "coreaudio-sys", +] + +[[package]] +name = "coreaudio-sys" +version = "0.2.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9b4739a805a62757a83e5654fa3faabec0442666b263bb2287d5a8185bfd953" +dependencies = [ + "bindgen", +] + [[package]] name = "cosmic-text" version = "0.15.0" @@ -1026,6 +1074,29 @@ dependencies = [ "unicode-segmentation", ] +[[package]] +name = "cpal" +version = "0.15.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "873dab07c8f743075e57f524c583985fbaf745602acbe916a01539364369a779" +dependencies = [ + "alsa", + "core-foundation-sys", + "coreaudio-rs", + "dasp_sample", + "jni 0.21.1", + "js-sys", + "libc", + "mach2", + "ndk 0.8.0", + "ndk-context", + "oboe", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", + "windows 0.54.0", +] + [[package]] name = "cpufeatures" version = "0.2.17" @@ -1228,6 +1299,12 @@ dependencies = [ "syn", ] +[[package]] +name = "dasp_sample" +version = "0.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c87e182de0887fd5361989c677c4e8f5000cd9491d6d563161a8f3a5519fc7f" + [[package]] name = "data-encoding" version = "2.11.0" @@ -2227,7 +2304,7 @@ dependencies = [ "http", "idna", "ipnet", - "jni", + "jni 0.22.4", "rand 0.10.1", "rustls", "thiserror 2.0.18", @@ -2247,7 +2324,7 @@ dependencies = [ "data-encoding", "idna", "ipnet", - "jni", + "jni 0.22.4", "once_cell", "prefix-trie", "rand 0.10.1", @@ -2270,7 +2347,7 @@ dependencies = [ "hickory-proto", "ipconfig", "ipnet", - "jni", + "jni 0.22.4", "moka", "ndk-context", "once_cell", @@ -3102,6 +3179,22 @@ version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" +[[package]] +name = "jni" +version = "0.21.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1a87aa2bb7d2af34197c04845522473242e1aa17c12f4935d5856491a7fb8c97" +dependencies = [ + "cesu8", + "cfg-if", + "combine", + "jni-sys 0.3.1", + "log", + "thiserror 1.0.69", + "walkdir", + "windows-sys 0.45.0", +] + [[package]] name = "jni" version = "0.22.4" @@ -3463,6 +3556,15 @@ version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d3d25b0e0b648a86960ac23b7ad4abb9717601dec6f66c165f5b037f3f03065f" +[[package]] +name = "mach2" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d640282b302c0bb0a2a8e0233ead9035e3bed871f0b7e81fe4a1ec829765db44" +dependencies = [ + "libc", +] + [[package]] name = "malloc_buf" version = "0.0.6" @@ -3596,7 +3698,7 @@ dependencies = [ "dispatch", "futures-channel", "futures-lite", - "jni", + "jni 0.22.4", "ndk-context", "objc2 0.6.4", "objc2-app-kit 0.3.2", @@ -3695,6 +3797,20 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "ndk" +version = "0.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2076a31b7010b17a38c01907c45b945e8f11495ee4dd588309718901b1f7a5b7" +dependencies = [ + "bitflags 2.11.1", + "jni-sys 0.3.1", + "log", + "ndk-sys 0.5.0+25.2.9519653", + "num_enum", + "thiserror 1.0.69", +] + [[package]] name = "ndk" version = "0.9.0" @@ -3704,7 +3820,7 @@ dependencies = [ "bitflags 2.11.1", "jni-sys 0.3.1", "log", - "ndk-sys", + "ndk-sys 0.6.0+11769913", "num_enum", "raw-window-handle", "thiserror 1.0.69", @@ -3716,6 +3832,15 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "27b02d87554356db9e9a873add8782d4ea6e3e58ea071a9adb9a2e8ddb884a8b" +[[package]] +name = "ndk-sys" +version = "0.5.0+25.2.9519653" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8c196769dd60fd4f363e11d948139556a344e79d451aeb2fa2fd040738ef7691" +dependencies = [ + "jni-sys 0.3.1", +] + [[package]] name = "ndk-sys" version = "0.6.0+11769913" @@ -4466,6 +4591,29 @@ dependencies = [ "objc2-foundation 0.2.2", ] +[[package]] +name = "oboe" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8b61bebd49e5d43f5f8cc7ee2891c16e0f41ec7954d36bcb6c14c5e0de867fb" +dependencies = [ + "jni 0.21.1", + "ndk 0.8.0", + "ndk-context", + "num-derive", + "num-traits", + "oboe-sys", +] + +[[package]] +name = "oboe-sys" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6c8bb09a4a2b1d668170cfe0a7d5bc103f8999fb316c98099b6a9939c9f2e79d" +dependencies = [ + "cc", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -4600,6 +4748,7 @@ dependencies = [ "async-trait", "base64", "bytes", + "cpal", "dirs", "iced", "image", @@ -5436,7 +5585,7 @@ checksum = "26d1e2536ce4f35f4846aa13bff16bd0ff40157cdb14cc056c7b14ba41233ba0" dependencies = [ "core-foundation 0.10.1", "core-foundation-sys", - "jni", + "jni 0.22.4", "log", "once_cell", "rustls", @@ -5877,7 +6026,7 @@ dependencies = [ "fastrand", "js-sys", "memmap2", - "ndk", + "ndk 0.9.0", "objc2 0.6.4", "objc2-core-foundation", "objc2-core-graphics", @@ -7162,7 +7311,7 @@ dependencies = [ "log", "metal", "naga", - "ndk-sys", + "ndk-sys 0.6.0+11769913", "objc", "once_cell", "ordered-float", @@ -7247,6 +7396,16 @@ dependencies = [ "thiserror 2.0.18", ] +[[package]] +name = "windows" +version = "0.54.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9252e5725dbed82865af151df558e754e4a3c2c30818359eb17465f1346a1b49" +dependencies = [ + "windows-core 0.54.0", + "windows-targets 0.52.6", +] + [[package]] name = "windows" version = "0.58.0" @@ -7254,7 +7413,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dd04d41d93c4992d421894c18c8b43496aa748dd4c081bac0dc93eb0489272b6" dependencies = [ "windows-core 0.58.0", - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -7278,6 +7437,16 @@ dependencies = [ "windows-core 0.62.2", ] +[[package]] +name = "windows-core" +version = "0.54.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "12661b9c89351d684a50a8a643ce5f608e20243b9fb84687800163429f161d65" +dependencies = [ + "windows-result 0.1.2", + "windows-targets 0.52.6", +] + [[package]] name = "windows-core" version = "0.58.0" @@ -7288,7 +7457,7 @@ dependencies = [ "windows-interface 0.58.0", "windows-result 0.2.0", "windows-strings 0.1.0", - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -7386,13 +7555,22 @@ dependencies = [ "windows-strings 0.5.1", ] +[[package]] +name = "windows-result" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e383302e8ec8515204254685643de10811af0ed97ea37210dc26fb0032647f8" +dependencies = [ + "windows-targets 0.52.6", +] + [[package]] name = "windows-result" version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d1043d8214f791817bab27572aaa8af63732e11bf84aa21a45a78d6c317ae0e" dependencies = [ - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -7411,7 +7589,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4cd9b125c486025df0eabcb585e62173c6c9eddcec5d117d3b6e8c30e2ee4d10" dependencies = [ "windows-result 0.2.0", - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -7423,13 +7601,22 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows-sys" +version = "0.45.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75283be5efb2831d37ea142365f009c02ec203cd29a3ebecbc093d52315b66d0" +dependencies = [ + "windows-targets 0.42.2", +] + [[package]] name = "windows-sys" version = "0.52.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "282be5f36a8ce781fad8c8ae18fa3f9beff57ec1b52cb3de0789201425d9a33d" dependencies = [ - "windows-targets", + "windows-targets 0.52.6", ] [[package]] @@ -7441,20 +7628,35 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows-targets" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e5180c00cd44c9b1c88adb3693291f1cd93605ded80c250a75d472756b4d071" +dependencies = [ + "windows_aarch64_gnullvm 0.42.2", + "windows_aarch64_msvc 0.42.2", + "windows_i686_gnu 0.42.2", + "windows_i686_msvc 0.42.2", + "windows_x86_64_gnu 0.42.2", + "windows_x86_64_gnullvm 0.42.2", + "windows_x86_64_msvc 0.42.2", +] + [[package]] name = "windows-targets" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b724f72796e036ab90c1021d4780d4d3d648aca59e491e6b98e725b84e99973" dependencies = [ - "windows_aarch64_gnullvm", - "windows_aarch64_msvc", - "windows_i686_gnu", + "windows_aarch64_gnullvm 0.52.6", + "windows_aarch64_msvc 0.52.6", + "windows_i686_gnu 0.52.6", "windows_i686_gnullvm", - "windows_i686_msvc", - "windows_x86_64_gnu", - "windows_x86_64_gnullvm", - "windows_x86_64_msvc", + "windows_i686_msvc 0.52.6", + "windows_x86_64_gnu 0.52.6", + "windows_x86_64_gnullvm 0.52.6", + "windows_x86_64_msvc 0.52.6", ] [[package]] @@ -7466,18 +7668,36 @@ dependencies = [ "windows-link", ] +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "597a5118570b68bc08d8d59125332c54f1ba9d9adeedeef5b99b02ba2b0698f8" + [[package]] name = "windows_aarch64_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" +[[package]] +name = "windows_aarch64_msvc" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e08e8864a60f06ef0d0ff4ba04124db8b0fb3be5776a5cd47641e942e58c4d43" + [[package]] name = "windows_aarch64_msvc" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" +[[package]] +name = "windows_i686_gnu" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c61d927d8da41da96a81f029489353e68739737d3beca43145c8afec9a31a84f" + [[package]] name = "windows_i686_gnu" version = "0.52.6" @@ -7490,24 +7710,48 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" +[[package]] +name = "windows_i686_msvc" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "44d840b6ec649f480a41c8d80f9c65108b92d89345dd94027bfe06ac444d1060" + [[package]] name = "windows_i686_msvc" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" +[[package]] +name = "windows_x86_64_gnu" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8de912b8b8feb55c064867cf047dda097f92d51efad5b491dfb98f6bbb70cb36" + [[package]] name = "windows_x86_64_gnu" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26d41b46a36d453748aedef1486d5c7a85db22e56aff34643984ea85514e94a3" + [[package]] name = "windows_x86_64_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" +[[package]] +name = "windows_x86_64_msvc" +version = "0.42.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9aec5da331524158c6d1a4ac0ab1541149c0b9505fde06423b02f5ef0106b9f0" + [[package]] name = "windows_x86_64_msvc" version = "0.52.6" @@ -7536,7 +7780,7 @@ dependencies = [ "js-sys", "libc", "memmap2", - "ndk", + "ndk 0.9.0", "objc2 0.5.2", "objc2-app-kit 0.2.2", "objc2-foundation 0.2.2", diff --git a/Cargo.toml b/Cargo.toml index 4cbfedc..f9bfde3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -30,17 +30,12 @@ bytes = "1.11.1" dirs = "6.0.0" iced = { version = "0.14.0", features = ["canvas", "image"] } # W4 custom avatars: decode/resize an arbitrary user image (png/jpeg only to keep -# the codec surface small) and a native file picker (xdg-portal backend, no GTK). +# the codec surface small). The matching native file picker (`rfd`) is platform- +# gated below — its backend differs per OS (xdg-portal on Linux, Win32 on Windows). image = { version = "0.25", default-features = false, features = ["png", "jpeg"] } -rfd = { version = "0.17", default-features = false, features = ["xdg-portal"] } iroh = "1.0.0-rc.0" iroh-gossip = "0.99.0" opus = "0.3.1" -# v0_3_49 exposes `Buffer::requested()` (the graph's per-cycle quantum), used by -# the playback RT callback to fill exactly what the device asks for instead of -# pinning the buffer to a hard-coded 1024-frame quantum (crackle on non-1024 -# hardware). The field has existed in libpipewire since 0.3.49 (2022). -pipewire = { version = "0.9", features = ["v0_3_49"] } rand = "0.10.1" ringbuf = "0.5.0" serde = { version = "1.0.228", features = ["derive"] } @@ -48,3 +43,24 @@ serde_json = "1.0.150" thiserror = "2.0.18" tokio = { version = "1.52.3", features = ["full"] } tokio-stream = "0.1.18" + +# --- Platform-specific dependencies ----------------------------------------- +# Audio and the native file-picker backends differ per OS. Everything else in the +# app talks to the `AudioBackend` trait and the `PlatformAudioBackend` alias (see +# `src/audio/mod.rs`), so platform selection is confined to these few lines. + +[target.'cfg(target_os = "linux")'.dependencies] +# Linux audio backend. v0_3_49 exposes `Buffer::requested()` (the graph's per-cycle +# quantum), used by the playback RT callback to fill exactly what the device asks +# for instead of a hard-coded 1024-frame quantum (crackle on non-1024 hardware). +# The field has existed in libpipewire since 0.3.49 (2022). +pipewire = { version = "0.9", features = ["v0_3_49"] } +# Native file picker via the XDG desktop portal (no GTK) on Linux. +rfd = { version = "0.17", default-features = false, features = ["xdg-portal"] } + +[target.'cfg(windows)'.dependencies] +# Native file picker using the built-in Win32 dialog backend on Windows. +rfd = { version = "0.17", default-features = false } +# Windows audio backend: cpal drives WASAPI for capture/playback behind the +# AudioBackend trait (src/audio/cpal_impl.rs). The Linux counterpart is pipewire. +cpal = "0.15" diff --git a/docs/WINDOWS.md b/docs/WINDOWS.md new file mode 100644 index 0000000..37ca748 --- /dev/null +++ b/docs/WINDOWS.md @@ -0,0 +1,77 @@ +# PeerSpeak on Windows + +Current status: the Windows port cross-compiles to `x86_64-pc-windows-gnu` and the `.exe` +launches under Wine. A real Windows/WASAPI host is still needed for the final audio-device +checks listed below. + +## What works today + +| Area | Status | +|---|---| +| GUI | Iced/wgpu builds and renders under Wine. | +| Networking | Iroh QUIC transport and gossip compile on Windows. | +| Audio backend | `cpal` drives WASAPI capture/playback behind `AudioBackend`. | +| Codec | Opus remains 48 kHz mono, 20 ms frames. | +| Identity | `ring` identity generation/load is platform-neutral. | +| Chimes | Windows uses PowerShell `System.Media.SoundPlayer` for WAV playback. | + +Windows paths are resolved through `dirs`: + +- Config: `%APPDATA%\peerspeak\config.json` +- Identity: `%APPDATA%\peerspeak\identity.key` +- Log: `%LOCALAPPDATA%\peerspeak\peerspeak.log` + +## Building + +### Native Windows + +Install MSVC Build Tools and CMake, then build normally: + +```powershell +cargo build --release +``` + +If CMake is 4.x or newer, the vendored `opus`/`libopus` build may need: + +```powershell +$env:CMAKE_POLICY_VERSION_MINIMUM = "3.5" +cargo build --release +``` + +### Cross-compile from Linux + +The current dev path cross-compiles from an Arch environment to the GNU Windows target: + +```sh +rustup target add x86_64-pc-windows-gnu +sudo pacman -S mingw-w64-gcc cmake +CMAKE_POLICY_VERSION_MINIMUM=3.5 cargo build --release --target x86_64-pc-windows-gnu --bin peerspeak +``` + +Wine is useful for launch/render smoke tests, but it is not a substitute for a real +Windows audio-device pass. The deeper migration plan (phases, decisions, the opus build +spike) lives in the maintainer's handoff docs, outside the repo. + +## First run and networking + +Expect a Windows Firewall prompt the first time the app opens network sockets. Allow it: +PeerSpeak uses UDP for QUIC, plus relay traffic when direct NAT traversal is not available. + +The default network mode keeps the n0 relay available for NAT traversal without publishing +presence to n0 DNS. Direct peer-to-peer paths may work when both networks allow them; relayed +connections are expected and valid. + +## Known gaps + +| Item | Status | +|---|---| +| Echo cancellation | Linux-only PipeWire feature. The Windows UI shows it disabled as unavailable. | +| Screen share | Requires a Windows `pixelpass.exe` on `PATH` or a configured override. | +| Chimes | Now routed through Windows `SoundPlayer`; needs a real Windows host to audibly verify. | +| Resampling/device format | Cross-compiled. cpal/WASAPI now chooses native 48 kHz when available and otherwise resamples/remaps at the device boundary; needs real Windows hardware audio verification. | +| Device persistence | Open. WASAPI friendly names may duplicate or change across driver/profile changes. | +| Playback pacing | Cross-compiled. The fixed playback target under WASAPI shared mode still needs real-hardware verification with `audio_probe`. | + +Before calling Windows support done, verify a real Windows machine can create/join a room, +capture mic audio, hear remote audio, select devices, restart with selections preserved, and +play notification chimes. diff --git a/src/app/mod.rs b/src/app/mod.rs index 8352bfc..fc6dfe5 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -2,7 +2,7 @@ use crate::core::{CoreController, messages::{CoreCommand, UiEvent}}; use crate::network::PeerState; use crate::notify::{self, Sound}; use crate::audio::eq::{EqSettings, EQ_GAIN_DB_MAX, EQ_GAIN_DB_MIN}; -use crate::audio::pw_cli::{AudioDevice, enumerate_audio_devices}; +use crate::audio::{AudioDevice, enumerate_audio_devices}; use crate::config::{AppConfig, NetworkMode, RecordingMode, RoomLayout}; use crate::hotkeys::{format_binding, HotkeyAction, HotkeyContext, KeyBinding}; use crate::presence::PresenceMode; @@ -553,11 +553,9 @@ pub fn run_gui() -> iced::Result { // the icon from the .desktop file matched by app_id instead). icon: window_icon(), // app_id must match the .desktop basename so Wayland compositors - // (e.g. KWin) attach our launcher icon to the window. - platform_specific: iced::window::settings::PlatformSpecific { - application_id: "peerspeak".to_string(), - ..Default::default() - }, + // (e.g. KWin) attach our launcher icon to the window. The field is + // Linux-only in iced (X11/Wayland); see platform_specific_settings(). + platform_specific: platform_specific_settings(), // We save the final size ourselves on CloseRequested, then exit. exit_on_close_request: false, ..Default::default() @@ -565,6 +563,22 @@ pub fn run_gui() -> iced::Result { .run() } +/// Window `PlatformSpecific` settings. `application_id` (used by X11/Wayland to +/// match our `.desktop` launcher icon) only exists in iced on Linux, so it is +/// set there and left at defaults on Windows. +#[cfg(target_os = "linux")] +fn platform_specific_settings() -> iced::window::settings::PlatformSpecific { + iced::window::settings::PlatformSpecific { + application_id: "peerspeak".to_string(), + ..Default::default() + } +} + +#[cfg(not(target_os = "linux"))] +fn platform_specific_settings() -> iced::window::settings::PlatformSpecific { + iced::window::settings::PlatformSpecific::default() +} + /// Build the window icon from an embedded 128×128 straight-RGBA blob rendered /// from `assets/icons/peerspeak.svg`. Using `from_rgba` (always available) keeps /// us off iced's heavy `image` feature — the blob is raw pixels, no decoder. @@ -1350,12 +1364,28 @@ fn update(state: &mut AppState, message: AppMessage) -> Task { // Defence in depth: only ever hand http(s) URLs to the opener. The // link span's href came from `linkify`, which only emits http/https, // but re-check here so this can't be widened into launching arbitrary - // schemes/args. `xdg-open` receives the URL as a single argv entry - // (no shell), so there's no injection surface. - if (url.starts_with("http://") || url.starts_with("https://")) - && let Err(e) = std::process::Command::new("xdg-open").arg(&url).spawn() - { - crate::log_msg(&format!("Failed to open URL {url:?}: {e}")); + // schemes/args. Each opener receives the URL as a single argv entry + // (no shell), so there's no injection surface: + // - Unix: `xdg-open `. + // - Windows: `rundll32 url.dll,FileProtocolHandler ` — opens the + // default browser without going through `cmd`/`start`, which would + // otherwise re-parse `&` in query strings. + if url.starts_with("http://") || url.starts_with("https://") { + let spawned = { + #[cfg(unix)] + { + std::process::Command::new("xdg-open").arg(&url).spawn() + } + #[cfg(windows)] + { + std::process::Command::new("rundll32") + .args(["url.dll,FileProtocolHandler", &url]) + .spawn() + } + }; + if let Err(e) = spawned { + crate::log_msg(&format!("Failed to open URL {url:?}: {e}")); + } } } AppMessage::ToggleMicTest(enabled) => { @@ -2541,10 +2571,28 @@ fn view(state: &AppState) -> Element<'_, AppMessage> { mic_meter, text("Drag the yellow handle to set the gate. While talking, place it just above your quiet-room level so silence is muted but your voice passes through.").size(11).color(color_subtext), vertical_space(4.0), - checkbox(state.config.echo_cancellation_enabled) - .label("Echo cancellation") - .on_toggle(AppMessage::ToggleEchoCancellation), - text("Cancels speaker echo + suppresses noise (PipeWire). Takes effect on your next room join.").size(11).color(color_subtext), + { + let control: Element<'_, AppMessage> = { + #[cfg(target_os = "linux")] + { + column![ + checkbox(state.config.echo_cancellation_enabled) + .label("Echo cancellation") + .on_toggle(AppMessage::ToggleEchoCancellation), + text("Cancels speaker echo + suppresses noise (PipeWire). Takes effect on your next room join.").size(11).color(color_subtext), + ].spacing(8).into() + } + #[cfg(not(target_os = "linux"))] + { + column![ + checkbox(false) + .label("Echo cancellation"), + text("Echo cancellation is not available on Windows yet.").size(11).color(color_subtext), + ].spacing(8).into() + } + }; + control + }, ].spacing(8).width(iced::Length::Fill), ] .spacing(10) @@ -3247,26 +3295,40 @@ fn view(state: &AppState) -> Element<'_, AppMessage> { column![] }, vertical_space(20.0), - // Echo cancellation — same flag + message as the Settings checkbox, so - // toggling here and there stay in sync automatically (single source of - // truth: config.echo_cancellation_enabled). Tooltip is explicit that it - // applies on the NEXT join (the PipeWire-module AEC is wired at join - // time, not hot-swappable mid-call). - tooltip( - checkbox(state.config.echo_cancellation_enabled) - .label("Echo cancellation") - .on_toggle(AppMessage::ToggleEchoCancellation), - container( - text("Cancels speaker echo + suppresses noise. Applies on your next room join.") - .size(11) - .color(color_text), - ) - .padding(8) - .max_width(260.0) - .style(c_style(color_crust, color_surface, 6.0)), - iced::widget::tooltip::Position::Top, - ) - .gap(8), + { + // Echo cancellation is wired at join time on Linux; other + // targets show an inert status row instead of a dead toggle. + let control: Element<'_, AppMessage> = { + #[cfg(target_os = "linux")] + { + tooltip( + checkbox(state.config.echo_cancellation_enabled) + .label("Echo cancellation") + .on_toggle(AppMessage::ToggleEchoCancellation), + container( + text("Cancels speaker echo + suppresses noise. Applies on your next room join.") + .size(11) + .color(color_text), + ) + .padding(8) + .max_width(260.0) + .style(c_style(color_crust, color_surface, 6.0)), + iced::widget::tooltip::Position::Top, + ) + .gap(8) + .into() + } + #[cfg(not(target_os = "linux"))] + { + column![ + checkbox(false) + .label("Echo cancellation"), + text("Not available on Windows yet.").size(11).color(color_subtext), + ].spacing(4).into() + } + }; + control + }, vertical_space(20.0), { let (rec_kind, rec_label, rec_bg, rec_hover, rec_fg) = if state.recording { diff --git a/src/audio/cpal_impl.rs b/src/audio/cpal_impl.rs new file mode 100644 index 0000000..4ae5b7d --- /dev/null +++ b/src/audio/cpal_impl.rs @@ -0,0 +1,1445 @@ +//! Windows audio backend — cpal / WASAPI (Phase 1). +//! +//! Implements [`AudioBackend`] on top of [`cpal`], which wraps WASAPI on Windows. +//! It is the Windows counterpart to `pipewire_impl.rs` and deliberately preserves +//! the exact same contract so the rest of the app (mixer, encoder, jitter buffer) +//! is unchanged: +//! +//! - **Capture**: mono, 48 kHz, S16 PCM, emitted as `Vec` frames of +//! [`CAPTURE_FRAME`] (960 = 20 ms) samples — matching the encoder/jitter frame. +//! The RT capture callback only downmixes and pushes samples into a lock-free +//! ring; the owning thread drains that ring, frames it, and sends — so the +//! callback never allocates, locks, or touches an mpsc channel. +//! - **Playback**: stereo interleaved ([`PLAYBACK_CHANNELS`]) S16 PCM at 48 kHz, +//! drained from a ring buffer that is paced to the device's hardware clock via +//! `ring_fill` exactly as the PipeWire backend does. +//! +//! ## Threading and the `!Send` stream +//! +//! `cpal::Stream` is `!Send` (some backends require it to be created and dropped +//! on the same thread), but [`AudioBackend`] is `Send + Sync` and the backend is +//! shared through an `Arc`. So the stream never lives in the struct: each of +//! `start_capture`/`start_playback` spawns one owning thread that builds the +//! stream, plays it, and keeps it alive until the per-worker `running` flag flips +//! (set by `stop`). The struct holds only `Send` handles (the flag + the join +//! handle). The stream's RT callback does the actual audio work; the owning +//! thread additionally feeds the playback ring (or drains the capture ring). +//! +//! `start_*` does not return until the owning thread reports back over a readiness +//! channel that the device resolved and the stream is built and playing — so a +//! device/format/WASAPI failure surfaces as a real `Err` to the caller instead of +//! leaving the UI in a joined-but-silent room. +//! +//! ## Sample rate and channel layout (W4) +//! +//! The whole pipeline runs internally at 48 kHz (Opus + the 960-sample frame) and +//! mono capture / stereo playback. We prefer a native-48 kHz device config so the +//! common case is conversion-free and bit-exact. When the device can't do 48 kHz +//! (commonly a 44.1 kHz-only endpoint) or can't do stereo output, we fall back to +//! the device's default config and convert at the boundary with the dep-free +//! [`super::resample`] linear resamplers instead of hard-erroring: +//! +//! - **Capture**: the device-rate mono stream is resampled to 48 kHz on the +//! capture drain thread (off the RT callback) before framing. +//! - **Playback**: the internal 48 kHz stereo bus is resampled to the device rate +//! and remapped to the device channel count inside the output RT callback, which +//! pulls internal frames from the ring on demand (allocation-free, so RT-safe). +//! The ring, prefill, and `ring_fill` pacing stay in internal 48 kHz-stereo +//! units, so the mixer is unchanged. +//! +//! Linear interpolation has no anti-aliasing filter (see [`super::resample`] docs); +//! it is adequate for speech and keeps the matching-rate path bit-exact, with the +//! seam ready for a higher-quality resampler later. + +use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, AtomicUsize, Ordering}; +use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender, channel}; +use std::sync::{Arc, Mutex}; +use std::thread::{self, JoinHandle}; +use std::time::{Duration, Instant}; + +use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; +use cpal::{ + Device, FromSample, Sample, SampleFormat, SampleRate, SizedSample, Stream, StreamConfig, +}; +use ringbuf::{ + HeapRb, + traits::{Consumer, Producer, Split}, +}; + +use super::resample::{PushResampler, StereoPullResampler}; +use super::{AudioBackend, AudioDevice, AudioError, PLAYBACK_CHANNELS, PLAYBACK_TARGET_SAMPLES}; + +/// The one sample rate the pipeline supports (Opus + the 20 ms frame). +const SAMPLE_RATE: u32 = 48_000; +/// Mono capture frame: 960 samples = 20 ms @ 48 kHz. Matches the PipeWire backend +/// and `core::jitter::FRAME_SAMPLES`. +const CAPTURE_FRAME: usize = 960; +/// Lock-free capture ring capacity (mono samples) between the RT callback and the +/// owning drain thread: 8 frames = 160 ms of headroom, so a scheduling hiccup on +/// the drain thread doesn't immediately overrun the RT producer. +const CAPTURE_RING_CAPACITY: usize = CAPTURE_FRAME * 8; +/// How long the capture drain thread sleeps when the ring is momentarily empty, +/// before polling again. Small enough to stay well under the 20 ms frame cadence. +const CAPTURE_POLL: Duration = Duration::from_millis(5); +/// Playback ring capacity in interleaved samples: 200 ms of stereo @ 48 kHz. +/// Comfortably above [`PLAYBACK_TARGET_SAMPLES`] so the clock-paced producer has +/// headroom and never has to drop frames in steady state. +const RING_CAPACITY: usize = 9600 * PLAYBACK_CHANNELS; +/// How often a blocked playback worker re-checks its `running` flag, bounding how +/// long `stop()` can take to join it (mirrors the PipeWire backend's `WORKER_POLL`). +const WORKER_POLL: Duration = Duration::from_millis(100); +/// How long a freshly-played stream has to prove itself (deliver its first RT +/// callbacks) before the start is treated as failed. cpal's `play()` only *queues* +/// the WASAPI `Start()`, so a queued-but-failed start would otherwise masquerade as +/// success and join the UI into a silent room (review W1). +const STREAM_START_TIMEOUT: Duration = Duration::from_secs(3); +/// Number of completed RT callbacks the owner waits for before declaring the stream +/// live. One callback isn't proof: a stream can fire once and immediately fail in +/// the same processing cycle, so requiring a couple of cycles (plus the terminal +/// error check) keeps a one-shot-then-dead stream from being reported Ok (B1). +const MIN_START_CALLBACKS: usize = 2; +/// Backstop for [`finish_start`]: bounds the WHOLE owner path (resolve + build + +/// play + the [`STREAM_START_TIMEOUT`] callback wait + any wedged-stream cleanup). +/// Sized as a generous setup budget plus the callback wait plus slack so a slow but +/// valid device (e.g. a Bluetooth endpoint that takes seconds to spin up) is not +/// falsely failed, while a driver that wedges before the owner can report is still +/// released eventually (review W6, B4). +const FINISH_START_TIMEOUT: Duration = Duration::from_secs(10); +/// Lowest / highest device sample rate the backend will drive. The floor bounds +/// the playback pull-resampler's input-pulls-per-output-frame (≈48000/rate) so a +/// pathological low rate can't blow the RT callback's deadline; the ceiling and a +/// nonzero floor also reject the 0 Hz / absurd values a misbehaving driver could +/// report, which would otherwise panic or spin (review W7). +const MIN_DEVICE_RATE: u32 = 8_000; +const MAX_DEVICE_RATE: u32 = 384_000; + +// Stream-error categories carried from the RT error callback to the owner thread +// through an `AtomicU8`, so the callback itself never allocates or logs — both of +// which it previously did via `format!`/`log_msg` on the time-critical stream +// thread (review W2). The owner/logger translates the code off the RT path. +const STREAM_ERR_NONE: u8 = 0; +const STREAM_ERR_DEVICE_UNAVAILABLE: u8 = 1; +const STREAM_ERR_BACKEND: u8 = 2; + +/// Map a cpal stream error to its [`STREAM_ERR_*`](STREAM_ERR_NONE) code. Pure + +/// allocation-free, so it is safe to call from the RT error callback. +fn stream_err_code(e: &cpal::StreamError) -> u8 { + match e { + cpal::StreamError::DeviceNotAvailable => STREAM_ERR_DEVICE_UNAVAILABLE, + _ => STREAM_ERR_BACKEND, + } +} + +/// Human-readable text for a [`STREAM_ERR_*`](STREAM_ERR_NONE) code, logged off the +/// RT path by the owner/health-logger thread. +fn stream_err_text(code: u8) -> &'static str { + match code { + STREAM_ERR_DEVICE_UNAVAILABLE => "audio device became unavailable", + _ => "audio backend stream error", + } +} + +/// Wait for a just-played stream to prove it actually started: its RT data +/// callback bumps `callbacks`, or an error callback sets `err_code`. Returns `Ok` +/// once [`MIN_START_CALLBACKS`] cycles have run with no error, `Err` on an +/// error-callback code or [`STREAM_START_TIMEOUT`], or a clean abort if `stop()` +/// cleared `running` mid-start. Polls a few cheap atomics on the owner thread — +/// never the RT thread (review W1). +/// +/// The error is **terminal and wins any race** with `callbacks`: a stream can run a +/// callback and then fail in the same processing cycle, so `err_code` is checked +/// first each loop AND re-checked before declaring success (review B1). +fn wait_for_stream_start( + callbacks: &AtomicUsize, + err_code: &AtomicU8, + running: &AtomicBool, +) -> Result<(), AudioError> { + let deadline = Instant::now() + STREAM_START_TIMEOUT; + let as_err = |code: u8| Err(AudioError::Stream(stream_err_text(code).to_string())); + loop { + let code = err_code.load(Ordering::Relaxed); + if code != STREAM_ERR_NONE { + return as_err(code); + } + if callbacks.load(Ordering::Relaxed) >= MIN_START_CALLBACKS { + // Re-check: a callback that pushed us to the threshold may have been the + // last before a same-cycle failure. Let a terminal error win. + let code = err_code.load(Ordering::Relaxed); + if code != STREAM_ERR_NONE { + return as_err(code); + } + return Ok(()); + } + if !running.load(Ordering::Relaxed) { + return Err(AudioError::Stream("stream start aborted".to_string())); + } + if Instant::now() >= deadline { + return Err(AudioError::Stream( + "stream did not start within timeout (no WASAPI callback)".to_string(), + )); + } + thread::sleep(Duration::from_millis(5)); + } +} + +/// Windows audio backend. See module docs. +pub struct CpalBackend { + capture: Mutex, + playback: Mutex, +} + +/// A spawned owning thread plus the flags that coordinate its lifetime: `running` +/// tells it to drop its stream and exit; `exited` is flipped true (by [`ExitGuard`] +/// in the thread body) when it actually returns, so a *detached* wedged start can be +/// detected as finished later (review B3). +struct StreamWorker { + running: Arc, + exited: Arc, + thread: JoinHandle<()>, +} + +/// Flips its flag true when dropped, marking a worker thread as exited. Lives at the +/// top of the worker closure so it fires on normal return, panic unwind, or whenever +/// a wedged driver call finally releases the thread — which is what lets a [`SlotState::Wedged`] +/// tombstone (B3) know its orphan is gone. +struct ExitGuard(Arc); +impl Drop for ExitGuard { + fn drop(&mut self) { + self.0.store(true, Ordering::Relaxed); + } +} + +/// The lifecycle state of a capture or playback slot. +enum SlotState { + /// No stream — a new start may proceed. + Idle, + /// A live, started stream owned by its worker thread. + Live(StreamWorker), + /// A start that timed out wedged in a driver call (review B3). Its worker thread + /// was *detached* rather than joined — joining would re-introduce the unbounded + /// hang [`FINISH_START_TIMEOUT`] exists to prevent — so it may still be alive, + /// holding the COM/device handle. `exited` flips true when that orphan finally + /// returns. New starts are rejected until then, so retries against a permanently + /// wedged device don't pile up more orphan threads. + Wedged { exited: Arc }, +} + +/// Inspect a slot before starting a stream into it. Clears a [`SlotState::Wedged`] +/// tombstone whose orphan has since exited (the slot becomes reusable), but rejects +/// a start while a wedged orphan is still alive or a live stream already owns the +/// slot. Pure w.r.t. the passed state, so the tombstone logic is unit-testable (B3). +fn ensure_idle(state: &mut SlotState, what: &str) -> Result<(), AudioError> { + match state { + SlotState::Idle => Ok(()), + SlotState::Live(_) => Err(AudioError::Stream(format!("{what} already started"))), + SlotState::Wedged { exited } => { + if exited.load(Ordering::Relaxed) { + *state = SlotState::Idle; + Ok(()) + } else { + Err(AudioError::Stream(format!( + "{what} is recovering from an unresponsive audio device; retry shortly" + ))) + } + } + } +} + +impl CpalBackend { + pub fn new() -> Self { + Self { + capture: Mutex::new(SlotState::Idle), + playback: Mutex::new(SlotState::Idle), + } + } +} + +impl Default for CpalBackend { + fn default() -> Self { + Self::new() + } +} + +impl AudioBackend for CpalBackend { + fn start_capture( + &self, + tx: Sender>, + target_node: Option, + ) -> Result<(), AudioError> { + let mut guard = self.capture.lock().unwrap(); + ensure_idle(&mut guard, "capture")?; + let running = Arc::new(AtomicBool::new(true)); + let exited = Arc::new(AtomicBool::new(false)); + let running_thread = running.clone(); + let exited_thread = exited.clone(); + let (ready_tx, ready_rx) = channel::>(); + let thread = thread::Builder::new() + .name("peerspeak-cpal-capture".to_string()) + .spawn(move || { + let _exit = ExitGuard(exited_thread); + run_capture(tx, target_node, running_thread, ready_tx); + }) + .map_err(|e| AudioError::Init(e.to_string()))?; + finish_start( + guard, + StreamWorker { + running, + exited, + thread, + }, + ready_rx, + "capture", + ) + } + + fn start_playback( + &self, + rx: Receiver>, + target_node: Option, + ring_fill: Arc, + ) -> Result<(), AudioError> { + let mut guard = self.playback.lock().unwrap(); + ensure_idle(&mut guard, "playback")?; + let running = Arc::new(AtomicBool::new(true)); + let exited = Arc::new(AtomicBool::new(false)); + let running_thread = running.clone(); + let exited_thread = exited.clone(); + let (ready_tx, ready_rx) = channel::>(); + let thread = thread::Builder::new() + .name("peerspeak-cpal-playback".to_string()) + .spawn(move || { + let _exit = ExitGuard(exited_thread); + run_playback(rx, target_node, ring_fill, running_thread, ready_tx); + }) + .map_err(|e| AudioError::Init(e.to_string()))?; + finish_start( + guard, + StreamWorker { + running, + exited, + thread, + }, + ready_rx, + "playback", + ) + } + + fn stop(&self) -> Result<(), AudioError> { + for slot in [&self.capture, &self.playback] { + let mut guard = slot.lock().unwrap(); + match std::mem::replace(&mut *guard, SlotState::Idle) { + SlotState::Live(worker) => { + worker.running.store(false, Ordering::Relaxed); + let _ = worker.thread.join(); + } + // A wedged orphan was detached and can't be joined. If it has since + // exited the slot is now clear; otherwise restore the tombstone so a + // later start still sees the device is recovering (B3). + SlotState::Wedged { exited } => { + if !exited.load(Ordering::Relaxed) { + *guard = SlotState::Wedged { exited }; + } + } + SlotState::Idle => {} + } + } + Ok(()) + } +} + +/// Block until the just-spawned worker reports (over `ready_rx`) that its stream +/// is built and playing, then either install it (`Ok`) or join it and surface the +/// real error. This is what makes `start_capture`/`start_playback` fail loudly +/// instead of returning `Ok` into a joined-but-silent room (Codex review W1). +fn finish_start( + mut guard: std::sync::MutexGuard<'_, SlotState>, + worker: StreamWorker, + ready_rx: Receiver>, + what: &str, +) -> Result<(), AudioError> { + // Bounded wait. An unbounded `recv()` here would hang `start_*` forever — and + // any concurrent `stop()` behind the same slot mutex — if a WASAPI/driver call + // wedged the worker before it could report (review W6). + match ready_rx.recv_timeout(FINISH_START_TIMEOUT) { + Ok(Ok(())) => { + *guard = SlotState::Live(worker); + Ok(()) + } + // Setup failed (Err) or the worker disconnected before reporting: either + // way it has stopped, so reap it and surface the error. + Ok(Err(e)) => { + worker.running.store(false, Ordering::Relaxed); + let _ = worker.thread.join(); + Err(e) + } + Err(RecvTimeoutError::Disconnected) => { + worker.running.store(false, Ordering::Relaxed); + let _ = worker.thread.join(); + Err(AudioError::Init(format!( + "cpal {what} worker exited before reporting readiness" + ))) + } + Err(RecvTimeoutError::Timeout) => { + // The worker is wedged in a driver call. Signal it to exit, but DETACH + // rather than join — joining would re-introduce the unbounded hang this + // timeout exists to prevent. Leave a Wedged tombstone so subsequent + // starts are rejected until the orphan's ExitGuard flips `exited`, rather + // than spawning more orphan threads against the same dead device (B3). + let StreamWorker { + running, + exited, + thread, + } = worker; + running.store(false, Ordering::Relaxed); + drop(thread); + *guard = SlotState::Wedged { exited }; + Err(AudioError::Init(format!( + "cpal {what} did not start within {FINISH_START_TIMEOUT:?}" + ))) + } + } +} + +// --------------------------------------------------------------------------- +// Device enumeration (for the settings device pickers) +// --------------------------------------------------------------------------- + +/// Enumerate WASAPI input/output devices via cpal, sorted by description to match +/// the PipeWire backend's stable UI ordering. +/// +/// cpal exposes a single friendly name per device, which is also what [`resolve`] +/// matches `target_node` against — so `name` and `description` are the same string +/// and a saved selection round-trips. Note: WASAPI device names are less stable +/// across driver/endpoint changes than PipeWire node names, so a saved device may +/// not always be found again; selection then falls back to the system default. +pub fn enumerate_audio_devices() -> Vec { + let host = cpal::default_host(); + let mut devices = Vec::new(); + + if let Ok(inputs) = host.input_devices() { + for device in inputs { + if let Ok(name) = device.name() { + devices.push(AudioDevice { + description: name.clone(), + name, + is_input: true, + }); + } + } + } + if let Ok(outputs) = host.output_devices() { + for device in outputs { + if let Ok(name) = device.name() { + devices.push(AudioDevice { + description: name.clone(), + name, + is_input: false, + }); + } + } + } + + devices.sort_by(|a, b| a.description.cmp(&b.description)); + devices +} + +// --------------------------------------------------------------------------- +// Device / config selection +// --------------------------------------------------------------------------- + +/// Resolve a device (by `target` name, else the system default) and a stream +/// config. We prefer a config running natively at [`SAMPLE_RATE`] (conversion-free); +/// if the device has none, we fall back to its default config and resample/remap at +/// the boundary (W4 — see module docs and [`choose_config`]). +fn resolve( + output: bool, + target: Option, +) -> Result<(Device, StreamConfig, SampleFormat), AudioError> { + let host = cpal::default_host(); + + let default = || { + if output { + host.default_output_device() + } else { + host.default_input_device() + } + }; + let device = match target { + // A saved device name that no longer resolves falls back to the system + // default — but log it, because WASAPI friendly names can change across + // driver/endpoint changes, so a silent fallback otherwise looks like + // "audio went to the wrong device for no reason" (review W7). + Some(ref name) => match find_device_by_name(&host, output, name) { + Some(dev) => Some(dev), + None => { + crate::log_msg(&format!( + "cpal: saved {} device '{name}' not found; using system default", + if output { "output" } else { "input" }, + )); + default() + } + }, + None => default(), + } + .ok_or_else(|| AudioError::Device("no audio device available".to_string()))?; + + let supported = choose_config(&device, output)?; + let sample_format = supported.sample_format(); + let config = supported.config(); + + // Validate the OS-reported geometry before any code divides by it or sizes a + // loop from it (review W7). Zero channels would panic `chunks_exact(0)` / + // `chunks_mut(0)`; a zero or absurd rate would yield an infinite/huge resample + // ratio. Reject up front with a real error instead of panicking or spinning. + if config.channels == 0 { + return Err(AudioError::Device( + "audio device reports zero channels".to_string(), + )); + } + if !(MIN_DEVICE_RATE..=MAX_DEVICE_RATE).contains(&config.sample_rate.0) { + return Err(AudioError::Device(format!( + "audio device sample rate {} Hz is outside the supported {MIN_DEVICE_RATE}–{MAX_DEVICE_RATE} Hz range", + config.sample_rate.0, + ))); + } + Ok((device, config, sample_format)) +} + +fn find_device_by_name(host: &cpal::Host, output: bool, name: &str) -> Option { + let devices = if output { + host.output_devices().ok()? + } else { + host.input_devices().ok()? + }; + devices + .into_iter() + .find(|d| d.name().is_ok_and(|n| n == name)) +} + +/// Whether the backend can actually open this config range. The workers only build +/// `F32`/`I16`/`U16` streams ([`build_input`]/[`build_output`] — every other sample +/// format hits the `other => Err(...)` arm), and a zero-channel range would later be +/// rejected by [`resolve`]'s geometry check. Filtering both here keeps +/// [`choose_config`] from *ranking* a range it can't drive ahead of a usable one and +/// then hard-failing the start instead of trying the next candidate (Codex B5 +/// re-review, P3). +fn usable_range(r: &cpal::SupportedStreamConfigRange) -> bool { + r.channels() > 0 && format_supported(r.sample_format()) +} + +/// Sample formats the capture/playback stream builders accept. Pure, so it's +/// unit-testable independently of the cpal range types. +fn format_supported(fmt: SampleFormat) -> bool { + matches!( + fmt, + SampleFormat::F32 | SampleFormat::I16 | SampleFormat::U16 + ) +} + +/// Pick a sample rate inside both a device's supported `[r_min, r_max]` span and the +/// backend's drivable `[MIN_DEVICE_RATE, MAX_DEVICE_RATE]` window, preferring +/// [`SAMPLE_RATE`] when it's reachable and otherwise the nearest in-window bound. +/// Returns `None` when the device span doesn't overlap the window at all. Pure and +/// integer-only, so the selection policy is unit-testable (review B5). +fn bounded_rate(r_min: u32, r_max: u32) -> Option { + let lo = r_min.max(MIN_DEVICE_RATE); + let hi = r_max.min(MAX_DEVICE_RATE); + (lo <= hi).then(|| SAMPLE_RATE.clamp(lo, hi)) +} + +/// Pick a stream config. Preference order, best (no conversion) first: +/// 1. exactly [`SAMPLE_RATE`] at the preferred layout (stereo out / mono in), +/// 2. exactly [`SAMPLE_RATE`] at any channel count (rate-exact, backend remaps), +/// 3. a supported config at a [`bounded_rate`] near 48 kHz (backend resamples + remaps), +/// 4. the device's default config (only if nothing above is drivable). +/// +/// Cases 3–4 incur resampling; the backend reads the returned config's rate and +/// channel count and converts at the boundary (W4). Case 3 (review B5) is what keeps +/// an oddball endpoint whose default rate is outside the drivable window — but which +/// also exposes a usable in-window config — from being rejected by [`resolve`]. A +/// device that exposes no config at all is still a hard error. +fn choose_config(device: &Device, output: bool) -> Result { + let ranges: Vec = if output { + device + .supported_output_configs() + .map_err(|e| AudioError::Device(e.to_string()))? + .collect() + } else { + device + .supported_input_configs() + .map_err(|e| AudioError::Device(e.to_string()))? + .collect() + }; + + // A range covers a sample-rate span and a fixed channel count. + let supports_48k = |r: &cpal::SupportedStreamConfigRange| { + r.min_sample_rate().0 <= SAMPLE_RATE && SAMPLE_RATE <= r.max_sample_rate().0 + }; + let pick = |channels: Option| { + ranges + .iter() + .find(|r| usable_range(r) && supports_48k(r) && channels.is_none_or(|c| r.channels() == c)) + .cloned() + }; + + // Cases 1 + 2: an exact-48 kHz config, preferring the native layout but + // accepting any channel count (the backend remaps without resampling). + let exact = if output { + pick(Some(PLAYBACK_CHANNELS as u16)).or_else(|| pick(None)) + } else { + pick(Some(1)).or_else(|| pick(None)) + }; + if let Some(r) = exact { + return Ok(r.with_sample_rate(SampleRate(SAMPLE_RATE))); + } + + // Case 3: no native 48 kHz. Before falling back to the device default — which + // resolve() rejects outright if its rate is outside the drivable window — look + // for a supported config whose rate range overlaps that window and drive it at a + // bounded rate, resampling at the boundary (review B5). Prefer the native layout, + // then the bounded rate closest to 48 kHz. + let pick_bounded = |channels: Option| -> Option<(cpal::SupportedStreamConfigRange, u32)> { + ranges + .iter() + .filter(|r| usable_range(r) && channels.is_none_or(|c| r.channels() == c)) + .filter_map(|r| { + bounded_rate(r.min_sample_rate().0, r.max_sample_rate().0) + .map(|rate| (r.clone(), rate)) + }) + .min_by_key(|(_, rate)| rate.abs_diff(SAMPLE_RATE)) + }; + let preferred_channels = if output { PLAYBACK_CHANNELS as u16 } else { 1 }; + if let Some((r, rate)) = pick_bounded(Some(preferred_channels)).or_else(|| pick_bounded(None)) { + crate::log_msg(&format!( + "cpal: device '{}' has no native {SAMPLE_RATE} Hz {} config; using bounded {rate} Hz / {} ch with linear resampling (W4/B5)", + device.name().unwrap_or_else(|_| "".to_string()), + if output { "output" } else { "input" }, + r.channels(), + )); + return Ok(r.with_sample_rate(SampleRate(rate))); + } + + // Case 4: last resort — the device's default config. If its rate is outside the + // drivable window, resolve() rejects it with a clear device error, which is the + // honest outcome: the device exposes nothing this backend can drive. + let def = if output { + device.default_output_config() + } else { + device.default_input_config() + } + .map_err(|e| AudioError::Device(e.to_string()))?; + crate::log_msg(&format!( + "cpal: device '{}' has no bounded {} config near {SAMPLE_RATE} Hz; falling back to default {} Hz / {} ch (W4)", + device.name().unwrap_or_else(|_| "".to_string()), + if output { "output" } else { "input" }, + def.sample_rate().0, + def.channels(), + )); + Ok(def) +} + +// --------------------------------------------------------------------------- +// Capture +// --------------------------------------------------------------------------- + +fn run_capture( + tx: Sender>, + target: Option, + running: Arc, + ready: Sender>, +) { + // The RT callback pushes mono samples into this lock-free ring; we drain it on + // this (non-RT) thread, so the callback never allocates or sends on a channel. + let rb = HeapRb::::new(CAPTURE_RING_CAPACITY); + let (producer, mut consumer) = rb.split(); + let overrun = Arc::new(AtomicU64::new(0)); + // Stream-liveness signals read by `wait_for_stream_start`: each RT data callback + // bumps `callbacks`; the RT error callback sets `err_code` (it does NOT log — + // that would allocate/syscall on the time-critical thread). See W1/W2/B1. + let callbacks = Arc::new(AtomicUsize::new(0)); + let err_code = Arc::new(AtomicU8::new(STREAM_ERR_NONE)); + + // Fallible device/stream setup: resolve, build, and *queue* the stream start. + let setup = || -> Result<(Stream, String, SampleFormat, usize, u32), AudioError> { + let (device, config, sample_format) = resolve(false, target)?; + let channels = config.channels as usize; + let device_rate = config.sample_rate.0; + let stream = match sample_format { + SampleFormat::F32 => build_input::( + &device, &config, producer, channels, overrun.clone(), callbacks.clone(), + err_code.clone(), + ), + SampleFormat::I16 => build_input::( + &device, &config, producer, channels, overrun.clone(), callbacks.clone(), + err_code.clone(), + ), + SampleFormat::U16 => build_input::( + &device, &config, producer, channels, overrun.clone(), callbacks.clone(), + err_code.clone(), + ), + other => Err(AudioError::Stream(format!( + "unsupported capture sample format: {other:?}" + ))), + }?; + stream + .play() + .map_err(|e| AudioError::Stream(e.to_string()))?; + let name = device.name().unwrap_or_else(|_| "".to_string()); + Ok((stream, name, sample_format, channels, device_rate)) + }; + + // Build + queue, then wait for the stream to actually prove it's live before + // reporting readiness. `play()` returning Ok only means WASAPI's `Start()` was + // queued; a later Start failure would otherwise leave us joined-but-silent (W1). + let (stream, dev_name, sample_format, channels, device_rate) = match setup() { + Ok(v) => match wait_for_stream_start(&callbacks, &err_code, &running) { + Ok(()) => { + let _ = ready.send(Ok(())); + v + } + Err(e) => { + // Drop the (possibly wedged) stream BEFORE reporting: cpal's + // Stream::drop joins its WASAPI worker, so if that wedges we want + // the Err withheld and finish_start's backstop to detach, rather + // than finish_start joining this owner forever (B2). + drop(v.0); + let _ = ready.send(Err(e)); + return; + } + }, + Err(e) => { + let _ = ready.send(Err(e)); + return; + } + }; + crate::log_msg(&format!( + "cpal capture started: device='{dev_name}' format={sample_format:?} channels={channels} device_rate={device_rate} Hz -> {SAMPLE_RATE} Hz" + )); + + // If the device isn't at 48 kHz, resample its mono stream up/down to 48 kHz on + // this (non-RT) thread before framing (W4). At 48 kHz this stays None and the + // samples pass straight through, bit-exact. + let mut resampler = + (device_rate != SAMPLE_RATE).then(|| PushResampler::new(device_rate, SAMPLE_RATE)); + // Reused scratch for a sample's resampled output (off-RT alloc; tiny — at most + // a couple of samples per input). Avoids a nested-closure borrow over `acc`/`tx`. + let mut resampled: Vec = Vec::new(); + + // Drain the RT ring on this thread: pop mono samples, (resample,) frame them + // (the `Vec` allocation lives here, off the RT path), and send completed + // frames. Keep `stream` alive until `stop()` flips the flag. + let mut acc = FrameAccumulator::new(CAPTURE_FRAME); + let mut last_overrun = 0u64; + let mut last_err = STREAM_ERR_NONE; + while running.load(Ordering::Relaxed) { + let mut drained = false; + while let Some(sample) = consumer.try_pop() { + drained = true; + resampled.clear(); + match resampler { + Some(ref mut rs) => { + rs.push(i16_to_f32(sample), |out| resampled.push(f32_to_i16(out))); + } + None => resampled.push(sample), + } + for s in resampled.drain(..) { + if let Some(frame) = acc.push(s) { + // Consumer gone (call ended) → stop feeding; the stream is + // dropped below on the way out. + if tx.send(frame).is_err() { + drop(stream); + return; + } + } + } + } + let o = overrun.load(Ordering::Relaxed); + if o != last_overrun { + crate::log_msg(&format!( + "cpal capture overrun: dropped {} samples (drain thread fell behind)", + o - last_overrun + )); + last_overrun = o; + } + // Surface a stream error the RT callback flagged (it can't log itself). + let ec = err_code.load(Ordering::Relaxed); + if ec != STREAM_ERR_NONE && ec != last_err { + crate::log_msg(&format!("cpal capture stream error: {}", stream_err_text(ec))); + last_err = ec; + } + if !drained { + thread::sleep(CAPTURE_POLL); + } + } + drop(stream); +} + +#[allow(clippy::too_many_arguments)] +fn build_input( + device: &Device, + config: &StreamConfig, + mut producer: P, + channels: usize, + overrun: Arc, + callbacks: Arc, + err_code: Arc, +) -> Result +where + T: SizedSample + Send + 'static, + i16: FromSample, + P: Producer + Send + 'static, +{ + // RT-safe error callback: record a category in an atomic only. Formatting + + // logging happen on the owner thread (the cpal/WASAPI error callback runs on + // the time-critical stream thread, where alloc/syscall are forbidden — W2). + let err_fn = move |e: cpal::StreamError| err_code.store(stream_err_code(&e), Ordering::Relaxed); + device + .build_input_stream::( + config, + move |data: &[T], _| { + // Count callbacks so the owner can confirm the stream is really + // running before reporting Ok (W1/B1). + callbacks.fetch_add(1, Ordering::Relaxed); + // RT-safe: downmix + wait-free push only. A full ring means the + // drain thread stalled; count the drop and keep going. + for frame in data.chunks_exact(channels) { + let mono = downmix_to_mono(frame); + if producer.try_push(mono).is_err() { + overrun.fetch_add(1, Ordering::Relaxed); + } + } + }, + err_fn, + None, + ) + .map_err(|e| AudioError::Stream(e.to_string())) +} + +/// Average a device frame's channels down to a single mono i16. For a 1-channel +/// device this is just the converted sample. +fn downmix_to_mono(frame: &[T]) -> i16 +where + T: Copy, + i16: FromSample, +{ + if frame.is_empty() { + return 0; + } + let sum: i32 = frame.iter().map(|&s| i16::from_sample(s) as i32).sum(); + (sum / frame.len() as i32) as i16 +} + +/// Scale an i16 PCM sample to f32 in roughly `[-1, 1]` for interpolation. +#[inline] +fn i16_to_f32(s: i16) -> f32 { + s as f32 / 32768.0 +} + +/// Convert an interpolated f32 sample back to i16, clamping to range. +#[inline] +fn f32_to_i16(x: f32) -> i16 { + (x * 32768.0).clamp(i16::MIN as f32, i16::MAX as f32) as i16 +} + +/// Accumulates mono samples into fixed-size [`CAPTURE_FRAME`] frames. Pulled out +/// of the RT callback so the framing is unit-testable. +struct FrameAccumulator { + buf: Vec, + frame_len: usize, +} + +impl FrameAccumulator { + fn new(frame_len: usize) -> Self { + Self { + buf: Vec::with_capacity(frame_len), + frame_len, + } + } + + /// Push one sample; returns a completed frame when the buffer fills. + fn push(&mut self, sample: i16) -> Option> { + self.buf.push(sample); + if self.buf.len() == self.frame_len { + Some(std::mem::replace( + &mut self.buf, + Vec::with_capacity(self.frame_len), + )) + } else { + None + } + } +} + +// --------------------------------------------------------------------------- +// Playback +// --------------------------------------------------------------------------- + +fn run_playback( + rx: Receiver>, + target: Option, + ring_fill: Arc, + running: Arc, + ready: Sender>, +) { + let rb = HeapRb::::new(RING_CAPACITY); + let (mut producer, consumer) = rb.split(); + + // Prefill to the steady-state depth so playout starts at target. `ring_fill` + // is an EXACT occupancy counter maintained by deltas (worker fetch_add on + // push, RT callback fetch_sub on pop) — not ringbuf's cached `occupied_len`, + // which is stale across the split halves and would lie high and starve the + // ring. See pipewire_impl.rs for the full rationale. + for _ in 0..PLAYBACK_TARGET_SAMPLES { + let _ = producer.try_push(0); + } + ring_fill.store(PLAYBACK_TARGET_SAMPLES, Ordering::Relaxed); + + // Diagnostics (mirrors the PipeWire backend's playout-health line). + let underrun = Arc::new(AtomicU64::new(0)); + let dropped = Arc::new(AtomicU64::new(0)); + // Largest single output-callback length seen (interleaved samples). WASAPI + // shared-mode picks its own period, so this can exceed the prefill target — + // which would force an underrun every cycle (review W2). The callback only + // does a wait-free fetch_max; the health logger reports/warns off the RT path. + let max_cb = Arc::new(AtomicUsize::new(0)); + // Stream-liveness signals (see the capture path / W1, W2, B1): each RT callback + // bumps `callbacks`; the RT error callback sets `err_code` without logging. + let callbacks = Arc::new(AtomicUsize::new(0)); + let err_code = Arc::new(AtomicU8::new(STREAM_ERR_NONE)); + + // Fallible device/stream setup. `consumer` is moved into the output callback. + let setup = || -> Result<(Stream, String, SampleFormat, usize, u32), AudioError> { + let (device, config, sample_format) = resolve(true, target)?; + let channels = config.channels as usize; + let device_rate = config.sample_rate.0; + let stream = match sample_format { + SampleFormat::F32 => build_output::( + &device, &config, consumer, ring_fill.clone(), underrun.clone(), + max_cb.clone(), callbacks.clone(), err_code.clone(), + ), + SampleFormat::I16 => build_output::( + &device, &config, consumer, ring_fill.clone(), underrun.clone(), + max_cb.clone(), callbacks.clone(), err_code.clone(), + ), + SampleFormat::U16 => build_output::( + &device, &config, consumer, ring_fill.clone(), underrun.clone(), + max_cb.clone(), callbacks.clone(), err_code.clone(), + ), + other => Err(AudioError::Stream(format!( + "unsupported playback sample format: {other:?}" + ))), + }?; + stream + .play() + .map_err(|e| AudioError::Stream(e.to_string()))?; + let name = device.name().unwrap_or_else(|_| "".to_string()); + Ok((stream, name, sample_format, channels, device_rate)) + }; + + // Build + queue, then wait for real callbacks before reporting readiness (W1). + let (stream, dev_name, sample_format, channels, device_rate) = match setup() { + Ok(v) => match wait_for_stream_start(&callbacks, &err_code, &running) { + Ok(()) => { + let _ = ready.send(Ok(())); + v + } + Err(e) => { + // Drop before reporting so a wedged Stream::drop withholds the Err + // and lets finish_start's backstop detach instead of hanging (B2). + drop(v.0); + let _ = ready.send(Err(e)); + return; + } + }, + Err(e) => { + let _ = ready.send(Err(e)); + return; + } + }; + crate::log_msg(&format!( + "cpal playback started: device='{dev_name}' format={sample_format:?} channels={channels} device_rate={device_rate} Hz <- {SAMPLE_RATE} Hz" + )); + + let logger = spawn_health_logger( + running.clone(), + ring_fill.clone(), + underrun.clone(), + dropped.clone(), + max_cb.clone(), + err_code.clone(), + ); + + // Feed the ring from the network mixer until `stop()` flips `running` or the + // sender disconnects (call ended). Clock-paced production keeps the ring near + // target, so the drop path below should never fire in steady state. + drain_loop(&rx, &running, |frame| { + if ring_fill.load(Ordering::Relaxed) + frame.len() > RING_CAPACITY { + dropped.fetch_add(1, Ordering::Relaxed); + return; + } + // Reserve occupancy BEFORE publishing samples, and publish the whole frame + // in one `push_slice` (review W3). Per-sample pushes let the RT consumer + // observe a half-written stereo pair (L without R) → channel tear, and a + // pop that raced the post-loop `fetch_add` could drive `ring_fill` below + // zero and wrap it to usize::MAX, wedging the mixer's pacing. Reserving + // first means the consumer can never pop a sample that isn't yet counted. + ring_fill.fetch_add(frame.len(), Ordering::Relaxed); + let pushed = producer.push_slice(&frame); + if pushed != frame.len() { + // The capacity check above should make this unreachable (the consumer + // only drains), but stay exact if it ever isn't. + ring_fill.fetch_sub(frame.len() - pushed, Ordering::Relaxed); + dropped.fetch_add(1, Ordering::Relaxed); + } + }); + + // We're shutting down (either stop() or disconnect). Ensure the logger sees it + // even on the disconnect path, then drop the stream. + running.store(false, Ordering::Relaxed); + let _ = logger.join(); + drop(stream); +} + +#[allow(clippy::too_many_arguments)] +fn build_output( + device: &Device, + config: &StreamConfig, + mut consumer: C, + ring_fill: Arc, + underrun: Arc, + max_cb: Arc, + callbacks: Arc, + err_code: Arc, +) -> Result +where + T: SizedSample + FromSample + Send + 'static, + C: Consumer + Send + 'static, +{ + // RT-safe error callback: atomic store only, no alloc/log (review W2). + let err_fn = move |e: cpal::StreamError| err_code.store(stream_err_code(&e), Ordering::Relaxed); + let device_rate = config.sample_rate.0; + let device_channels = config.channels as usize; + if device_rate == SAMPLE_RATE && device_channels == PLAYBACK_CHANNELS { + device + .build_output_stream::( + config, + move |data: &mut [T], _| { + callbacks.fetch_add(1, Ordering::Relaxed); + // Record demand in INTERNAL 48 kHz-stereo samples (not raw + // device samples) so the health logger's prefill-target + // comparison is apples-to-apples for any rate/layout (W4). + max_cb.fetch_max( + internal_demand(data.len(), device_channels, device_rate), + Ordering::Relaxed, + ); + let (popped, starved) = fill_output(&mut consumer, data); + if starved > 0 { + underrun.fetch_add(starved, Ordering::Relaxed); + } + if popped > 0 { + // Decrement the exact occupancy by what we actually pulled + // (underruns removed nothing) so the mixer paces against the + // true ring depth. + ring_fill.fetch_sub(popped, Ordering::Relaxed); + } + }, + err_fn, + None, + ) + .map_err(|e| AudioError::Stream(e.to_string())) + } else { + let mut resampler = StereoPullResampler::new(SAMPLE_RATE, device_rate); + device + .build_output_stream::( + config, + move |data: &mut [T], _| { + callbacks.fetch_add(1, Ordering::Relaxed); + max_cb.fetch_max( + internal_demand(data.len(), device_channels, device_rate), + Ordering::Relaxed, + ); + let (popped, starved) = + fill_output_remap(&mut consumer, data, device_channels, &mut resampler); + if starved > 0 { + underrun.fetch_add(starved, Ordering::Relaxed); + } + if popped > 0 { + // Decrement the exact occupancy by what we actually pulled + // (underruns removed nothing) so the mixer paces against the + // true ring depth. + ring_fill.fetch_sub(popped, Ordering::Relaxed); + } + }, + err_fn, + None, + ) + .map_err(|e| AudioError::Stream(e.to_string())) + } +} + +/// Convert an output callback's raw device-sample length into the equivalent +/// internal 48 kHz-stereo sample demand, so the prefill-target comparison stays +/// meaningful regardless of the device's rate/channel layout (W4 diagnostic fix). +/// Integer-only and allocation-free, so it is safe on the RT callback thread. +#[inline] +fn internal_demand(device_len: usize, device_channels: usize, device_rate: u32) -> usize { + let device_frames = device_len / device_channels.max(1); + let need_frames = + (device_frames as u64 * SAMPLE_RATE as u64).div_ceil(device_rate.max(1) as u64) as usize; + need_frames * PLAYBACK_CHANNELS +} + +/// Drain the ring into the device buffer, substituting silence on underrun. +/// Returns `(samples_popped, samples_starved)`. RT-safe (wait-free `try_pop`). +fn fill_output(consumer: &mut C, out: &mut [T]) -> (usize, u64) +where + T: Sample + FromSample, + C: Consumer, +{ + let mut popped = 0usize; + let mut starved = 0u64; + for slot in out.iter_mut() { + match consumer.try_pop() { + Some(v) => { + *slot = T::from_sample(v); + popped += 1; + } + None => { + *slot = T::from_sample(0i16); + starved += 1; + } + } + } + (popped, starved) +} + +/// Resample/remap internal 48 kHz stereo ring samples into the device buffer. +/// Returns `(internal_samples_popped, device_samples_starved)`. RT-safe. +fn fill_output_remap( + consumer: &mut C, + out: &mut [T], + device_channels: usize, + resampler: &mut StereoPullResampler, +) -> (usize, u64) +where + T: Sample + FromSample, + C: Consumer, +{ + let mut popped = 0usize; + let mut starved = 0u64; + for frame in out.chunks_mut(device_channels) { + match resampler.next(|| { + let l = match consumer.try_pop() { + Some(v) => { + popped += 1; + v + } + None => return None, + }; + let r = match consumer.try_pop() { + Some(v) => { + popped += 1; + v + } + None => return None, + }; + Some((i16_to_f32(l), i16_to_f32(r))) + }) { + Some((l, r)) => { + if device_channels == 1 { + frame[0] = T::from_sample(f32_to_i16((l + r) * 0.5)); + } else { + frame[0] = T::from_sample(f32_to_i16(l)); + frame[1] = T::from_sample(f32_to_i16(r)); + for slot in &mut frame[2..] { + *slot = T::from_sample(0i16); + } + } + } + None => { + for slot in frame { + *slot = T::from_sample(0i16); + } + starved += device_channels as u64; + } + } + } + (popped, starved) +} + +/// Once-per-second playout-health line (mirrors the PipeWire backend). Quiet +/// unless a second actually glitched, or `PEERSPEAK_AUDIO_VERBOSE` is set. +fn spawn_health_logger( + running: Arc, + ring_fill: Arc, + underrun: Arc, + dropped: Arc, + max_cb: Arc, + err_code: Arc, +) -> JoinHandle<()> { + let verbose = std::env::var_os("PEERSPEAK_AUDIO_VERBOSE").is_some(); + thread::spawn(move || { + let (mut last_u, mut last_d) = (0u64, 0u64); + let mut reported_cb = 0usize; + let mut last_err = STREAM_ERR_NONE; + while running.load(Ordering::Relaxed) { + thread::sleep(Duration::from_secs(1)); + let u = underrun.load(Ordering::Relaxed); + let d = dropped.load(Ordering::Relaxed); + let fill = ring_fill.load(Ordering::Relaxed); + let (du, dd) = (u - last_u, d - last_d); + last_u = u; + last_d = d; + if verbose || du > 0 || dd > 0 { + crate::log_msg(&format!( + "playout-health: fill={fill} samples (~{}ms) | underrun +{du} samples/s (total {u}) | dropped +{dd} frames/s (total {d})", + fill / (48 * PLAYBACK_CHANNELS), + )); + } + // Surface a stream error the RT callback flagged (it can't log itself). + let ec = err_code.load(Ordering::Relaxed); + if ec != STREAM_ERR_NONE && ec != last_err { + crate::log_msg(&format!("cpal playback stream error: {}", stream_err_text(ec))); + last_err = ec; + } + // Report the device's per-cycle demand (in internal 48 kHz-stereo + // samples) the first time it's seen, and on any new high. If a callback + // demands more than the prefill target, the ring can't satisfy it and + // underruns every cycle — the W2 signature; warn so a real-host log + // shows whether it's biting. + let cb = max_cb.load(Ordering::Relaxed); + if cb > reported_cb { + reported_cb = cb; + let ms = cb / (48 * PLAYBACK_CHANNELS); + if cb > PLAYBACK_TARGET_SAMPLES { + crate::log_msg(&format!( + "cpal output callback demands up to {cb} internal samples/cycle (~{ms}ms) EXCEEDS prefill target {PLAYBACK_TARGET_SAMPLES} — expect periodic underruns; needs a larger target or a fixed buffer size (review W2)", + )); + } else if verbose { + crate::log_msg(&format!( + "cpal output callback demands up to {cb} internal samples/cycle (~{ms}ms), target {PLAYBACK_TARGET_SAMPLES}", + )); + } + } + } + }) +} + +/// Pump frames from `rx` to `on_frame` until `running` goes false or the sender +/// disconnects. The timed receive re-checks `running` at least every +/// [`WORKER_POLL`], so `stop()` can join the worker promptly instead of hanging +/// on a parked blocking `recv()` (same A7 fix as the PipeWire backend). Pure +/// w.r.t. its inputs, so it's unit-testable. +fn drain_loop(rx: &Receiver>, running: &AtomicBool, mut on_frame: impl FnMut(Vec)) { + while running.load(Ordering::Relaxed) { + match rx.recv_timeout(WORKER_POLL) { + Ok(frame) => on_frame(frame), + Err(RecvTimeoutError::Timeout) => continue, + Err(RecvTimeoutError::Disconnected) => return, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::mpsc; + + #[test] + fn downmix_averages_channels() { + assert_eq!(downmix_to_mono::(&[100, 100]), 100); + assert_eq!(downmix_to_mono::(&[100, -100]), 0); + assert_eq!(downmix_to_mono::(&[50]), 50); + assert_eq!(downmix_to_mono::(&[]), 0); + // 4-channel average rounds toward zero (integer division). + assert_eq!(downmix_to_mono::(&[10, 20, 30, 41]), 25); + } + + #[test] + fn frame_accumulator_emits_full_frames() { + let mut acc = FrameAccumulator::new(3); + assert_eq!(acc.push(1), None); + assert_eq!(acc.push(2), None); + assert_eq!(acc.push(3), Some(vec![1, 2, 3])); + // Resets for the next frame. + assert_eq!(acc.push(4), None); + assert_eq!(acc.push(5), None); + assert_eq!(acc.push(6), Some(vec![4, 5, 6])); + } + + #[test] + fn fill_output_pops_then_substitutes_silence() { + let rb = HeapRb::::new(8); + let (mut prod, mut cons) = rb.split(); + for v in [1, 2, 3] { + prod.try_push(v).unwrap(); + } + let mut out = [0i16; 5]; + let (popped, starved) = fill_output(&mut cons, &mut out); + assert_eq!(popped, 3); + assert_eq!(starved, 2); + assert_eq!(out, [1, 2, 3, 0, 0]); + } + + #[test] + fn fill_output_remap_downmixes_to_mono() { + let rb = HeapRb::::new(8); + let (mut prod, mut cons) = rb.split(); + for v in [100, 300, 500, -100, 7, 9] { + prod.try_push(v).unwrap(); + } + let mut out = [0i16; 2]; + let mut resampler = StereoPullResampler::new(SAMPLE_RATE, SAMPLE_RATE); + let (popped, starved) = fill_output_remap(&mut cons, &mut out, 1, &mut resampler); + assert_eq!(popped, 6); + assert_eq!(starved, 0); + assert_eq!(out, [200, 200]); + } + + #[test] + fn fill_output_remap_silences_underrun() { + let rb = HeapRb::::new(8); + let (_prod, mut cons) = rb.split(); + let mut out = [11i16; 4]; + let mut resampler = StereoPullResampler::new(SAMPLE_RATE, SAMPLE_RATE); + let (popped, starved) = fill_output_remap(&mut cons, &mut out, 2, &mut resampler); + assert_eq!(popped, 0); + assert_eq!(starved, out.len() as u64); + assert_eq!(out, [0, 0, 0, 0]); + } + + #[test] + fn fill_output_remap_copies_stereo_at_matching_rate() { + let rb = HeapRb::::new(8); + let (mut prod, mut cons) = rb.split(); + for v in [1, -1, 2, -2, 3, -3] { + prod.try_push(v).unwrap(); + } + let mut out = [0i16; 4]; + let mut resampler = StereoPullResampler::new(SAMPLE_RATE, SAMPLE_RATE); + let (popped, starved) = fill_output_remap(&mut cons, &mut out, 2, &mut resampler); + assert_eq!(popped, 6); + assert_eq!(starved, 0); + assert_eq!(out, [1, -1, 2, -2]); + } + + #[test] + fn drain_loop_exits_when_running_flips_even_with_sender_alive() { + let (tx, rx) = mpsc::channel::>(); + let running = Arc::new(AtomicBool::new(true)); + let r2 = running.clone(); + let h = thread::spawn(move || drain_loop(&rx, &r2, |_| {})); + thread::sleep(Duration::from_millis(50)); + running.store(false, Ordering::Relaxed); + thread::sleep(WORKER_POLL + Duration::from_millis(150)); + assert!( + h.is_finished(), + "drain_loop must exit after running=false even while the sender is alive" + ); + drop(tx); + h.join().unwrap(); + } + + #[test] + fn drain_loop_returns_on_disconnect() { + let (tx, rx) = mpsc::channel::>(); + let running = Arc::new(AtomicBool::new(true)); + drop(tx); + drain_loop(&rx, &running, |_| panic!("no frame should arrive")); + } + + #[test] + fn bounded_rate_prefers_48k_when_in_window() { + // A device span that contains 48 kHz resolves exactly. + assert_eq!(bounded_rate(44_100, 96_000), Some(SAMPLE_RATE)); + assert_eq!( + bounded_rate(MIN_DEVICE_RATE, MAX_DEVICE_RATE), + Some(SAMPLE_RATE) + ); + } + + #[test] + fn bounded_rate_clamps_to_nearest_in_window_bound() { + // Entirely below 48 kHz → the top bound (closest reachable to 48 kHz). + assert_eq!(bounded_rate(8_000, 16_000), Some(16_000)); + // Entirely above 48 kHz → the bottom bound. + assert_eq!(bounded_rate(88_200, 192_000), Some(88_200)); + } + + #[test] + fn bounded_rate_rejects_spans_outside_the_window() { + assert_eq!(bounded_rate(1_000, 4_000), None); // below the floor + assert_eq!(bounded_rate(400_000, 500_000), None); // above the ceiling + } + + #[test] + fn bounded_rate_intersects_window_edges() { + // Overlaps only the floor: [4k, 8k] ∩ [8k, 384k] = {8k}. + assert_eq!(bounded_rate(4_000, MIN_DEVICE_RATE), Some(MIN_DEVICE_RATE)); + // Overlaps only the ceiling. + assert_eq!( + bounded_rate(MAX_DEVICE_RATE, 500_000), + Some(MAX_DEVICE_RATE) + ); + } + + #[test] + fn format_supported_matches_the_stream_builders() { + // Exactly the three the build_input/build_output match arms accept. + for f in [SampleFormat::F32, SampleFormat::I16, SampleFormat::U16] { + assert!(format_supported(f), "{f:?} should be drivable"); + } + // Everything else cpal can expose must be filtered out before ranking, or a + // start could pick it and then hit the `unsupported sample format` arm (P3). + for f in [ + SampleFormat::I8, + SampleFormat::U8, + SampleFormat::I32, + SampleFormat::U32, + SampleFormat::I64, + SampleFormat::U64, + SampleFormat::F64, + ] { + assert!(!format_supported(f), "{f:?} must not be reported drivable"); + } + } + + #[test] + fn ensure_idle_allows_an_idle_slot() { + let mut s = SlotState::Idle; + assert!(ensure_idle(&mut s, "capture").is_ok()); + assert!(matches!(s, SlotState::Idle)); + } + + #[test] + fn ensure_idle_rejects_a_live_wedged_orphan_then_clears_when_it_exits() { + let exited = Arc::new(AtomicBool::new(false)); + let mut s = SlotState::Wedged { + exited: exited.clone(), + }; + // Orphan still alive → reject, tombstone preserved. + assert!(ensure_idle(&mut s, "playback").is_err()); + assert!(matches!(s, SlotState::Wedged { .. })); + // Orphan's ExitGuard fired → the next start clears the tombstone and proceeds. + exited.store(true, Ordering::Relaxed); + assert!(ensure_idle(&mut s, "playback").is_ok()); + assert!(matches!(s, SlotState::Idle)); + } + + #[test] + fn drain_loop_delivers_frames() { + let (tx, rx) = mpsc::channel::>(); + let running = Arc::new(AtomicBool::new(true)); + let r2 = running.clone(); + let got = Arc::new(Mutex::new(Vec::new())); + let g2 = got.clone(); + let h = thread::spawn(move || drain_loop(&rx, &r2, |f| g2.lock().unwrap().push(f))); + tx.send(vec![1, 2, 3]).unwrap(); + tx.send(vec![4, 5]).unwrap(); + thread::sleep(Duration::from_millis(50)); + running.store(false, Ordering::Relaxed); + drop(tx); + h.join().unwrap(); + assert_eq!(*got.lock().unwrap(), vec![vec![1, 2, 3], vec![4, 5]]); + } +} diff --git a/src/audio/mod.rs b/src/audio/mod.rs index 4a83742..1d068e0 100644 --- a/src/audio/mod.rs +++ b/src/audio/mod.rs @@ -56,12 +56,60 @@ pub trait AudioBackend: Send + Sync { fn stop(&self) -> Result<(), AudioError>; } -pub mod echo_cancel; pub mod eq; pub mod gate; pub mod limiter; pub mod multitrack; pub mod pan; +// Linear resamplers used by the Windows/cpal backend (W4). Platform-neutral and +// pure, so it builds (and its tests run) everywhere even though only the cpal +// backend wires it in. +pub mod resample; +#[cfg(target_os = "linux")] +pub mod echo_cancel; +#[cfg(target_os = "linux")] pub mod pipewire_impl; +#[cfg(windows)] +pub mod cpal_impl; +#[cfg(target_os = "linux")] pub mod pw_cli; pub mod recorder; + +/// A selectable audio device for the input/output pickers. `name` is the stable +/// identifier the backend uses to request the device (`target_node`); +/// `description` is the human-facing label shown in the UI. The two may be equal +/// (cpal/WASAPI) or differ (PipeWire node name vs. description). +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct AudioDevice { + pub name: String, + pub description: String, + pub is_input: bool, +} + +impl std::fmt::Display for AudioDevice { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "{}", self.description) + } +} + +// Enumerate audio input/output devices for the pickers (sorted by description), +// returning the same `AudioDevice` shape regardless of platform: PipeWire +// (`pw-cli`) on Linux, cpal/WASAPI on Windows. +#[cfg(target_os = "linux")] +pub use pw_cli::enumerate_audio_devices; +#[cfg(windows)] +pub use cpal_impl::enumerate_audio_devices; + +/// The audio backend implementation for the current platform. +/// +/// The whole app constructs and threads this alias (via +/// `PlatformAudioBackend::new()`) rather than any concrete backend type, so +/// platform selection lives entirely here. Both implementations satisfy the +/// [`AudioBackend`] trait, which is the only interface the core talks to. +/// +/// - Linux → PipeWire ([`pipewire_impl::PipeWireBackend`]). +/// - Windows → cpal/WASAPI ([`cpal_impl::CpalBackend`]). +#[cfg(target_os = "linux")] +pub type PlatformAudioBackend = pipewire_impl::PipeWireBackend; +#[cfg(windows)] +pub type PlatformAudioBackend = cpal_impl::CpalBackend; diff --git a/src/audio/pw_cli.rs b/src/audio/pw_cli.rs index 51d3ed6..f1e4e61 100644 --- a/src/audio/pw_cli.rs +++ b/src/audio/pw_cli.rs @@ -1,18 +1,6 @@ +use super::AudioDevice; use std::process::Command; -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct AudioDevice { - pub name: String, - pub description: String, - pub is_input: bool, -} - -impl std::fmt::Display for AudioDevice { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{}", self.description) - } -} - pub fn enumerate_audio_devices() -> Vec { let output = Command::new("pw-cli") .arg("list-objects") diff --git a/src/audio/resample.rs b/src/audio/resample.rs new file mode 100644 index 0000000..4f5723f --- /dev/null +++ b/src/audio/resample.rs @@ -0,0 +1,307 @@ +//! Dep-free linear-interpolation resamplers for the Windows/cpal backend (W4). +//! +//! The pipeline runs internally at 48 kHz (Opus + the 20 ms frame), but a WASAPI +//! endpoint may run at a different rate (commonly 44.1 kHz) and/or a non-stereo +//! channel layout. These convert at the device boundary so such a device plays and +//! captures instead of hard-erroring (the W4 limitation in the Windows port). +//! +//! ## Where each is used +//! - [`PushResampler`] (single channel) converts **capture** from the device rate +//! to 48 kHz on the capture drain thread — off the RT callback. +//! - [`StereoPullResampler`] converts **playback** from the internal 48 kHz stereo +//! bus to the device rate inside the output RT callback, pulling internal frames +//! from the ring on demand. It allocates nothing in `next`, so it is RT-safe. +//! +//! ## Quality +//! This is plain linear interpolation with no anti-aliasing filter: correct, +//! allocation-free, and adequate for speech, but it adds some aliasing when +//! downsampling. The seam is intentionally tiny so a higher-quality polyphase/FIR +//! resampler (e.g. the `rubato` crate, pending a supply-chain decision) can later +//! replace the internals without touching the cpal backend. The matching-rate / +//! matching-layout path in the backend bypasses these entirely and stays bit-exact. + +/// Linear interpolation between `a` and `b` at fractional position `frac` in `[0, 1)`. +#[inline] +fn lerp(a: f32, b: f32, frac: f32) -> f32 { + a + (b - a) * frac +} + +/// Stateful single-channel **push** resampler: feed input samples at `in_rate`, +/// receive output samples at `out_rate` through an `emit` callback. It carries the +/// fractional read position and the previous input sample across calls, so feeding +/// the stream block-by-block joins seamlessly. Neither [`push`](Self::push) nor +/// [`process`](Self::process) allocates. +pub struct PushResampler { + /// Input samples consumed per output sample (`in_rate / out_rate`). + step: f64, + /// Position of the next output sample, in input-sample units, measured from the + /// index of `prev` (the most recent input). Always advanced to stay `< 1.0` + /// after each input is consumed. + next: f64, + /// The previous input sample (left edge of the current interpolation segment). + prev: f32, + /// Whether any input has been seen yet (anchors the first output at input[0]). + started: bool, +} + +impl PushResampler { + /// Build a resampler from `in_rate` to `out_rate` (both in Hz). Rates are + /// clamped to `>= 1` so `step` is always finite and non-zero: a zero `step` + /// would make [`push`](Self::push)'s `while self.next < 1.0` loop forever. The + /// cpal backend's `resolve()` also rejects such rates up front, so this is + /// belt-and-suspenders against a future caller (review W7). + pub fn new(in_rate: u32, out_rate: u32) -> Self { + Self { + step: in_rate.max(1) as f64 / out_rate.max(1) as f64, + next: 0.0, + prev: 0.0, + started: false, + } + } + + /// Feed one input sample; `emit` is called for each output sample produced + /// (zero or more, depending on the rate ratio). + pub fn push(&mut self, cur: f32, mut emit: impl FnMut(f32)) { + if !self.started { + // First sample: just establish the left edge. Linear interpolation + // needs the next input as the right edge, so the first output is + // produced on the next push. This gives exact alignment + // (`output[k] == input[k]` at equal rates) with one input-sample of + // latency — negligible (~20 µs at 48 kHz). + self.started = true; + self.prev = cur; + self.next = 0.0; + return; + } + // `prev` sits at position 0 of this segment and `cur` at position 1; emit + // every output whose position falls in [0, 1). + while self.next < 1.0 { + emit(lerp(self.prev, cur, self.next as f32)); + self.next += self.step; + } + self.next -= 1.0; + self.prev = cur; + } + + /// Convenience for tests / batch callers: push a whole slice. + pub fn process(&mut self, input: &[f32], mut emit: impl FnMut(f32)) { + for &s in input { + self.push(s, &mut emit); + } + } +} + +/// Stateful stereo **pull** resampler: produce output frames at `out_rate` by +/// pulling input frames at `in_rate` from a closure on demand. Call +/// [`next`](Self::next) once per output frame; it pulls as many input frames as the +/// ratio requires and returns the interpolated `(left, right)`, or `None` when the +/// puller runs dry (an underrun). Allocates nothing, so it is safe in an RT output +/// callback. +pub struct StereoPullResampler { + /// Input frames consumed per output frame (`in_rate / out_rate`). + step: f64, + /// Position of the next output frame within `[prev, cur)`, in `[0, 1)`. + frac: f64, + /// Left edge of the current interpolation segment. + prev: (f32, f32), + /// Right edge of the current interpolation segment. + cur: (f32, f32), + /// Whether `prev`/`cur` have been primed from the puller yet. + primed: bool, +} + +impl StereoPullResampler { + /// Build a resampler from `in_rate` to `out_rate` (both in Hz). Rates are + /// clamped to `>= 1` so `step` is finite and non-zero — otherwise + /// [`next`](Self::next)'s `while self.frac >= 1.0` could spin (review W7). + pub fn new(in_rate: u32, out_rate: u32) -> Self { + Self { + step: in_rate.max(1) as f64 / out_rate.max(1) as f64, + frac: 0.0, + prev: (0.0, 0.0), + cur: (0.0, 0.0), + primed: false, + } + } + + /// Produce the next output frame, pulling input frames via `pull` as needed. + /// Returns `None` if `pull` returns `None` before the frame can be formed + /// (underrun); the caller should substitute silence for that frame. + pub fn next(&mut self, mut pull: impl FnMut() -> Option<(f32, f32)>) -> Option<(f32, f32)> { + if !self.primed { + // Prime both edges from two pulls so the first output frame aligns + // exactly with the first input frame (`out[0] == in[0]` at equal + // rates). Needs two frames available to start, which the prefilled + // playback ring always has. + self.prev = pull()?; + self.cur = pull()?; + self.primed = true; + self.frac = 0.0; + } + // Advance the segment until the read position lands inside [prev, cur). + while self.frac >= 1.0 { + self.prev = self.cur; + self.cur = pull()?; + self.frac -= 1.0; + } + let f = self.frac as f32; + let out = (lerp(self.prev.0, self.cur.0, f), lerp(self.prev.1, self.cur.1, f)); + self.frac += self.step; + Some(out) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Equal rates align exactly: `output[k] == input[k]`. The final input lands on + /// the next push (one-sample streaming latency), so we get `n - 1` outputs. + #[test] + fn push_identity_when_rates_match() { + let mut r = PushResampler::new(48_000, 48_000); + let input = [0.0, 0.1, 0.2, 0.3, 0.4]; + let mut out = Vec::new(); + r.process(&input, |s| out.push(s)); + assert_eq!(out.len(), input.len() - 1); + for (a, b) in out.iter().zip(input.iter()) { + assert!((a - b).abs() < 1e-6, "{a} vs {b}"); + } + } + + /// Upsampling 2x roughly doubles the output count and the midpoints interpolate. + #[test] + fn push_upsample_2x_interpolates_midpoints() { + let mut r = PushResampler::new(24_000, 48_000); // step = 0.5 + let input = [0.0, 1.0, 2.0, 3.0]; + let mut out = Vec::new(); + r.process(&input, |s| out.push(s)); + // (n - 1) segments at 2 outputs each = 6. + assert_eq!(out.len(), 6, "out {out:?}"); + // A half-step between 1.0 and 2.0 must appear near 1.5. + assert!( + out.iter().any(|&s| (s - 1.5).abs() < 1e-3), + "expected a ~1.5 midpoint in {out:?}" + ); + } + + /// Downsampling drops the rate: fewer outputs than inputs, monotonic ramp preserved. + #[test] + fn push_downsample_reduces_count() { + let mut r = PushResampler::new(48_000, 44_100); // step ~1.088 + let input: Vec = (0..441).map(|i| i as f32).collect(); + let mut out = Vec::new(); + r.process(&input, |s| out.push(s)); + // 441 in @ 48k -> ~405 out @ 44.1k. + assert!( + (390..=410).contains(&out.len()), + "expected ~405 outputs, got {}", + out.len() + ); + // Output stays within the input's value range and is non-decreasing. + for w in out.windows(2) { + assert!(w[1] >= w[0] - 1e-3, "ramp should not reverse: {w:?}"); + } + assert!(*out.last().unwrap() <= 440.0 + 1e-3); + } + + /// Pull resampler at equal rates returns each input frame in order, aligned. + /// Two-pull priming uses one frame of lookahead, so `n` inputs yield `n - 1` + /// outputs (the last frame emits once a successor arrives). + #[test] + fn pull_identity_when_rates_match() { + let mut r = StereoPullResampler::new(48_000, 48_000); + let frames = [(0.0, 9.0), (1.0, 8.0), (2.0, 7.0), (3.0, 6.0)]; + let mut idx = 0; + let mut out = Vec::new(); + while let Some(f) = r.next(|| { + let v = frames.get(idx).copied(); + idx += 1; + v + }) { + out.push(f); + } + assert_eq!(out.len(), frames.len() - 1, "out {out:?}"); + for (got, want) in out.iter().zip(frames.iter()) { + assert!((got.0 - want.0).abs() < 1e-6 && (got.1 - want.1).abs() < 1e-6); + } + } + + /// Pull resampler reports underrun (`None`) once the source is exhausted. + #[test] + fn pull_returns_none_on_underrun() { + let mut r = StereoPullResampler::new(48_000, 44_100); // step ~1.088 -> pulls >1 per out + let frames = [(0.0, 0.0), (1.0, -1.0)]; + let mut idx = 0; + let mut pull = || { + let v = frames.get(idx).copied(); + idx += 1; + v + }; + // First frame primes + emits; subsequent calls eventually exhaust the source. + let mut produced = 0; + let mut hit_none = false; + for _ in 0..10 { + if r.next(&mut pull).is_some() { + produced += 1; + } else { + hit_none = true; + break; + } + } + assert!(produced >= 1, "should produce at least the primed frame"); + assert!(hit_none, "should report underrun once the puller is dry"); + } + + /// Downsampling via pull consumes more input frames than it emits output frames. + #[test] + fn pull_downsample_consumes_more_than_it_emits() { + let mut r = StereoPullResampler::new(48_000, 24_000); // step = 2.0 + let input: Vec<(f32, f32)> = (0..100).map(|i| (i as f32, -(i as f32))).collect(); + let mut idx = 0; + let mut emitted = 0; + for _ in 0..40 { + let f = r.next(|| { + let v = input.get(idx).copied(); + idx += 1; + v + }); + if f.is_some() { + emitted += 1; + } else { + break; + } + } + // At step 2.0 we consume ~2 input frames per output frame. + assert!(idx > emitted, "consumed {idx} input, emitted {emitted} output"); + } + + /// A zero rate must not produce a zero `step` (which would spin `push`'s inner + /// `while self.next < 1.0` forever). Clamping makes the call terminate (W7). + #[test] + fn push_zero_rate_does_not_spin() { + let mut r = PushResampler::new(0, 48_000); + let mut count = 0usize; + // Feed two samples; with a clamped non-zero step this returns promptly. + r.push(0.0, |_| count += 1); + r.push(1.0, |_| count += 1); + // Reaching here at all is the assertion (no hang); some output is produced. + assert!(count >= 1); + } + + /// A zero output rate must not make the pull resampler's segment-advance loop + /// spin. Clamping keeps `step` finite so `next` terminates (W7). + #[test] + fn pull_zero_out_rate_does_not_spin() { + let mut r = StereoPullResampler::new(48_000, 0); + let frames = [(0.0, 0.0), (1.0, 1.0), (2.0, 2.0)]; + let mut idx = 0; + let got = r.next(|| { + let v = frames.get(idx).copied(); + idx += 1; + v + }); + // Terminates and yields the primed frame instead of hanging. + assert!(got.is_some()); + } +} diff --git a/src/bin/audio_probe.rs b/src/bin/audio_probe.rs index e655622..f7429a3 100644 --- a/src/bin/audio_probe.rs +++ b/src/bin/audio_probe.rs @@ -1,11 +1,11 @@ //! Audio playout diagnostic probe. //! -//! Drives a phase-continuous sine tone through the *real* PipeWire playback path -//! (`PipeWireBackend::start_playback`), using the *same* fill-paced production -//! the production mixer uses (`core/mod.rs`): generate a frame only while the +//! Drives a phase-continuous sine tone through the *real* playback path +//! (PipeWire on Linux, cpal/WASAPI on Windows), using the *same* fill-paced +//! production the production mixer uses (`core/mod.rs`): generate a frame only while the //! playback ring is below `PLAYBACK_TARGET_SAMPLES`, so production tracks the -//! PipeWire hardware clock. No network, no microphone — this isolates the local -//! output path so we can confirm the clock-paced playout is glitch-free. +//! hardware clock. No network, no microphone — this isolates the local output +//! path so we can confirm the clock-paced playout is glitch-free. //! //! Use your ears on the tone (any click/pop is a glitch) together with the //! `playout-health:` lines tailed to stdout: @@ -17,105 +17,235 @@ //! //! Run: cargo run --bin audio_probe -- [freq_hz] [seconds] [target_node] //! e.g. cargo run --release --bin audio_probe -- 440 30 +//! +//! This probe exercises the platform playback backend directly: PipeWire on Linux +//! and cpal/WASAPI on Windows. Other targets use a stub that explains the limitation. -use std::io::{BufRead, BufReader, Seek, SeekFrom}; -use std::sync::Arc; -use std::sync::atomic::AtomicUsize; -use std::sync::mpsc; -use std::time::Duration; - -use peerspeak::audio::AudioBackend; -use peerspeak::audio::pipewire_impl::PipeWireBackend; -use peerspeak::core::jitter::FRAME_SAMPLES; // 960 mono frames = 20ms @ 48kHz - -const SAMPLE_RATE: f32 = 48_000.0; - -#[tokio::main] -async fn main() { - let mut args = std::env::args().skip(1); - let freq: f32 = args.next().and_then(|s| s.parse().ok()).unwrap_or(440.0); - let secs: u64 = args.next().and_then(|s| s.parse().ok()).unwrap_or(30); - let target_node: Option = args.next(); - - // The playout-health logger is quiet in normal operation (it only logs - // glitches); ask it for the full once-per-second heartbeat so the probe can - // show the steady-state numbers. - // SAFETY: set before any playback thread starts, so no concurrent env read. - unsafe { std::env::set_var("PEERSPEAK_AUDIO_VERBOSE", "1") }; - - println!("audio_probe: {freq} Hz tone for {secs}s through the real playback path."); - println!("Listen for clicks/pops; watch the playout-health lines below.\n"); - - // Tail the app log (where playout-health lines land) to stdout in the - // background so it's all in one terminal. - spawn_log_tailer(); - - let backend = PipeWireBackend::new(); - let (tx, rx) = mpsc::channel::>(); - let ring_fill = Arc::new(AtomicUsize::new(0)); - if let Err(e) = backend.start_playback(rx, target_node, ring_fill.clone()) { - eprintln!("failed to start playback: {e}"); - return; - } - - // Phase-continuous sine, generated one 20ms frame at a time, fill-paced - // exactly like the production mixer: only produce while the ring is below - // target, so production tracks the PipeWire hardware clock. - use std::sync::atomic::Ordering; - let deadline = tokio::time::Instant::now() + Duration::from_secs(secs); - let mut n: u64 = 0; // running sample index keeps phase continuous across frames - while tokio::time::Instant::now() < deadline { - if ring_fill.load(Ordering::Relaxed) >= peerspeak::audio::PLAYBACK_TARGET_SAMPLES { - tokio::time::sleep(Duration::from_millis(2)).await; - continue; - } - let mut frame = Vec::with_capacity(FRAME_SAMPLES * peerspeak::audio::PLAYBACK_CHANNELS); - for _ in 0..FRAME_SAMPLES { - let t = n as f32 / SAMPLE_RATE; - // 0.25 amplitude: clearly audible but not harsh. - let sample = (0.25 * i16::MAX as f32 * (2.0 * std::f32::consts::PI * freq * t).sin()) as i16; - // Stereo playback bus: duplicate the probe tone to L/R. - frame.push(sample); - frame.push(sample); - n += 1; - } - if tx.send(frame).is_err() { - eprintln!("playback channel closed early"); - break; - } - } - - // Let the ring drain, then stop. - tokio::time::sleep(Duration::from_millis(300)).await; - let _ = backend.stop(); - println!("\naudio_probe: done."); +#[cfg(target_os = "linux")] +fn main() { + unix_probe::run(); } -/// Open the app log, seek to the end, and echo new lines (the `playout-health:` -/// reports) to stdout once they appear. -fn spawn_log_tailer() { - let path = peerspeak::log_file_path(); - std::thread::spawn(move || { - // Wait for the file to exist (first log_msg creates it). - let file = loop { - if let Ok(f) = std::fs::File::open(&path) { - break f; +#[cfg(windows)] +fn main() { + win_probe::run(); +} + +#[cfg(not(any(target_os = "linux", windows)))] +fn main() { + eprintln!( + "audio_probe is only supported on Linux and Windows builds (it drives the platform playback backend directly)." + ); +} + +#[cfg(target_os = "linux")] +mod unix_probe { + use std::io::{BufRead, BufReader, Seek, SeekFrom}; + use std::sync::Arc; + use std::sync::atomic::AtomicUsize; + use std::sync::mpsc; + use std::time::Duration; + + use peerspeak::audio::AudioBackend; + use peerspeak::audio::pipewire_impl::PipeWireBackend; + use peerspeak::core::jitter::FRAME_SAMPLES; // 960 mono frames = 20ms @ 48kHz + + const SAMPLE_RATE: f32 = 48_000.0; + + #[tokio::main] + pub async fn run() { + let mut args = std::env::args().skip(1); + let freq: f32 = args.next().and_then(|s| s.parse().ok()).unwrap_or(440.0); + let secs: u64 = args.next().and_then(|s| s.parse().ok()).unwrap_or(30); + let target_node: Option = args.next(); + + // The playout-health logger is quiet in normal operation (it only logs + // glitches); ask it for the full once-per-second heartbeat so the probe can + // show the steady-state numbers. + // SAFETY: set before any playback thread starts, so no concurrent env read. + unsafe { std::env::set_var("PEERSPEAK_AUDIO_VERBOSE", "1") }; + + println!("audio_probe: {freq} Hz tone for {secs}s through the real playback path."); + println!("Listen for clicks/pops; watch the playout-health lines below.\n"); + + // Tail the app log (where playout-health lines land) to stdout in the + // background so it's all in one terminal. + spawn_log_tailer(); + + let backend = PipeWireBackend::new(); + let (tx, rx) = mpsc::channel::>(); + let ring_fill = Arc::new(AtomicUsize::new(0)); + if let Err(e) = backend.start_playback(rx, target_node, ring_fill.clone()) { + eprintln!("failed to start playback: {e}"); + return; + } + + // Phase-continuous sine, generated one 20ms frame at a time, fill-paced + // exactly like the production mixer: only produce while the ring is below + // target, so production tracks the PipeWire hardware clock. + use std::sync::atomic::Ordering; + let deadline = tokio::time::Instant::now() + Duration::from_secs(secs); + let mut n: u64 = 0; // running sample index keeps phase continuous across frames + while tokio::time::Instant::now() < deadline { + if ring_fill.load(Ordering::Relaxed) >= peerspeak::audio::PLAYBACK_TARGET_SAMPLES { + tokio::time::sleep(Duration::from_millis(2)).await; + continue; } - std::thread::sleep(Duration::from_millis(100)); - }; - let mut reader = BufReader::new(file); - let _ = reader.seek(SeekFrom::End(0)); - loop { - let mut line = String::new(); - match reader.read_line(&mut line) { - Ok(0) => std::thread::sleep(Duration::from_millis(150)), - Ok(_) => { - if line.contains("playout-health:") { - print!("{line}"); - } + let mut frame = Vec::with_capacity(FRAME_SAMPLES * peerspeak::audio::PLAYBACK_CHANNELS); + for _ in 0..FRAME_SAMPLES { + let t = n as f32 / SAMPLE_RATE; + // 0.25 amplitude: clearly audible but not harsh. + let sample = + (0.25 * i16::MAX as f32 * (2.0 * std::f32::consts::PI * freq * t).sin()) as i16; + // Stereo playback bus: duplicate the probe tone to L/R. + frame.push(sample); + frame.push(sample); + n += 1; + } + if tx.send(frame).is_err() { + eprintln!("playback channel closed early"); + break; + } + } + + // Let the ring drain, then stop. + tokio::time::sleep(Duration::from_millis(300)).await; + let _ = backend.stop(); + println!("\naudio_probe: done."); + } + + /// Open the app log, seek to the end, and echo new lines (the `playout-health:` + /// reports) to stdout once they appear. + fn spawn_log_tailer() { + let path = peerspeak::log_file_path(); + std::thread::spawn(move || { + // Wait for the file to exist (first log_msg creates it). + let file = loop { + if let Ok(f) = std::fs::File::open(&path) { + break f; } - Err(_) => std::thread::sleep(Duration::from_millis(150)), + std::thread::sleep(Duration::from_millis(100)); + }; + let mut reader = BufReader::new(file); + let _ = reader.seek(SeekFrom::End(0)); + loop { + let mut line = String::new(); + match reader.read_line(&mut line) { + Ok(0) => std::thread::sleep(Duration::from_millis(150)), + Ok(_) => { + if line.contains("playout-health:") { + print!("{line}"); + } + } + Err(_) => std::thread::sleep(Duration::from_millis(150)), + } + } + }); + } +} + +#[cfg(windows)] +mod win_probe { + use std::io::{BufRead, BufReader, Seek, SeekFrom}; + use std::sync::Arc; + use std::sync::atomic::AtomicUsize; + use std::sync::mpsc; + use std::time::Duration; + + use peerspeak::audio::AudioBackend; + use peerspeak::audio::cpal_impl::CpalBackend; + use peerspeak::core::jitter::FRAME_SAMPLES; // 960 mono frames = 20ms @ 48kHz + + const SAMPLE_RATE: f32 = 48_000.0; + + #[tokio::main] + pub async fn run() { + let mut args = std::env::args().skip(1); + let freq: f32 = args.next().and_then(|s| s.parse().ok()).unwrap_or(440.0); + let secs: u64 = args.next().and_then(|s| s.parse().ok()).unwrap_or(30); + let target_node: Option = args.next(); + + // The playout-health logger is quiet in normal operation (it only logs + // glitches); ask it for the full once-per-second heartbeat so the probe can + // show the steady-state numbers. + // SAFETY: set before any playback thread starts, so no concurrent env read. + unsafe { std::env::set_var("PEERSPEAK_AUDIO_VERBOSE", "1") }; + + println!("audio_probe: {freq} Hz tone for {secs}s through the real playback path."); + println!("Listen for clicks/pops; watch the playout-health lines below.\n"); + + // Tail the app log (where playout-health lines land) to stdout in the + // background so it's all in one terminal. + spawn_log_tailer(); + + let backend = CpalBackend::new(); + let (tx, rx) = mpsc::channel::>(); + let ring_fill = Arc::new(AtomicUsize::new(0)); + if let Err(e) = backend.start_playback(rx, target_node, ring_fill.clone()) { + eprintln!("failed to start playback: {e}"); + return; + } + + // Phase-continuous sine, generated one 20ms frame at a time, fill-paced + // exactly like the production mixer: only produce while the ring is below + // target, so production tracks the cpal/WASAPI hardware clock. + use std::sync::atomic::Ordering; + let deadline = tokio::time::Instant::now() + Duration::from_secs(secs); + let mut n: u64 = 0; // running sample index keeps phase continuous across frames + while tokio::time::Instant::now() < deadline { + if ring_fill.load(Ordering::Relaxed) >= peerspeak::audio::PLAYBACK_TARGET_SAMPLES { + tokio::time::sleep(Duration::from_millis(2)).await; + continue; + } + let mut frame = Vec::with_capacity(FRAME_SAMPLES * peerspeak::audio::PLAYBACK_CHANNELS); + for _ in 0..FRAME_SAMPLES { + let t = n as f32 / SAMPLE_RATE; + // 0.25 amplitude: clearly audible but not harsh. + let sample = + (0.25 * i16::MAX as f32 * (2.0 * std::f32::consts::PI * freq * t).sin()) as i16; + // Stereo playback bus: duplicate the probe tone to L/R. + frame.push(sample); + frame.push(sample); + n += 1; + } + if tx.send(frame).is_err() { + eprintln!("playback channel closed early"); + break; } } - }); + + // Let the ring drain, then stop. + tokio::time::sleep(Duration::from_millis(300)).await; + let _ = backend.stop(); + println!("\naudio_probe: done."); + } + + /// Open the app log, seek to the end, and echo new lines (the `playout-health:` + /// reports) to stdout once they appear. + fn spawn_log_tailer() { + let path = peerspeak::log_file_path(); + std::thread::spawn(move || { + // Wait for the file to exist (first log_msg creates it). + let file = loop { + if let Ok(f) = std::fs::File::open(&path) { + break f; + } + std::thread::sleep(Duration::from_millis(100)); + }; + let mut reader = BufReader::new(file); + let _ = reader.seek(SeekFrom::End(0)); + loop { + let mut line = String::new(); + match reader.read_line(&mut line) { + Ok(0) => std::thread::sleep(Duration::from_millis(150)), + Ok(_) => { + if line.contains("playout-health:") { + print!("{line}"); + } + } + Err(_) => std::thread::sleep(Duration::from_millis(150)), + } + } + }); + } } diff --git a/src/core/mod.rs b/src/core/mod.rs index 2a6b58d..afd8341 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -1,7 +1,7 @@ pub mod messages; pub mod jitter; -use crate::audio::{AudioBackend, pipewire_impl::PipeWireBackend}; +use crate::audio::{AudioBackend, PlatformAudioBackend}; use crate::audio::eq::{Eq, EqSettings}; use crate::codec::{AudioEncoder, opus_impl::OpusEncoder}; use crate::core::jitter::{JitterBuffer, FRAME_SAMPLES}; @@ -237,7 +237,7 @@ fn run_mic_monitor( /// Stops a standalone mic monitor if one is running. MUST NOT be called while a /// room session is active — `backend.stop()` would also tear down the call's /// capture/playback. Monitor and session are mutually exclusive by construction. -fn stop_mic_monitor(backend: &PipeWireBackend, monitor: Option) { +fn stop_mic_monitor(backend: &PlatformAudioBackend, monitor: Option) { if let Some(m) = monitor { let _ = backend.stop(); let _ = m.thread.join(); @@ -390,6 +390,7 @@ struct ActiveSession { grace_timers: GraceTimers, transport: Arc, /// Loaded PipeWire echo-cancel module (if enabled); unloads on drop. + #[cfg(target_os = "linux")] echo_cancel: Option, /// Our pixelpass screen-share host child while sharing (`kill_on_drop`, so it /// also dies if the session is dropped without an explicit stop). @@ -400,7 +401,7 @@ struct ActiveSession { } impl ActiveSession { - async fn shutdown(mut self, audio_backend: Arc) { + async fn shutdown(mut self, audio_backend: Arc) { crate::log_msg("ActiveSession::shutdown started"); // Tear down any screen-share children first so the host stops streaming // promptly (kill_on_drop is the backstop, but kill explicitly so viewers @@ -431,6 +432,7 @@ impl ActiveSession { // Unload the echo-cancel module now that the audio streams releasing its // virtual nodes have stopped. (Dropping the guard runs `pactl unload`.) + #[cfg(target_os = "linux")] drop(self.echo_cancel); crate::log_msg("Leaving room..."); @@ -731,7 +733,7 @@ async fn run_core_loop( let known_peers: Arc>>> = Arc::new(std::sync::Mutex::new(HashMap::new())); - let audio_backend = Arc::new(PipeWireBackend::new()); + let audio_backend = Arc::new(PlatformAudioBackend::new()); let is_muted = Arc::new(AtomicBool::new(false)); let is_deafened = Arc::new(AtomicBool::new(false)); @@ -1095,7 +1097,9 @@ async fn run_core_loop( // The guard unloads the module on drop — including the early-return // paths below, since it's a local until moved into the session. On // any failure, warn and fall back to the direct devices. + #[cfg(target_os = "linux")] let mut echo_cancel_guard = None; + #[cfg(target_os = "linux")] let (capture_target, playback_target) = if echo_cancellation { match crate::audio::echo_cancel::enable( input_device.as_deref(), @@ -1122,6 +1126,10 @@ async fn run_core_loop( } else { (input_device.clone(), output_device.clone()) }; + #[cfg(not(target_os = "linux"))] + let _ = echo_cancellation; + #[cfg(not(target_os = "linux"))] + let (capture_target, playback_target) = (input_device.clone(), output_device.clone()); if let Err(e) = audio_backend.start_capture(capture_tx, capture_target) { let _ = ui_tx.send(UiEvent::Error(format!("Failed to start capture: {}", e))).await; @@ -1632,6 +1640,7 @@ async fn run_core_loop( conn_event_task, grace_timers, transport: transport.clone(), + #[cfg(target_os = "linux")] echo_cancel: echo_cancel_guard, screenshare_host: None, screenshare_viewers: Vec::new(), @@ -1807,17 +1816,26 @@ async fn run_core_loop( } CoreCommand::SetNetworkMode(mode) => { - network_mode = mode; - // Rebuild the persistent stack to the new posture immediately if - // idle; if a call is active, defer to the next Leave/Join so the - // live call isn't disrupted (preserves "applies on next join"). - if active_session.is_none() { - let lookup = net.memory_lookup.clone(); - net.shutdown().await; - let publish = presence_mode.lock().unwrap().publishes_to_discovery(); - net = build_net_stack(secret_key.clone(), network_mode, lookup, friends_handler.clone(), publish).await?; - } else { - net_rebuild_pending = true; + // Skip when the posture is unchanged. The GUI re-sends the saved + // network mode as part of its startup config-sync, and that mode + // usually already matches the freshly-built stack — rebuilding the + // iroh endpoint for an identical posture just churns the network + // and adds a needless ~1s teardown+rebuild bounce at every launch + // (seen on both Linux and Windows/Wine). A real change still + // rebuilds exactly as before. + if mode != network_mode { + network_mode = mode; + // Rebuild the persistent stack to the new posture immediately if + // idle; if a call is active, defer to the next Leave/Join so the + // live call isn't disrupted (preserves "applies on next join"). + if active_session.is_none() { + let lookup = net.memory_lookup.clone(); + net.shutdown().await; + let publish = presence_mode.lock().unwrap().publishes_to_discovery(); + net = build_net_stack(secret_key.clone(), network_mode, lookup, friends_handler.clone(), publish).await?; + } else { + net_rebuild_pending = true; + } } } diff --git a/src/lib.rs b/src/lib.rs index e6369ed..b110029 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -24,6 +24,9 @@ use std::path::{Path, PathBuf}; use std::sync::OnceLock; const LOG_MAX_BYTES: u64 = 5 * 1024 * 1024; +// Owner-only log permissions are a Unix concept (mode bits); on Windows the log +// inherits the directory's default ACL. Only referenced under `cfg(unix)`. +#[cfg(unix)] const LOG_MODE: u32 = 0o600; /// Resolves the log file path once: `$XDG_STATE_HOME/peerspeak/peerspeak.log` @@ -84,8 +87,6 @@ fn prepare_log_file(path: &Path) -> std::io::Result { } fn prepare_log_file_with_limit(path: &Path, max_bytes: u64) -> std::io::Result { - use std::os::unix::fs::{OpenOptionsExt, PermissionsExt}; - if let Some(parent) = path.parent() { let _ = std::fs::create_dir_all(parent); } @@ -98,12 +99,23 @@ fn prepare_log_file_with_limit(path: &Path, max_bytes: u64) -> std::io::Result PathBuf { @@ -145,6 +158,9 @@ mod tests { assert_eq!(redact_for_log(" "), ""); } + // Owner-only log perms are a Unix concept; on Windows the file inherits the + // directory ACL and there's no mode to assert. + #[cfg(unix)] #[test] fn log_file_is_created_private() { let dir = temp_log_dir(); diff --git a/src/notify.rs b/src/notify.rs index 0e7baf9..58355d9 100644 --- a/src/notify.rs +++ b/src/notify.rs @@ -3,11 +3,12 @@ //! //! The WAVs are embedded in the binary (`include_bytes!`) so a deployed single //! binary is self-contained — no asset directory to ship alongside it. On first -//! use each sound is written once to a temp file, then played fire-and-forget -//! via `pw-play` (PipeWire-native; falls back to `paplay`/`aplay`). Playback runs -//! on a detached thread that waits on the child, so it never blocks the UI and -//! never leaves a zombie. Any failure (no player, no audio) is silent by design — -//! a missing chime should never disrupt a call. +//! use each sound is written once to a temp file, then played fire-and-forget. +//! Linux uses `pw-play` (PipeWire-native; falls back to `paplay`/`aplay`); +//! Windows uses PowerShell's `System.Media.SoundPlayer`. Playback runs on a +//! detached thread that waits on the child, so it never blocks the UI and never +//! leaves a zombie. Any failure (no player, no audio) is silent by design — a +//! missing chime should never disrupt a call. use std::collections::HashMap; use std::path::{Path, PathBuf}; @@ -202,8 +203,14 @@ fn cached_path(sound: Sound) -> Option { Some(path) } +#[cfg(any(windows, test))] +fn escape_powershell_single_quoted(s: &str) -> String { + s.replace('\'', "''") +} + /// Try each available player in turn, waiting on the first that starts (which /// reaps the child). Runs on a detached thread, so the wait is harmless. +#[cfg(not(windows))] fn spawn_player(path: &Path) { for player in ["pw-play", "paplay", "aplay"] { let started = Command::new(player) @@ -221,6 +228,23 @@ fn spawn_player(path: &Path) { } } +/// Play through Windows' built-in WAV player. Runs on a detached thread, so +/// `PlaySync()` blocking for the sound duration is fine. +#[cfg(windows)] +fn spawn_player(path: &Path) { + let path = escape_powershell_single_quoted(&path.display().to_string()); + let command = format!("(New-Object System.Media.SoundPlayer '{path}').PlaySync()"); + let _ = Command::new("powershell") + .arg("-NoProfile") + .arg("-NonInteractive") + .arg("-Command") + .arg(command) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status(); +} + #[cfg(test)] mod tests { use super::*; @@ -234,6 +258,18 @@ mod tests { assert!(!should_play(false, false)); } + #[test] + fn test_powershell_single_quote_escape() { + assert_eq!( + escape_powershell_single_quoted(r"C:\Users\O'Brien\chime.wav"), + r"C:\Users\O''Brien\chime.wav" + ); + assert_eq!( + escape_powershell_single_quoted("a'b'c"), + "a''b''c" + ); + } + #[test] fn test_sound_indices_unique_and_match_all() { // `index()` must be a 0..COUNT bijection in `ALL` order, or the flag diff --git a/src/screenshare/mod.rs b/src/screenshare/mod.rs index 2029220..6187673 100644 --- a/src/screenshare/mod.rs +++ b/src/screenshare/mod.rs @@ -25,6 +25,16 @@ use tokio::process::{Child, Command}; /// points elsewhere. const PIXELPASS_BIN: &str = "pixelpass"; +#[cfg(windows)] +fn pixelpass_path_candidates(dir: &Path) -> [PathBuf; 2] { + [dir.join(PIXELPASS_BIN), dir.join("pixelpass.exe")] +} + +#[cfg(not(windows))] +fn pixelpass_path_candidates(dir: &Path) -> [PathBuf; 1] { + [dir.join(PIXELPASS_BIN)] +} + /// Pixelpass endpoint tickets are normally ~140 chars. Leave headroom for format /// growth, but reject unbounded gossip payloads before the UI offers "Watch". const MAX_TICKET_LEN: usize = 512; @@ -143,7 +153,7 @@ pub fn pixelpass_path(config_override: Option<&str>) -> Option { } let path_var = std::env::var_os("PATH")?; std::env::split_paths(&path_var) - .map(|dir| dir.join(PIXELPASS_BIN)) + .flat_map(|dir| pixelpass_path_candidates(&dir)) .find(|c| c.is_file()) } @@ -513,4 +523,14 @@ mod tests { // only assert it doesn't return the empty path as a match. assert_ne!(pixelpass_path(Some(" ")).as_deref(), Some(Path::new(""))); } + + #[test] + fn pixelpass_path_candidates_are_platform_specific() { + let dir = Path::new("bin"); + let candidates: Vec = pixelpass_path_candidates(dir).into_iter().collect(); + #[cfg(windows)] + assert_eq!(candidates, vec![dir.join("pixelpass"), dir.join("pixelpass.exe")]); + #[cfg(not(windows))] + assert_eq!(candidates, vec![dir.join("pixelpass")]); + } }