feat(voice): replaced miniaudio with native platform audio backends

- Replaced the miniaudio dependency with custom OS audio device abstractions and backends.
- Implemented platform-specific audio playback and capture for macOS (Audio Queue), Windows (WASAPI), and Linux (PulseAudio/ALSA).
- Added a fallback stub backend returning errors for unsupported platforms.
- Updated audio stream handling with reliable fill guard wakeups and streamlined rate validation.
This commit is contained in:
can1357
2026-08-07 23:34:16 +02:00
parent bc5eee4e24
commit 7cae7ef3f5
17 changed files with 2555 additions and 14386 deletions
Generated
+2 -22
View File
@@ -576,8 +576,6 @@ dependencies = [
"cexpr",
"clang-sys",
"itertools 0.13.0",
"log",
"prettyplease",
"proc-macro2",
"quote",
"regex",
@@ -3726,25 +3724,6 @@ dependencies = [
"web_atoms",
]
[[package]]
name = "maudio"
version = "0.1.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ab8050247e9de440204a366e45b93068a5772a7e5be1c9dcc33ee5366214ecd0"
dependencies = [
"maudio-sys",
]
[[package]]
name = "maudio-sys"
version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "06b3b2d00f3a690dfea26eed932022c3f0366a9cb15785f2bf84717de8277d35"
dependencies = [
"bindgen",
"cc",
]
[[package]]
name = "md-5"
version = "0.10.6"
@@ -5144,11 +5123,12 @@ dependencies = [
"audiopus_sys",
"bytes",
"flume",
"maudio",
"libc",
"opus",
"parking_lot",
"tokio",
"webrtc",
"windows-sys 0.61.2",
]
[[package]]
+4 -21
View File
@@ -33,20 +33,6 @@ strip = true
# (see crates/pi-natives/src/crash_handler.rs).
panic = "unwind"
# rustc ICEs codegenning `MaybeUninit<maudio_sys::ffi::ma_fence>` — a union with
# ScalarPair backend repr on x86_64-pc-windows-msvc — once the MIR opts that run
# from opt-level 2 upward const-fold it into an `Uninit` operand
# (`rustc_codegen_ssa/src/mir/block.rs`: "codegen_argument: OperandRef(Uninit …)
# invalid for pair argument"). Hit by the rustup toolchain this repo pins
# (nightly-2026-07-28) on the local napi path (`packages/natives` →
# `build:bindings`); Bazel's older pin (nightly/2026-04-29, MODULE.bazel) is
# unaffected, which is why CI never saw it. opt-level 1 keeps those MIR opts off.
# maudio is a thin FFI shim — miniaudio itself is C compiled by `cc` at its own
# optimization level — so the cost is noise. Drop this once the rustup pin moves
# past the fix. `ci` and `local` inherit this override from `release`.
[profile.release.package.maudio]
opt-level = 1
[profile.ci]
inherits = "release"
lto = "thin"
@@ -75,9 +61,6 @@ split-debuginfo = "unpacked"
opt-level = 2
debug = false
[profile.dev.package.maudio]
opt-level = 1
[workspace.lints.rust]
# ──────────────────────────────────────────────────────────────────────────────
# Rust Lint Levels
@@ -256,14 +239,10 @@ smallvec = { version = "1.15.1", features = [
# ──────────────────────────────────────────────────────────────────────────────
xxhash-rust = { version = "0.8", features = ["xxh64"] }
# ──────────────────────────────────────────────────────────────────────────────
# Audio & Realtime Media
# ──────────────────────────────────────────────────────────────────────────────
audiopus_sys = { version = "0.2.2", features = ["static"] }
# Pregenerated bindings (crate default): keeps bindgen/libclang out of the
# build graph so Bazel actions stay hermetic.
maudio = "0.1.6"
opus = "0.3.1"
webrtc = "0.17.2"
@@ -274,12 +253,16 @@ libc = "0.2"
os_pipe = "1"
windows-sys = { version = "0.61", features = [
"Win32_Foundation",
"Win32_Media_Audio",
"Win32_Media_Multimedia",
"Win32_Security",
"Win32_Storage_FileSystem",
"Win32_Storage_ProjectedFileSystem",
"Win32_System_Com",
"Win32_System_LibraryLoader",
"Win32_System_IO",
"Win32_System_Ioctl",
"Win32_System_Threading",
] }
winreg = "0.56"
+2 -2
View File
@@ -34,8 +34,8 @@ ENV BUN_INSTALL=/opt/bun \
PATH=/opt/bun/bin:/usr/local/cargo/bin:/usr/local/bin:/usr/bin:/bin \
CARGO_TERM_COLOR=never
# clang/libclang-dev: bindgen for maudio-sys (miniaudio); cmake/make/ninja-build:
# audiopus_sys builds bundled libopus via CMake (native audio stack, 17.1.1+).
# clang/libclang-dev: bindgen for pipewire-sys/libspa-sys (Linux desktop capture);
# cmake/make/ninja-build: audiopus_sys builds bundled libopus via CMake.
# bazelisk: hermetic bazel launcher for the native addon build (17.1.5+).
RUN apt-get update \
&& apt-get install -y --no-install-recommends \
-11
View File
@@ -199,17 +199,6 @@ crate.annotation(
build_script_env = {"CFLAGS": "-UNDEBUG"},
)
# `generate-bindings` is enabled for darwin so miniaudio's Rust bindings match
# its CoreAudio C layout. Build scripts execute for the host, however, so a
# macOS-hosted Linux cross-build also sees that feature. Use Cargo's TARGET to
# keep bindgen darwin-only and select checked GNU/musl bindings by compilation
# target, including ARM64's larger glibc pthread layouts.
crate.annotation(
crate = "maudio-sys",
patches = ["//bazel/patches:maudio-sys-target-bindings.patch"],
patch_args = ["-p1"],
)
use_repo(crate, "crates")
File diff suppressed because it is too large Load Diff
+2 -11
View File
@@ -5,8 +5,8 @@ package(default_visibility = ["//visibility:public"])
exports_files(["Cargo.toml"])
# napi-free voice engine (miniaudio capture/playback + WebRTC/Opus live peer).
# Keeping the webrtc/opus/maudio graph behind this rlib means the pi_natives
# napi-free voice engine (in-house capture/playback backends + WebRTC/Opus live
# peer). Keeping the webrtc/opus graph behind this rlib means the pi_natives
# leaf — which recompiles on every release (version-sentinel edit in lib.rs) —
# no longer pays for the voice stack's monomorphization.
rust_library(
@@ -28,12 +28,3 @@ rust_test(
deps = all_crate_deps(normal_dev = True),
)
# Direct C/Rust ABI comparison for miniaudio's opaque device state on the host.
# The patched crate also carries paired target-gated C and Rust static assertions
# so every cross-built Linux binding gets equivalent compile-time validation.
rust_test(
name = "maudio_layout_test",
srcs = ["bazel/maudio_layout.rs"],
edition = "2024",
deps = ["@crates//:maudio"],
)
+7 -12
View File
@@ -14,20 +14,15 @@ workspace = true
audiopus_sys.workspace = true
bytes.workspace = true
flume.workspace = true
maudio.workspace = true
opus.workspace = true
parking_lot.workspace = true
tokio.workspace = true
webrtc.workspace = true
[target.'cfg(target_os = "macos")'.dependencies]
# Regenerate miniaudio FFI bindings from the compiled headers on macOS only.
# The maudio-sys pregenerated `unix.rs` bindings are Linux-shaped: their
# `ma_device`/`ma_context` backend-state unions omit the `#ifdef
# MA_SUPPORT_COREAUDIO` members that the macOS build of miniaudio.c actually
# contains, so shipping them on darwin is an ABI mismatch that silently breaks
# CoreAudio capture (device init/start succeed but no render callbacks fire).
# darwin builds with host Xcode (libclang available), so bindgen is fine here;
# the hermetic zig/musl/windows toolchains keep the bindgen-free pregen path,
# whose bindings are correct for those platforms.
maudio = { workspace = true, features = ["generate-bindings"] }
[target.'cfg(target_os = "linux")'.dependencies]
# dlopen/dlsym for the PulseAudio and ALSA backends: no link-time dependency
# on libpulse/libasound, so the prebuilt addon loads on systems without them.
libc.workspace = true
[target.'cfg(target_os = "windows")'.dependencies]
windows-sys.workspace = true
-39
View File
@@ -1,39 +0,0 @@
use std::mem::{align_of, size_of};
use maudio::maudio_sys::ffi::{ma_context, ma_device, ma_mutex};
unsafe extern "C" {
fn omp_maudio_sizeof_mutex() -> usize;
fn omp_maudio_alignof_mutex() -> usize;
fn omp_maudio_sizeof_context() -> usize;
fn omp_maudio_alignof_context() -> usize;
fn omp_maudio_sizeof_device() -> usize;
fn omp_maudio_alignof_device() -> usize;
}
fn c_layout(
size: unsafe extern "C" fn() -> usize,
align: unsafe extern "C" fn() -> usize,
) -> (usize, usize) {
// SAFETY: The patched maudio-sys C object defines both argument-free probes.
unsafe { (size(), align()) }
}
#[test]
fn rust_bindings_match_compiled_miniaudio_layouts() {
assert_eq!(
(size_of::<ma_mutex>(), align_of::<ma_mutex>()),
c_layout(omp_maudio_sizeof_mutex, omp_maudio_alignof_mutex),
"ma_mutex Rust/C layout mismatch",
);
assert_eq!(
(size_of::<ma_context>(), align_of::<ma_context>()),
c_layout(omp_maudio_sizeof_context, omp_maudio_alignof_context),
"ma_context Rust/C layout mismatch",
);
assert_eq!(
(size_of::<ma_device>(), align_of::<ma_device>()),
c_layout(omp_maudio_sizeof_device, omp_maudio_alignof_device),
"ma_device Rust/C layout mismatch",
);
}
+59 -70
View File
@@ -1,9 +1,10 @@
//! Cross-platform microphone capture and streaming speaker playback.
//!
//! miniaudio owns platform device discovery, format conversion, channel mixing,
//! and resampling. The engine exposes one stable mono `f32` contract: the
//! N-API classes in pi-natives adapt it to TypeScript, and [`crate::live`]
//! shares [`PlaybackStream`] for remote-audio rendering.
//! The per-platform backends in [`crate::device`] own device access, format
//! conversion, channel mixing, and resampling. The engine exposes one stable
//! mono `f32` contract: the N-API classes in pi-natives adapt it to
//! TypeScript, and [`crate::live`] shares [`PlaybackStream`] for remote-audio
//! rendering.
use std::sync::{
Arc,
@@ -11,46 +12,32 @@ use std::sync::{
};
use flume::TryRecvError;
use maudio::{
audio::{performance::PerformanceProfile, sample_rate::SampleRate},
backend::Backend,
device::{
Device,
device_builder::{DeviceBuilder, DeviceBuilderOps},
},
};
use tokio::sync::Notify;
use crate::VoiceResult;
use crate::{
VoiceResult,
device::{CaptureDevice, DeviceConfig, PlaybackDevice},
};
const AUDIO_CHANNELS: u32 = 1;
// PulseAudio TCP playback stutters with a 20 ms target buffer; 50 ms absorbs
// transport jitter while preserving interactive latency.
#[cfg(target_os = "linux")]
const PLAYBACK_PERIOD_MS: u32 = 50;
#[cfg(not(target_os = "linux"))]
const PLAYBACK_PERIOD_MS: u32 = 20;
// miniaudio's PulseAudio backend reserves three periods. Android's OpenSL ES
// source emits 125 ms fragments, so Linux capture needs at least 150 ms queued.
#[cfg(target_os = "linux")]
const CAPTURE_PERIOD_MS: u32 = 50;
#[cfg(not(target_os = "linux"))]
const CAPTURE_PERIOD_MS: u32 = 20;
// PulseAudio can retain its default three periods after the producer closes.
// Wait for all of them before stopping the device so the tail reaches the sink.
#[cfg(target_os = "linux")]
const PLAYBACK_DRAIN_CALLBACKS: usize = 3;
#[cfg(not(target_os = "linux"))]
const PLAYBACK_DRAIN_CALLBACKS: usize = 2;
#[cfg(target_os = "macos")]
const AUDIO_BACKENDS: &[Backend] = &[Backend::CoreAudio];
#[cfg(target_os = "windows")]
const AUDIO_BACKENDS: &[Backend] = &[Backend::Wasapi];
#[cfg(target_os = "linux")]
const AUDIO_BACKENDS: &[Backend] = &[Backend::PulseAudio, Backend::Alsa, Backend::Jack];
#[cfg(not(any(target_os = "macos", target_os = "windows", target_os = "linux")))]
const AUDIO_BACKENDS: &[Backend] = &[Backend::Sndio, Backend::Audio4, Backend::Oss];
// Backends queue up to three periods (AudioQueue buffers, WASAPI padding cap,
// Pulse `maxlength`/ALSA buffer). Draining needs three silence periods
// COMMITTED to the OS behind the tail: once the third is accepted into a
// three-period FIFO, everything ahead of it has played. The callback that
// marks drained races teardown — on Linux the delivery gate may cancel that
// callback's own write after `wait_for_drain` wakes — so count one extra
// empty callback: the racy, possibly-uncommitted write is always the fourth,
// which is margin rather than accounted flush.
const PLAYBACK_DRAIN_CALLBACKS: usize = 4;
/// Shared render-time state for one playback device: gain, drain, stop.
///
@@ -105,6 +92,18 @@ impl PlaybackState {
}
}
/// Wakes drain waiters when the backend drops the fill callback (device loss
/// or stop) so `wait_for_drain` can never outlive the render path.
struct FillGuard {
state: Arc<PlaybackState>,
}
impl Drop for FillGuard {
fn drop(&mut self) {
self.state.mark_stopped();
}
}
/// Producer endpoint for one native playback device. Cloned into the WebRTC
/// remote-audio decoder so it can feed the same speaker stream.
#[derive(Clone)]
@@ -131,7 +130,7 @@ impl PlaybackWriter {
/// Running mono playback stream shared by N-API playback and native WebRTC.
pub struct PlaybackStream {
device: Option<Device<f32>>,
device: Option<PlaybackDevice>,
writer: Option<PlaybackWriter>,
state: Arc<PlaybackState>,
}
@@ -146,15 +145,15 @@ impl PlaybackStream {
let mut current = Vec::new();
let mut cursor = 0;
let mut empty_callbacks = 0;
let mut builder = DeviceBuilder::playback().f32();
builder
.sample_rate(sample_rate)
.playback_channels(AUDIO_CHANNELS)
.period_size_millis(PLAYBACK_PERIOD_MS)
.performance_profile(PerformanceProfile::LowLatency)
.backends(AUDIO_BACKENDS);
let mut device = builder
.with_callback(move |_device, output| {
let config = DeviceConfig { sample_rate, period_ms: PLAYBACK_PERIOD_MS };
// The guard travels inside the fill closure: if the backend drops the
// callback for any reason (worker exit on device loss, stop), waiters
// blocked in `wait_for_drain` wake instead of hanging forever.
let guard = FillGuard { state: Arc::clone(&state) };
let device = PlaybackDevice::start(
config,
Box::new(move |output| {
let _ = &guard;
fill_playback(
&rx,
&mut current,
@@ -163,11 +162,9 @@ impl PlaybackStream {
&callback_state,
&mut empty_callbacks,
);
})
}),
)
.map_err(|error| format!("Failed to open the default speaker: {error}"))?;
device
.device_start()
.map_err(|error| format!("Failed to start speaker playback: {error}"))?;
Ok(Self {
device: Some(device),
@@ -212,9 +209,7 @@ impl PlaybackStream {
let Some(mut device) = self.device.take() else {
return Ok(());
};
device
.device_stop()
.map_err(|error| format!("Failed to stop speaker playback: {error}"))
device.stop()
}
}
@@ -224,9 +219,12 @@ impl Drop for PlaybackStream {
}
}
fn audio_sample_rate(sample_rate: u32) -> VoiceResult<SampleRate> {
SampleRate::try_from(sample_rate)
.map_err(|error| format!("Unsupported audio sample rate {sample_rate}: {error}"))
/// Bounds the logical rate to what OS converters accept before device open.
fn audio_sample_rate(sample_rate: u32) -> VoiceResult<u32> {
if !(8_000..=384_000).contains(&sample_rate) {
return Err(format!("Unsupported audio sample rate {sample_rate}"));
}
Ok(sample_rate)
}
fn fill_playback(
@@ -282,10 +280,10 @@ fn fill_playback(
}
/// Running default-microphone capture delivering low-latency mono `f32`
/// chunks to its callback. Wraps the miniaudio device so N-API callers never
/// see maudio types.
/// chunks to its callback. Wraps the platform device so N-API callers never
/// see backend types.
pub struct CaptureStream {
device: Option<Device<f32>>,
device: Option<CaptureDevice>,
}
impl CaptureStream {
@@ -296,23 +294,16 @@ impl CaptureStream {
C: FnMut(&[f32]) + Send + 'static,
{
let sample_rate = audio_sample_rate(sample_rate)?;
let mut builder = DeviceBuilder::capture().f32();
builder
.sample_rate(sample_rate)
.capture_channels(AUDIO_CHANNELS)
.period_size_millis(CAPTURE_PERIOD_MS)
.performance_profile(PerformanceProfile::LowLatency)
.backends(AUDIO_BACKENDS);
let mut device = builder
.with_callback(move |_device, samples| {
let config = DeviceConfig { sample_rate, period_ms: CAPTURE_PERIOD_MS };
let device = CaptureDevice::start(
config,
Box::new(move |samples| {
if !samples.is_empty() {
on_audio(samples);
}
})
}),
)
.map_err(|error| format!("Failed to open the default microphone: {error}"))?;
device
.device_start()
.map_err(|error| format!("Failed to start microphone capture: {error}"))?;
Ok(Self { device: Some(device) })
}
@@ -321,9 +312,7 @@ impl CaptureStream {
let Some(mut device) = self.device.take() else {
return Ok(());
};
device
.device_stop()
.map_err(|error| format!("Failed to stop microphone capture: {error}"))
device.stop()
}
}
+475
View File
@@ -0,0 +1,475 @@
//! macOS default-device audio backend using `AudioToolbox` Audio Queues.
use std::{
ffi::c_void,
mem::size_of,
ptr, slice,
sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
};
use super::{CaptureSink, DeviceConfig, PlaybackFill};
use crate::VoiceResult;
const BUFFER_COUNT: usize = 3;
const LINEAR_PCM: u32 = 0x6c70_636d;
const FORMAT_FLAGS: u32 = 0x9;
type AudioQueueRef = *mut AudioQueueOpaque;
type AudioTimeStamp = c_void;
type AudioStreamPacketDescription = c_void;
#[repr(C)]
struct AudioQueueOpaque {
_private: [u8; 0],
}
struct QueueHandle(AudioQueueRef);
// SAFETY: AudioQueue control functions explicitly support calls from arbitrary
// threads.
unsafe impl Send for QueueHandle {}
impl QueueHandle {
fn stop_and_dispose(self) -> VoiceResult<()> {
// SAFETY: This handle owns a live queue; an immediate stop waits for queue
// activity.
let stop_status = unsafe { AudioQueueStop(self.0, 1) };
// SAFETY: The synchronous stop completed, and this handle exclusively owns the
// queue.
let dispose_status = unsafe { AudioQueueDispose(self.0, 1) };
if stop_status != 0 {
return Err(format!("CoreAudio queue stop failed (OSStatus {stop_status})"));
}
if dispose_status != 0 {
return Err(format!("CoreAudio queue dispose failed (OSStatus {dispose_status})"));
}
Ok(())
}
}
#[repr(C)]
#[allow(non_snake_case, reason = "fields must match the CoreAudio C ABI")]
struct AudioStreamBasicDescription {
mSampleRate: f64,
mFormatID: u32,
mFormatFlags: u32,
mBytesPerPacket: u32,
mFramesPerPacket: u32,
mBytesPerFrame: u32,
mChannelsPerFrame: u32,
mBitsPerChannel: u32,
mReserved: u32,
}
#[repr(C)]
#[allow(non_snake_case, reason = "fields must match the AudioQueue C ABI")]
struct AudioQueueBuffer {
mAudioDataBytesCapacity: u32,
mAudioData: *mut c_void,
mAudioDataByteSize: u32,
mUserData: *mut c_void,
mPacketDescriptionCapacity: u32,
mPacketDescriptions: *mut c_void,
mPacketDescriptionCount: u32,
}
#[link(name = "AudioToolbox", kind = "framework")]
unsafe extern "C" {
fn AudioQueueNewOutput(
format: *const AudioStreamBasicDescription,
callback: unsafe extern "C" fn(*mut c_void, AudioQueueRef, *mut AudioQueueBuffer),
user_data: *mut c_void,
callback_run_loop: *const c_void,
callback_run_loop_mode: *const c_void,
flags: u32,
queue: *mut AudioQueueRef,
) -> i32;
fn AudioQueueNewInput(
format: *const AudioStreamBasicDescription,
callback: unsafe extern "C" fn(
*mut c_void,
AudioQueueRef,
*mut AudioQueueBuffer,
*const AudioTimeStamp,
u32,
*const AudioStreamPacketDescription,
),
user_data: *mut c_void,
callback_run_loop: *const c_void,
callback_run_loop_mode: *const c_void,
flags: u32,
queue: *mut AudioQueueRef,
) -> i32;
fn AudioQueueAllocateBuffer(
queue: AudioQueueRef,
buffer_byte_size: u32,
buffer: *mut *mut AudioQueueBuffer,
) -> i32;
fn AudioQueueEnqueueBuffer(
queue: AudioQueueRef,
buffer: *mut AudioQueueBuffer,
packet_description_count: u32,
packet_descriptions: *const AudioStreamPacketDescription,
) -> i32;
fn AudioQueueStart(queue: AudioQueueRef, start_time: *const AudioTimeStamp) -> i32;
fn AudioQueueStop(queue: AudioQueueRef, immediate: u8) -> i32;
fn AudioQueueDispose(queue: AudioQueueRef, immediate: u8) -> i32;
}
unsafe extern "C" {
fn pthread_self() -> usize;
}
struct PlaybackContext {
fill: PlaybackFill,
stopped: Arc<AtomicBool>,
callback_thread: Arc<AtomicUsize>,
}
struct CaptureContext {
sink: CaptureSink,
stopped: Arc<AtomicBool>,
callback_thread: Arc<AtomicUsize>,
}
fn stream_format(sample_rate: u32) -> AudioStreamBasicDescription {
AudioStreamBasicDescription {
mSampleRate: f64::from(sample_rate),
mFormatID: LINEAR_PCM,
mFormatFlags: FORMAT_FLAGS,
mBytesPerPacket: size_of::<f32>() as u32,
mFramesPerPacket: 1,
mBytesPerFrame: size_of::<f32>() as u32,
mChannelsPerFrame: 1,
mBitsPerChannel: 32,
mReserved: 0,
}
}
fn buffer_size(config: DeviceConfig) -> VoiceResult<u32> {
config
.period_samples()
.checked_mul(size_of::<f32>())
.and_then(|bytes| u32::try_from(bytes).ok())
.ok_or_else(|| "CoreAudio period buffer is too large".to_owned())
}
fn dispose_failed_start(queue: AudioQueueRef, operation: &str, status: i32) -> String {
if !queue.is_null() {
// SAFETY: `queue` was returned by AudioQueueNewInput/Output and has not been
// disposed.
unsafe { AudioQueueDispose(queue, 1) };
}
format!("CoreAudio {operation} failed (OSStatus {status})")
}
unsafe extern "C" fn playback_callback(
user_data: *mut c_void,
queue: AudioQueueRef,
buffer: *mut AudioQueueBuffer,
) {
if user_data.is_null() || queue.is_null() || buffer.is_null() {
return;
}
let context = user_data.cast::<PlaybackContext>();
// SAFETY: AudioQueue passes the live context pointer supplied when the queue
// was created.
unsafe {
(*context)
.callback_thread
.store(pthread_self(), Ordering::Release);
};
// SAFETY: The callback only projects the independently allocated atomic flag
// from `context`.
if unsafe { (*context).stopped.load(Ordering::Acquire) } {
// SAFETY: The context remains live until synchronous queue disposal completes.
unsafe { (*context).callback_thread.store(0, Ordering::Release) };
return;
}
// SAFETY: AudioQueue passes one of its allocated buffers exclusively to this
// callback.
let buffer = unsafe { &mut *buffer };
let sample_count = buffer.mAudioDataBytesCapacity as usize / size_of::<f32>();
// SAFETY: AudioQueue allocated `mAudioData` with the reported capacity for
// linear PCM data.
let samples =
unsafe { slice::from_raw_parts_mut(buffer.mAudioData.cast::<f32>(), sample_count) };
// SAFETY: AudioQueue serializes callbacks, so only this callback borrows the
// `fill` field.
unsafe { ((*context).fill)(samples) };
// SAFETY: The callback only projects the independently allocated atomic flag
// from `context`.
if !unsafe { (*context).stopped.load(Ordering::Acquire) } {
buffer.mAudioDataByteSize = buffer.mAudioDataBytesCapacity;
// SAFETY: The queue and buffer belong to this callback and remain live while it
// returns.
let _ = unsafe { AudioQueueEnqueueBuffer(queue, buffer, 0, ptr::null()) };
}
// SAFETY: The context remains live until synchronous queue disposal completes.
unsafe { (*context).callback_thread.store(0, Ordering::Release) };
}
unsafe extern "C" fn capture_callback(
user_data: *mut c_void,
queue: AudioQueueRef,
buffer: *mut AudioQueueBuffer,
_start_time: *const AudioTimeStamp,
_packet_count: u32,
_packet_descriptions: *const AudioStreamPacketDescription,
) {
if user_data.is_null() || queue.is_null() || buffer.is_null() {
return;
}
let context = user_data.cast::<CaptureContext>();
// SAFETY: AudioQueue passes the live context pointer supplied when the queue
// was created.
unsafe {
(*context)
.callback_thread
.store(pthread_self(), Ordering::Release);
};
// SAFETY: The callback only projects the independently allocated atomic flag
// from `context`.
if unsafe { (*context).stopped.load(Ordering::Acquire) } {
// SAFETY: The context remains live until synchronous queue disposal completes.
unsafe { (*context).callback_thread.store(0, Ordering::Release) };
return;
}
// SAFETY: AudioQueue passes one of its allocated buffers exclusively to this
// callback.
let buffer = unsafe { &mut *buffer };
let byte_size = buffer.mAudioDataByteSize as usize;
if byte_size != 0 && byte_size.is_multiple_of(size_of::<f32>()) {
// SAFETY: AudioQueue filled `mAudioDataByteSize` bytes within this allocated
// buffer.
let samples = unsafe {
slice::from_raw_parts(buffer.mAudioData.cast::<f32>(), byte_size / size_of::<f32>())
};
// SAFETY: AudioQueue serializes callbacks, so only this callback borrows the
// `sink` field.
unsafe { ((*context).sink)(samples) };
}
// SAFETY: The callback only projects the independently allocated atomic flag
// from `context`.
if !unsafe { (*context).stopped.load(Ordering::Acquire) } {
// SAFETY: The queue and buffer belong to this callback and remain live while it
// returns.
let _ = unsafe { AudioQueueEnqueueBuffer(queue, buffer, 0, ptr::null()) };
}
// SAFETY: The context remains live until synchronous queue disposal completes.
unsafe { (*context).callback_thread.store(0, Ordering::Release) };
}
/// Running `CoreAudio` default-speaker queue.
pub(crate) struct PlaybackDevice {
queue: Option<QueueHandle>,
context: Option<Box<PlaybackContext>>,
stopped: Arc<AtomicBool>,
callback_thread: Arc<AtomicUsize>,
}
// SAFETY: AudioQueue control functions may be called from any thread, and the
// callback is `Send`.
unsafe impl Send for PlaybackDevice {}
impl PlaybackDevice {
/// Open and start the default speaker queue.
pub fn start(config: DeviceConfig, fill: PlaybackFill) -> VoiceResult<Self> {
let byte_size = buffer_size(config)?;
let format = stream_format(config.sample_rate);
let stopped = Arc::new(AtomicBool::new(false));
let callback_thread = Arc::new(AtomicUsize::new(0));
let mut context = Box::new(PlaybackContext {
fill,
stopped: Arc::clone(&stopped),
callback_thread: Arc::clone(&callback_thread),
});
let user_data = ptr::from_mut(&mut *context).cast::<c_void>();
let mut queue = ptr::null_mut();
// SAFETY: All pointers are valid for the call; the boxed context outlives the
// queue.
let status = unsafe {
AudioQueueNewOutput(
&format,
playback_callback,
user_data,
ptr::null(),
ptr::null(),
0,
&mut queue,
)
};
if status != 0 {
return Err(dispose_failed_start(queue, "queue creation", status));
}
for _ in 0..BUFFER_COUNT {
let mut buffer = ptr::null_mut();
// SAFETY: `queue` is live and `buffer` points to writable storage for the
// result.
let status = unsafe { AudioQueueAllocateBuffer(queue, byte_size, &mut buffer) };
if status != 0 {
return Err(dispose_failed_start(queue, "buffer allocation", status));
}
// SAFETY: AudioQueue returned a valid buffer with at least `byte_size` writable
// bytes.
let buffer_ref = unsafe { &mut *buffer };
let sample_count = buffer_ref.mAudioDataBytesCapacity as usize / size_of::<f32>();
// SAFETY: AudioQueue allocated the data pointer with the reported capacity.
let samples =
unsafe { slice::from_raw_parts_mut(buffer_ref.mAudioData.cast::<f32>(), sample_count) };
(context.fill)(samples);
buffer_ref.mAudioDataByteSize = buffer_ref.mAudioDataBytesCapacity;
// SAFETY: `queue` and `buffer` are live, and PCM requires no packet
// descriptions.
let status = unsafe { AudioQueueEnqueueBuffer(queue, buffer, 0, ptr::null()) };
if status != 0 {
return Err(dispose_failed_start(queue, "buffer enqueue", status));
}
}
// SAFETY: `queue` is live and null requests immediate start.
let status = unsafe { AudioQueueStart(queue, ptr::null()) };
if status != 0 {
return Err(dispose_failed_start(queue, "queue start", status));
}
Ok(Self { queue: Some(QueueHandle(queue)), context: Some(context), stopped, callback_thread })
}
/// Stop playback and dispose the queue, handing off teardown from its
/// callback thread.
pub fn stop(&mut self) -> VoiceResult<()> {
self.stopped.store(true, Ordering::Release);
let Some(queue) = self.queue.take() else {
return Ok(());
};
let Some(context) = self.context.take() else {
self.queue = Some(queue);
return Err("CoreAudio playback queue lost its callback context".to_owned());
};
// SAFETY: `pthread_self` returns the stable identifier for the calling thread.
let current_thread = unsafe { pthread_self() };
if current_thread != 0 && self.callback_thread.load(Ordering::Acquire) == current_thread {
let stopped = Arc::clone(&self.stopped);
drop(std::thread::spawn(move || {
stopped.store(true, Ordering::Release);
let _ = queue.stop_and_dispose();
drop(context);
}));
return Ok(());
}
let result = queue.stop_and_dispose();
drop(context);
result
}
}
impl Drop for PlaybackDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
/// Running `CoreAudio` default-microphone queue.
pub(crate) struct CaptureDevice {
queue: Option<QueueHandle>,
context: Option<Box<CaptureContext>>,
stopped: Arc<AtomicBool>,
callback_thread: Arc<AtomicUsize>,
}
// SAFETY: AudioQueue control functions may be called from any thread, and the
// callback is `Send`.
unsafe impl Send for CaptureDevice {}
impl CaptureDevice {
/// Open and start the default microphone queue.
pub fn start(config: DeviceConfig, sink: CaptureSink) -> VoiceResult<Self> {
let byte_size = buffer_size(config)?;
let format = stream_format(config.sample_rate);
let stopped = Arc::new(AtomicBool::new(false));
let callback_thread = Arc::new(AtomicUsize::new(0));
let mut context = Box::new(CaptureContext {
sink,
stopped: Arc::clone(&stopped),
callback_thread: Arc::clone(&callback_thread),
});
let user_data = ptr::from_mut(&mut *context).cast::<c_void>();
let mut queue = ptr::null_mut();
// SAFETY: All pointers are valid for the call; the boxed context outlives the
// queue.
let status = unsafe {
AudioQueueNewInput(
&format,
capture_callback,
user_data,
ptr::null(),
ptr::null(),
0,
&mut queue,
)
};
if status != 0 {
return Err(dispose_failed_start(queue, "queue creation", status));
}
for _ in 0..BUFFER_COUNT {
let mut buffer = ptr::null_mut();
// SAFETY: `queue` is live and `buffer` points to writable storage for the
// result.
let status = unsafe { AudioQueueAllocateBuffer(queue, byte_size, &mut buffer) };
if status != 0 {
return Err(dispose_failed_start(queue, "buffer allocation", status));
}
// SAFETY: `queue` and `buffer` are live, and input PCM requires no packet
// descriptions.
let status = unsafe { AudioQueueEnqueueBuffer(queue, buffer, 0, ptr::null()) };
if status != 0 {
return Err(dispose_failed_start(queue, "buffer enqueue", status));
}
}
// SAFETY: `queue` is live and null requests immediate start.
let status = unsafe { AudioQueueStart(queue, ptr::null()) };
if status != 0 {
return Err(dispose_failed_start(queue, "queue start", status));
}
Ok(Self { queue: Some(QueueHandle(queue)), context: Some(context), stopped, callback_thread })
}
/// Stop capture and dispose the queue, handing off teardown from its
/// callback thread.
pub fn stop(&mut self) -> VoiceResult<()> {
self.stopped.store(true, Ordering::Release);
let Some(queue) = self.queue.take() else {
return Ok(());
};
let Some(context) = self.context.take() else {
self.queue = Some(queue);
return Err("CoreAudio capture queue lost its callback context".to_owned());
};
// SAFETY: `pthread_self` returns the stable identifier for the calling thread.
let current_thread = unsafe { pthread_self() };
if current_thread != 0 && self.callback_thread.load(Ordering::Acquire) == current_thread {
let stopped = Arc::clone(&self.stopped);
drop(std::thread::spawn(move || {
stopped.store(true, Ordering::Release);
let _ = queue.stop_and_dispose();
drop(context);
}));
return Ok(());
}
let result = queue.stop_and_dispose();
drop(context);
result
}
}
impl Drop for CaptureDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
+919
View File
@@ -0,0 +1,919 @@
//! Linux default-device audio through runtime-loaded `PulseAudio` or ALSA.
use std::{
ffi::{CStr, c_char, c_int, c_long, c_uint, c_void},
ptr,
sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, Ordering},
mpsc,
},
thread::{self, JoinHandle},
time::Duration,
};
use super::{CaptureSink, DeviceConfig, PlaybackFill};
use crate::VoiceResult;
const PA_STREAM_PLAYBACK: c_int = 1;
const PA_STREAM_RECORD: c_int = 2;
#[cfg(target_endian = "little")]
const PA_SAMPLE_FLOAT32_NATIVE: c_int = 5;
#[cfg(target_endian = "big")]
const PA_SAMPLE_FLOAT32_NATIVE: c_int = 6;
const SND_PCM_STREAM_PLAYBACK: c_int = 0;
const SND_PCM_STREAM_CAPTURE: c_int = 1;
const SND_PCM_NONBLOCK: c_int = 1;
const SND_PCM_ACCESS_RW_INTERLEAVED: c_int = 3;
#[cfg(target_endian = "little")]
const SND_PCM_FORMAT_FLOAT_NATIVE: c_int = 14;
#[cfg(target_endian = "big")]
const SND_PCM_FORMAT_FLOAT_NATIVE: c_int = 15;
#[repr(C)]
struct PaSampleSpec {
format: c_int,
rate: u32,
channels: u8,
}
#[repr(C)]
struct PaBufferAttr {
maxlength: u32,
tlength: u32,
prebuf: u32,
minreq: u32,
fragsize: u32,
}
type PaSimpleNew = unsafe extern "C" fn(
*const c_char,
*const c_char,
c_int,
*const c_char,
*const c_char,
*const PaSampleSpec,
*const c_void,
*const PaBufferAttr,
*mut c_int,
) -> *mut c_void;
type PaSimpleFree = unsafe extern "C" fn(*mut c_void);
type PaSimpleWrite = unsafe extern "C" fn(*mut c_void, *const c_void, usize, *mut c_int) -> c_int;
type PaSimpleRead = unsafe extern "C" fn(*mut c_void, *mut c_void, usize, *mut c_int) -> c_int;
type PaStrerror = unsafe extern "C" fn(c_int) -> *const c_char;
struct PulseApi {
simple_new: PaSimpleNew,
simple_free: PaSimpleFree,
simple_write: PaSimpleWrite,
simple_read: PaSimpleRead,
strerror: PaStrerror,
}
static PULSE_API: OnceLock<Result<&'static PulseApi, String>> = OnceLock::new();
impl PulseApi {
fn get() -> Result<&'static Self, String> {
PULSE_API.get_or_init(Self::load).clone()
}
fn load() -> Result<&'static Self, String> {
let simple = open_library(c"libpulse-simple.so.0", libc::RTLD_NOW | libc::RTLD_GLOBAL)?;
let pulse = open_library(c"libpulse.so.0", libc::RTLD_NOW | libc::RTLD_GLOBAL)?;
// SAFETY: each symbol is resolved from the library defining this exact C API.
let simple_new = unsafe {
std::mem::transmute::<*mut c_void, PaSimpleNew>(symbol(simple, c"pa_simple_new")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let simple_free = unsafe {
std::mem::transmute::<*mut c_void, PaSimpleFree>(symbol(simple, c"pa_simple_free")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let simple_write = unsafe {
std::mem::transmute::<*mut c_void, PaSimpleWrite>(symbol(simple, c"pa_simple_write")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let simple_read = unsafe {
std::mem::transmute::<*mut c_void, PaSimpleRead>(symbol(simple, c"pa_simple_read")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let strerror =
unsafe { std::mem::transmute::<*mut c_void, PaStrerror>(symbol(pulse, c"pa_strerror")?) };
Ok(Box::leak(Box::new(Self { simple_new, simple_free, simple_write, simple_read, strerror })))
}
fn error(&self, code: c_int) -> String {
// SAFETY: pa_strerror accepts every PulseAudio error code and returns a static
// string.
let message = unsafe { (self.strerror)(code) };
cstring_lossy(message, "unknown PulseAudio error")
}
}
type SndPcmOpen = unsafe extern "C" fn(*mut *mut c_void, *const c_char, c_int, c_int) -> c_int;
type SndPcmSetParams =
unsafe extern "C" fn(*mut c_void, c_int, c_int, c_uint, c_uint, c_int, c_uint) -> c_int;
type SndPcmIo = unsafe extern "C" fn(*mut c_void, *mut c_void, c_long) -> c_long;
type SndPcmRecover = unsafe extern "C" fn(*mut c_void, c_int, c_int) -> c_int;
type SndPcmWait = unsafe extern "C" fn(*mut c_void, c_int) -> c_int;
type SndPcmControl = unsafe extern "C" fn(*mut c_void) -> c_int;
type SndStrerror = unsafe extern "C" fn(c_int) -> *const c_char;
struct AlsaApi {
pcm_open: SndPcmOpen,
pcm_set_params: SndPcmSetParams,
pcm_writei: SndPcmIo,
pcm_readi: SndPcmIo,
pcm_recover: SndPcmRecover,
pcm_wait: SndPcmWait,
pcm_start: SndPcmControl,
pcm_close: SndPcmControl,
strerror: SndStrerror,
}
static ALSA_API: OnceLock<Result<&'static AlsaApi, String>> = OnceLock::new();
impl AlsaApi {
fn get() -> Result<&'static Self, String> {
ALSA_API.get_or_init(Self::load).clone()
}
fn load() -> Result<&'static Self, String> {
let library = open_library(c"libasound.so.2", libc::RTLD_NOW | libc::RTLD_GLOBAL)?;
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_open = unsafe {
std::mem::transmute::<*mut c_void, SndPcmOpen>(symbol(library, c"snd_pcm_open")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_set_params = unsafe {
std::mem::transmute::<*mut c_void, SndPcmSetParams>(symbol(
library,
c"snd_pcm_set_params",
)?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_writei = unsafe {
std::mem::transmute::<*mut c_void, SndPcmIo>(symbol(library, c"snd_pcm_writei")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_readi = unsafe {
std::mem::transmute::<*mut c_void, SndPcmIo>(symbol(library, c"snd_pcm_readi")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_recover = unsafe {
std::mem::transmute::<*mut c_void, SndPcmRecover>(symbol(library, c"snd_pcm_recover")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_wait = unsafe {
std::mem::transmute::<*mut c_void, SndPcmWait>(symbol(library, c"snd_pcm_wait")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_start = unsafe {
std::mem::transmute::<*mut c_void, SndPcmControl>(symbol(library, c"snd_pcm_start")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let pcm_close = unsafe {
std::mem::transmute::<*mut c_void, SndPcmControl>(symbol(library, c"snd_pcm_close")?)
};
// SAFETY: each symbol is resolved from the library defining this exact C API.
let strerror = unsafe {
std::mem::transmute::<*mut c_void, SndStrerror>(symbol(library, c"snd_strerror")?)
};
Ok(Box::leak(Box::new(Self {
pcm_open,
pcm_set_params,
pcm_writei,
pcm_readi,
pcm_recover,
pcm_wait,
pcm_start,
pcm_close,
strerror,
})))
}
fn error(&self, code: c_int) -> String {
// SAFETY: snd_strerror accepts ALSA status codes and returns a static string.
let message = unsafe { (self.strerror)(code) };
cstring_lossy(message, "unknown ALSA error")
}
}
fn open_library(name: &CStr, flags: c_int) -> Result<*mut c_void, String> {
// SAFETY: name is a valid NUL-terminated string and flags are supported by
// dlopen.
let library = unsafe { libc::dlopen(name.as_ptr(), flags) };
if library.is_null() {
Err(format!("could not load {}: {}", name.to_string_lossy(), dlerror()))
} else {
Ok(library)
}
}
fn symbol(library: *mut c_void, name: &CStr) -> Result<*mut c_void, String> {
// SAFETY: library is a live handle deliberately retained for process lifetime.
unsafe { libc::dlerror() };
// SAFETY: library is live and name is a valid NUL-terminated symbol name.
let address = unsafe { libc::dlsym(library, name.as_ptr()) };
// SAFETY: dlerror reads and clears this thread's dynamic-loader error state.
let error = unsafe { libc::dlerror() };
if error.is_null() {
Ok(address)
} else {
Err(format!(
"could not resolve {}: {}",
name.to_string_lossy(),
cstring_lossy(error, "unknown dynamic-loader error")
))
}
}
fn dlerror() -> String {
// SAFETY: dlerror returns either null or a thread-local NUL-terminated string.
let error = unsafe { libc::dlerror() };
cstring_lossy(error, "unknown dynamic-loader error")
}
fn cstring_lossy(value: *const c_char, fallback: &str) -> String {
if value.is_null() {
fallback.to_owned()
} else {
// SAFETY: callers only pass pointers returned by APIs specifying NUL-terminated
// strings.
unsafe { CStr::from_ptr(value) }
.to_string_lossy()
.into_owned()
}
}
struct PulseStream(*mut c_void);
// SAFETY: the pointer is transferred to and exclusively dereferenced by its
// worker thread.
unsafe impl Send for PulseStream {}
impl PulseStream {
fn open(
api: &PulseApi,
config: DeviceConfig,
direction: c_int,
attr: &PaBufferAttr,
) -> Result<Self, String> {
let spec = PaSampleSpec {
format: PA_SAMPLE_FLOAT32_NATIVE,
rate: config.sample_rate,
channels: 1,
};
let mut error = 0;
// SAFETY: all pointers reference valid values for the duration of
// pa_simple_new.
let stream = unsafe {
(api.simple_new)(
ptr::null(),
c"oh-my-pi".as_ptr(),
direction,
ptr::null(),
c"voice".as_ptr(),
&raw const spec,
ptr::null(),
&raw const *attr,
&raw mut error,
)
};
if stream.is_null() {
Err(format!("PulseAudio open failed: {}", api.error(error)))
} else {
Ok(Self(stream))
}
}
}
struct AlsaStream(*mut c_void);
// SAFETY: the pointer is transferred to and exclusively dereferenced by its
// worker thread.
unsafe impl Send for AlsaStream {}
impl AlsaStream {
fn open(api: &AlsaApi, config: DeviceConfig, direction: c_int) -> Result<Self, String> {
let latency = config
.period_ms
.checked_mul(3)
.and_then(|value| value.checked_mul(1000))
.ok_or_else(|| "audio period is too large".to_owned())?;
match Self::open_named(api, config, direction, latency, c"default") {
Ok(stream) => Ok(stream),
Err(AlsaOpenError::Open(error)) => Err(error),
Err(AlsaOpenError::Params(default_error)) => {
Self::open_named(api, config, direction, latency, c"plug:default").map_err(|error| {
format!("{default_error}; plug:default fallback failed: {}", error.message())
})
},
}
}
fn open_named(
api: &AlsaApi,
config: DeviceConfig,
direction: c_int,
latency: u32,
name: &CStr,
) -> Result<Self, AlsaOpenError> {
let mut pcm = ptr::null_mut();
// SAFETY: pcm is valid output storage and name is NUL-terminated.
let status =
unsafe { (api.pcm_open)(&raw mut pcm, name.as_ptr(), direction, SND_PCM_NONBLOCK) };
if status < 0 {
return Err(AlsaOpenError::Open(format!(
"ALSA open of {} failed: {}",
name.to_string_lossy(),
api.error(status)
)));
}
let stream = Self(pcm);
// SAFETY: pcm is an open handle owned by this thread and all enum values match
// ALSA.
let status = unsafe {
(api.pcm_set_params)(
stream.0,
SND_PCM_FORMAT_FLOAT_NATIVE,
SND_PCM_ACCESS_RW_INTERLEAVED,
1,
config.sample_rate,
1,
latency,
)
};
if status < 0 {
// SAFETY: stream owns this open handle and it has not been closed.
unsafe { (api.pcm_close)(stream.0) };
Err(AlsaOpenError::Params(format!(
"ALSA parameter setup on {} failed: {}",
name.to_string_lossy(),
api.error(status)
)))
} else {
Ok(stream)
}
}
}
enum AlsaOpenError {
Open(String),
Params(String),
}
impl AlsaOpenError {
fn message(self) -> String {
match self {
Self::Open(message) | Self::Params(message) => message,
}
}
}
fn pulse_attr(config: DeviceConfig) -> Result<PaBufferAttr, String> {
let period_bytes = config
.period_samples()
.checked_mul(size_of::<f32>())
.and_then(|value| u32::try_from(value).ok())
.ok_or_else(|| "audio period is too large".to_owned())?;
let target_bytes = period_bytes
.checked_mul(3)
.ok_or_else(|| "audio period is too large".to_owned())?;
Ok(PaBufferAttr {
maxlength: target_bytes,
tlength: period_bytes,
prebuf: u32::MAX,
minreq: u32::MAX,
fragsize: period_bytes,
})
}
fn remember_error(slot: &Mutex<Option<String>>, error: String) {
if let Ok(mut stored) = slot.lock() {
*stored = Some(error);
}
}
type DeliveryGate = (AtomicBool, parking_lot::Mutex<()>);
fn fill_if_armed(gate: &DeliveryGate, fill: &mut PlaybackFill, buffer: &mut [f32]) -> bool {
if !gate.0.load(Ordering::Acquire) {
return false;
}
let _delivery = gate.1.lock();
if !gate.0.load(Ordering::Acquire) {
return false;
}
fill(buffer);
gate.0.load(Ordering::Acquire)
}
fn sink_if_armed(gate: &DeliveryGate, sink: &mut CaptureSink, buffer: &[f32]) -> bool {
if !gate.0.load(Ordering::Acquire) {
return false;
}
let _delivery = gate.1.lock();
if !gate.0.load(Ordering::Acquire) {
return false;
}
sink(buffer);
gate.0.load(Ordering::Acquire)
}
fn pulse_playback_loop(
api: &PulseApi,
stream: &PulseStream,
stop: &AtomicBool,
gate: &DeliveryGate,
fill: &mut PlaybackFill,
samples: usize,
) -> Result<(), String> {
let mut buffer = vec![0.0_f32; samples];
while !stop.load(Ordering::Acquire) {
if !fill_if_armed(gate, fill, &mut buffer) {
break;
}
let mut error = 0;
// SAFETY: stream is open and buffer contains exactly the supplied byte count.
let status = unsafe {
(api.simple_write)(
stream.0,
buffer.as_ptr().cast(),
size_of_val(buffer.as_slice()),
&raw mut error,
)
};
if status < 0 {
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.simple_free)(stream.0) };
return Err(format!("PulseAudio playback failed: {}", api.error(error)));
}
}
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.simple_free)(stream.0) };
Ok(())
}
fn pulse_capture_loop(
api: &PulseApi,
stream: &PulseStream,
stop: &AtomicBool,
gate: &DeliveryGate,
sink: &mut CaptureSink,
samples: usize,
) -> Result<(), String> {
let mut buffer = vec![0.0_f32; samples];
while !stop.load(Ordering::Acquire) {
let mut error = 0;
// SAFETY: stream is open and buffer has writable storage for the supplied byte
// count.
let status = unsafe {
(api.simple_read)(
stream.0,
buffer.as_mut_ptr().cast(),
size_of_val(buffer.as_slice()),
&raw mut error,
)
};
if status < 0 {
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.simple_free)(stream.0) };
return Err(format!("PulseAudio capture failed: {}", api.error(error)));
}
if !sink_if_armed(gate, sink, &buffer) {
break;
}
}
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.simple_free)(stream.0) };
Ok(())
}
fn close_alsa_with_error(api: &AlsaApi, stream: &AlsaStream, error: String) -> Result<(), String> {
// SAFETY: stream is open and exclusively owned by this thread.
unsafe { (api.pcm_close)(stream.0) };
Err(error)
}
fn recover_alsa(
api: &AlsaApi,
stream: &AlsaStream,
error: c_int,
context: &str,
) -> Result<(), String> {
// SAFETY: stream is open and error came from an operation on this stream.
let recovered = unsafe { (api.pcm_recover)(stream.0, error, 1) };
if recovered < 0 {
return Err(format!("{context}: {}", api.error(recovered)));
}
// SAFETY: stream is prepared after successful recovery and remains
// worker-owned.
let _ = unsafe { (api.pcm_start)(stream.0) };
Ok(())
}
fn wait_for_alsa(api: &AlsaApi, stream: &AlsaStream, timeout_ms: c_int) -> Result<(), String> {
// SAFETY: stream is open and timeout_ms is a valid non-negative timeout.
let status = unsafe { (api.pcm_wait)(stream.0, timeout_ms) };
if status >= 0 || status == -libc::EINTR {
Ok(())
} else {
recover_alsa(api, stream, status, "ALSA wait recovery failed")
}
}
fn alsa_playback_loop(
api: &AlsaApi,
stream: &AlsaStream,
stop: &AtomicBool,
gate: &DeliveryGate,
fill: &mut PlaybackFill,
samples: usize,
timeout_ms: c_int,
) -> Result<(), String> {
let mut buffer = vec![0.0_f32; samples];
while !stop.load(Ordering::Acquire) {
if !fill_if_armed(gate, fill, &mut buffer) {
break;
}
let mut offset = 0;
while offset < samples && !stop.load(Ordering::Acquire) {
let Ok(frames) = c_long::try_from(samples - offset) else {
return close_alsa_with_error(
api,
stream,
"audio period exceeds ALSA frame range".to_owned(),
);
};
// SAFETY: stream is open and the remaining buffer contains frames of mono f32
// audio.
let status = unsafe {
(api.pcm_writei)(stream.0, buffer.as_ptr().add(offset).cast_mut().cast(), frames)
};
if status == -c_long::from(libc::EAGAIN) {
if let Err(error) = wait_for_alsa(api, stream, timeout_ms) {
return close_alsa_with_error(api, stream, error);
}
} else if status < 0 {
let Ok(error) = c_int::try_from(status) else {
return close_alsa_with_error(
api,
stream,
format!("ALSA returned an invalid status: {status}"),
);
};
if let Err(error) = recover_alsa(api, stream, error, "ALSA playback recovery failed") {
return close_alsa_with_error(api, stream, error);
}
} else if status > 0 {
let Ok(written) = usize::try_from(status) else {
return close_alsa_with_error(
api,
stream,
format!("ALSA returned an invalid frame count: {status}"),
);
};
if written > samples - offset {
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.pcm_close)(stream.0) };
return Err(format!(
"ALSA wrote {written} frames after receiving {}",
samples - offset
));
}
offset += written;
}
}
}
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.pcm_close)(stream.0) };
Ok(())
}
fn alsa_capture_loop(
api: &AlsaApi,
stream: &AlsaStream,
stop: &AtomicBool,
gate: &DeliveryGate,
sink: &mut CaptureSink,
samples: usize,
timeout_ms: c_int,
) -> Result<(), String> {
let mut buffer = vec![0.0_f32; samples];
let Ok(frames) = c_long::try_from(samples) else {
return close_alsa_with_error(
api,
stream,
"audio period exceeds ALSA frame range".to_owned(),
);
};
while !stop.load(Ordering::Acquire) {
// SAFETY: stream is open and buffer has writable storage for frames mono f32
// frames.
let status = unsafe { (api.pcm_readi)(stream.0, buffer.as_mut_ptr().cast(), frames) };
if status == -c_long::from(libc::EAGAIN) {
if let Err(error) = wait_for_alsa(api, stream, timeout_ms) {
return close_alsa_with_error(api, stream, error);
}
continue;
}
if status < 0 {
let Ok(error) = c_int::try_from(status) else {
return close_alsa_with_error(
api,
stream,
format!("ALSA returned an invalid status: {status}"),
);
};
if let Err(error) = recover_alsa(api, stream, error, "ALSA capture recovery failed") {
return close_alsa_with_error(api, stream, error);
}
continue;
}
if status > 0 {
let Ok(captured) = usize::try_from(status) else {
return close_alsa_with_error(
api,
stream,
format!("ALSA returned an invalid frame count: {status}"),
);
};
if captured > samples {
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.pcm_close)(stream.0) };
return Err(format!("ALSA captured {captured} frames into a {samples}-frame buffer"));
}
if !sink_if_armed(gate, sink, &buffer[..captured]) {
break;
}
}
}
// SAFETY: stream is still open and exclusively owned by this thread.
unsafe { (api.pcm_close)(stream.0) };
Ok(())
}
struct ThreadDone(mpsc::Sender<()>);
impl Drop for ThreadDone {
fn drop(&mut self) {
let _ = self.0.send(());
}
}
struct RunningDevice {
stop: AtomicBool,
delivery: Arc<DeliveryGate>,
worker_id: OnceLock<thread::ThreadId>,
error: Mutex<Option<String>>,
}
fn finish(
device: &Arc<RunningDevice>,
thread: &mut Option<JoinHandle<()>>,
done: &mut Option<mpsc::Receiver<()>>,
) -> VoiceResult<()> {
device.delivery.0.store(false, Ordering::Release);
device.stop.store(true, Ordering::Release);
if device
.worker_id
.get()
.is_some_and(|worker_id| *worker_id == thread::current().id())
{
// A callback-thread stop cannot wait on its own delivery lock or join
// itself. Disarming is sufficient; the worker exits after the callback.
return Ok(());
}
// The gate waits out an in-flight callback and prevents every future one,
// keeping the bounded PulseAudio detach path contract-clean.
drop(device.delivery.1.lock());
if let Some(handle) = thread.take() {
let completed = done.take().is_none_or(|receiver| {
match receiver.recv_timeout(Duration::from_millis(500)) {
Ok(()) | Err(mpsc::RecvTimeoutError::Disconnected) => true,
Err(mpsc::RecvTimeoutError::Timeout) => false,
}
});
if completed {
handle
.join()
.map_err(|_| "audio worker thread panicked".to_owned())?;
} else {
// A pathological PulseAudio server can stall pa_simple I/O forever.
// Detaching keeps stop/Drop bounded; the worker owns and eventually frees
// the handle if the server ever unblocks. The delivery gate prevents callbacks.
drop(handle);
}
}
device
.error
.lock()
.map_err(|_| "audio worker error state was poisoned".to_owned())?
.take()
.map_or(Ok(()), Err)
}
/// Running `PulseAudio` or ALSA playback worker.
pub(crate) struct PlaybackDevice {
device: Arc<RunningDevice>,
thread: Option<JoinHandle<()>>,
done: Option<mpsc::Receiver<()>>,
}
impl PlaybackDevice {
/// Opens the default playback device and starts its worker thread.
pub fn start(config: DeviceConfig, mut fill: PlaybackFill) -> VoiceResult<Self> {
let samples = config.period_samples();
let attr = pulse_attr(config)?;
let timeout_ms = c_int::try_from(config.period_ms)
.unwrap_or(c_int::MAX)
.max(1);
let delivery = Arc::new((AtomicBool::new(true), parking_lot::Mutex::new(())));
let device = Arc::new(RunningDevice {
stop: AtomicBool::new(false),
delivery,
error: Mutex::new(None),
worker_id: OnceLock::new(),
});
let worker_device = Arc::clone(&device);
let (opened_tx, opened_rx) = mpsc::sync_channel(1);
let (done_tx, done_rx) = mpsc::channel();
let thread = thread::Builder::new()
.name("pi-voice-playback".to_owned())
.spawn(move || {
let _done = ThreadDone(done_tx);
let _ = worker_device.worker_id.set(thread::current().id());
let pulse_error = match PulseApi::get().and_then(|api| {
PulseStream::open(api, config, PA_STREAM_PLAYBACK, &attr).map(|stream| (api, stream))
}) {
Ok((api, stream)) => {
let _ = opened_tx.send(Ok(()));
if let Err(error) = pulse_playback_loop(
api,
&stream,
&worker_device.stop,
worker_device.delivery.as_ref(),
&mut fill,
samples,
) {
remember_error(&worker_device.error, error);
}
return;
},
Err(error) => error,
};
match AlsaApi::get().and_then(|api| {
AlsaStream::open(api, config, SND_PCM_STREAM_PLAYBACK).map(|stream| (api, stream))
}) {
Ok((api, stream)) => {
let _ = opened_tx.send(Ok(()));
if let Err(error) = alsa_playback_loop(
api,
&stream,
&worker_device.stop,
worker_device.delivery.as_ref(),
&mut fill,
samples,
timeout_ms,
) {
remember_error(&worker_device.error, error);
}
},
Err(alsa_error) => {
let _ = opened_tx.send(Err(format!(
"no Linux playback backend available; PulseAudio: {pulse_error}; ALSA: \
{alsa_error}"
)));
},
}
})
.map_err(|error| format!("could not start playback worker: {error}"))?;
match opened_rx.recv() {
Ok(Ok(())) => Ok(Self { device, thread: Some(thread), done: Some(done_rx) }),
Ok(Err(error)) => {
let _ = thread.join();
Err(error)
},
Err(error) => {
let _ = thread.join();
Err(format!("playback worker exited during startup: {error}"))
},
}
}
/// Stops playback, waiting out delivery when called off the worker thread.
pub fn stop(&mut self) -> VoiceResult<()> {
finish(&self.device, &mut self.thread, &mut self.done)
}
}
impl Drop for PlaybackDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
/// Running `PulseAudio` or ALSA capture worker.
pub(crate) struct CaptureDevice {
device: Arc<RunningDevice>,
thread: Option<JoinHandle<()>>,
done: Option<mpsc::Receiver<()>>,
}
impl CaptureDevice {
/// Opens the default capture device and starts its worker thread.
pub fn start(config: DeviceConfig, mut sink: CaptureSink) -> VoiceResult<Self> {
let samples = config.period_samples();
let attr = pulse_attr(config)?;
let timeout_ms = c_int::try_from(config.period_ms)
.unwrap_or(c_int::MAX)
.max(1);
let delivery = Arc::new((AtomicBool::new(true), parking_lot::Mutex::new(())));
let device = Arc::new(RunningDevice {
stop: AtomicBool::new(false),
delivery,
error: Mutex::new(None),
worker_id: OnceLock::new(),
});
let worker_device = Arc::clone(&device);
let (opened_tx, opened_rx) = mpsc::sync_channel(1);
let (done_tx, done_rx) = mpsc::channel();
let thread = thread::Builder::new()
.name("pi-voice-capture".to_owned())
.spawn(move || {
let _done = ThreadDone(done_tx);
let _ = worker_device.worker_id.set(thread::current().id());
let pulse_error = match PulseApi::get().and_then(|api| {
PulseStream::open(api, config, PA_STREAM_RECORD, &attr).map(|stream| (api, stream))
}) {
Ok((api, stream)) => {
let _ = opened_tx.send(Ok(()));
if let Err(error) = pulse_capture_loop(
api,
&stream,
&worker_device.stop,
worker_device.delivery.as_ref(),
&mut sink,
samples,
) {
remember_error(&worker_device.error, error);
}
return;
},
Err(error) => error,
};
match AlsaApi::get().and_then(|api| {
AlsaStream::open(api, config, SND_PCM_STREAM_CAPTURE).map(|stream| (api, stream))
}) {
Ok((api, stream)) => {
let _ = opened_tx.send(Ok(()));
if let Err(error) = alsa_capture_loop(
api,
&stream,
&worker_device.stop,
worker_device.delivery.as_ref(),
&mut sink,
samples,
timeout_ms,
) {
remember_error(&worker_device.error, error);
}
},
Err(alsa_error) => {
let _ = opened_tx.send(Err(format!(
"no Linux capture backend available; PulseAudio: {pulse_error}; ALSA: \
{alsa_error}"
)));
},
}
})
.map_err(|error| format!("could not start capture worker: {error}"))?;
match opened_rx.recv() {
Ok(Ok(())) => Ok(Self { device, thread: Some(thread), done: Some(done_rx) }),
Ok(Err(error)) => {
let _ = thread.join();
Err(error)
},
Err(error) => {
let _ = thread.join();
Err(format!("capture worker exited during startup: {error}"))
},
}
}
/// Stops capture, waiting out delivery when called off the worker thread.
pub fn stop(&mut self) -> VoiceResult<()> {
finish(&self.device, &mut self.thread, &mut self.done)
}
}
impl Drop for CaptureDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
+106
View File
@@ -0,0 +1,106 @@
//! In-house default-device audio backends.
//!
//! One backend per platform, each implementing the same two devices against
//! the OS audio API directly: `CoreAudio` `AudioQueue` on macOS, shared-mode
//! WASAPI with automatic format conversion on Windows, and the `PulseAudio`
//! simple API with an ALSA fallback (both loaded via `dlopen`) on Linux.
//! Every backend delegates format conversion, channel mixing, and resampling
//! to the OS so the engine keeps a single mono `f32` contract at the
//! requested logical sample rate.
//!
//! # Contract
//! - [`PlaybackDevice::start`] opens the default speaker and invokes `fill`
//! with a mono `f32` buffer roughly every [`DeviceConfig::period_ms`]. The
//! callback runs on a backend-owned audio thread and must not block. Backends
//! keep at most three periods queued OS-side; the engine's drain accounting
//! (`PLAYBACK_DRAIN_CALLBACKS`) depends on that bound.
//! - [`CaptureDevice::start`] opens the default microphone and invokes `sink`
//! with non-empty mono `f32` chunks at the requested sample rate, also from a
//! backend-owned thread.
//! - `stop` is idempotent, callable from any thread, and guarantees no callback
//! is running or will run after it returns — except when invoked from within
//! a device callback itself: the calling thread cannot await its own
//! cessation, so backends defer teardown to another thread and the
//! post-return guarantee applies only to external callers. The engine never
//! stops from callbacks; the carve-out exists for contract soundness, not for
//! use. Dropping a device stops it.
#[cfg(target_os = "macos")]
mod coreaudio;
#[cfg(target_os = "macos")]
use coreaudio as imp;
#[cfg(target_os = "windows")]
mod wasapi;
#[cfg(target_os = "windows")]
use wasapi as imp;
#[cfg(target_os = "linux")]
mod linux;
#[cfg(target_os = "linux")]
use linux as imp;
#[cfg(not(any(target_os = "macos", target_os = "windows", target_os = "linux")))]
mod unsupported;
#[cfg(not(any(target_os = "macos", target_os = "windows", target_os = "linux")))]
use unsupported as imp;
use crate::VoiceResult;
/// Render callback: fill the whole output buffer with mono `f32` samples.
pub(crate) type PlaybackFill = Box<dyn FnMut(&mut [f32]) + Send + 'static>;
/// Capture callback: consume a non-empty chunk of mono `f32` samples.
pub(crate) type CaptureSink = Box<dyn FnMut(&[f32]) + Send + 'static>;
/// Stream parameters shared by both device directions.
#[derive(Clone, Copy)]
pub(crate) struct DeviceConfig {
/// Logical client-side sample rate in Hz; the OS converts to hardware.
pub sample_rate: u32,
/// Target callback period in milliseconds.
pub period_ms: u32,
}
impl DeviceConfig {
/// Samples per callback period at the logical rate (never zero).
pub fn period_samples(self) -> usize {
((self.sample_rate as usize * self.period_ms as usize) / 1000).max(1)
}
}
/// Running default-speaker playback stream driven by a fill callback.
pub(crate) struct PlaybackDevice {
inner: imp::PlaybackDevice,
}
impl PlaybackDevice {
/// Open and start the default speaker; `fill` runs on the audio thread.
pub fn start(config: DeviceConfig, fill: PlaybackFill) -> VoiceResult<Self> {
Ok(Self { inner: imp::PlaybackDevice::start(config, fill)? })
}
/// Stop playback and release the device. Idempotent; no callback runs
/// after this returns.
pub fn stop(&mut self) -> VoiceResult<()> {
self.inner.stop()
}
}
/// Running default-microphone capture stream driven by a sink callback.
pub(crate) struct CaptureDevice {
inner: imp::CaptureDevice,
}
impl CaptureDevice {
/// Open and start the default microphone; `sink` runs on the audio thread.
pub fn start(config: DeviceConfig, sink: CaptureSink) -> VoiceResult<Self> {
Ok(Self { inner: imp::CaptureDevice::start(config, sink)? })
}
/// Stop capture and release the device. Idempotent; no callback runs
/// after this returns.
pub fn stop(&mut self) -> VoiceResult<()> {
self.inner.stop()
}
}
+30
View File
@@ -0,0 +1,30 @@
//! Stub backend for platforms without a native audio implementation.
use super::{CaptureSink, DeviceConfig, PlaybackFill};
use crate::VoiceResult;
const UNSUPPORTED: &str = "Native audio is not supported on this platform";
pub(crate) struct PlaybackDevice;
impl PlaybackDevice {
pub fn start(_config: DeviceConfig, _fill: PlaybackFill) -> VoiceResult<Self> {
Err(UNSUPPORTED.to_owned())
}
pub fn stop(&mut self) -> VoiceResult<()> {
Ok(())
}
}
pub(crate) struct CaptureDevice;
impl CaptureDevice {
pub fn start(_config: DeviceConfig, _sink: CaptureSink) -> VoiceResult<Self> {
Err(UNSUPPORTED.to_owned())
}
pub fn stop(&mut self) -> VoiceResult<()> {
Ok(())
}
}
+934
View File
@@ -0,0 +1,934 @@
//! Shared-mode WASAPI playback and capture for the default Windows devices.
use std::{
ffi::c_void,
ptr::{null, null_mut},
slice,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
mpsc::{self, Sender},
},
thread::{self, JoinHandle},
time::Duration,
};
use windows_sys::{
Win32::{
Foundation::{CloseHandle, HANDLE, WAIT_OBJECT_0, WAIT_TIMEOUT},
Media::{
Audio::{
AUDCLNT_BUFFERFLAGS_SILENT, AUDCLNT_E_DEVICE_INVALIDATED, AUDCLNT_SHAREMODE_SHARED,
AUDCLNT_STREAMFLAGS_AUTOCONVERTPCM, AUDCLNT_STREAMFLAGS_EVENTCALLBACK,
AUDCLNT_STREAMFLAGS_SRC_DEFAULT_QUALITY, WAVEFORMATEX, eCapture, eConsole, eRender,
},
Multimedia::WAVE_FORMAT_IEEE_FLOAT,
},
System::{
Com::{
CLSCTX_ALL, COINIT_MULTITHREADED, CoCreateInstance, CoInitializeEx, CoUninitialize,
},
Threading::{CreateEventW, SetEvent, WaitForSingleObject},
},
},
core::{GUID, HRESULT, IUnknown_Vtbl},
};
use super::{CaptureSink, DeviceConfig, PlaybackFill};
use crate::VoiceResult;
const WAIT_TIMEOUT_MS: u32 = 2_000;
const REOPEN_ATTEMPTS: usize = 4;
const REOPEN_BACKOFF: Duration = Duration::from_millis(200);
const CLSID_MMDEVICE_ENUMERATOR: GUID = GUID::from_u128(0xbcde_0395_e52f_467c_8e3d_c457_9291_692e);
const IID_IMMDEVICE_ENUMERATOR: GUID = GUID::from_u128(0xa956_64d2_9614_4f35_a746_de8d_b636_17e6);
const IID_IAUDIO_CLIENT: GUID = GUID::from_u128(0x1cb9_ad4c_dbfa_4c32_b178_c2f5_68a7_03b2);
const IID_IAUDIO_RENDER_CLIENT: GUID = GUID::from_u128(0xf294_acfc_3146_4483_a7bf_addc_a7c2_60e2);
const IID_IAUDIO_CAPTURE_CLIENT: GUID = GUID::from_u128(0xc8ad_bd64_e71e_48a0_a4de_185c_395c_d317);
#[repr(C)]
struct RawComInterface<V> {
vtable: *const V,
}
trait ComVtable {
fn unknown(&self) -> &IUnknown_Vtbl;
}
#[repr(C)]
#[allow(dead_code, reason = "all slots are required to preserve the COM vtable layout")]
struct MmDeviceEnumeratorVtable {
base: IUnknown_Vtbl,
enum_audio_endpoints:
unsafe extern "system" fn(*mut c_void, i32, u32, *mut *mut c_void) -> HRESULT,
get_default_audio_endpoint:
unsafe extern "system" fn(*mut c_void, i32, i32, *mut *mut c_void) -> HRESULT,
get_device: unsafe extern "system" fn(*mut c_void, *const u16, *mut *mut c_void) -> HRESULT,
register_endpoint_notification_callback:
unsafe extern "system" fn(*mut c_void, *mut c_void) -> HRESULT,
unregister_endpoint_notification_callback:
unsafe extern "system" fn(*mut c_void, *mut c_void) -> HRESULT,
}
impl ComVtable for MmDeviceEnumeratorVtable {
fn unknown(&self) -> &IUnknown_Vtbl {
&self.base
}
}
#[repr(C)]
#[allow(dead_code, reason = "all slots are required to preserve the COM vtable layout")]
struct MmDeviceVtable {
base: IUnknown_Vtbl,
activate: unsafe extern "system" fn(
*mut c_void,
*const GUID,
u32,
*const c_void,
*mut *mut c_void,
) -> HRESULT,
open_property_store: unsafe extern "system" fn(*mut c_void, u32, *mut *mut c_void) -> HRESULT,
get_id: unsafe extern "system" fn(*mut c_void, *mut *mut u16) -> HRESULT,
get_state: unsafe extern "system" fn(*mut c_void, *mut u32) -> HRESULT,
}
impl ComVtable for MmDeviceVtable {
fn unknown(&self) -> &IUnknown_Vtbl {
&self.base
}
}
#[repr(C)]
#[allow(dead_code, reason = "all slots are required to preserve the COM vtable layout")]
struct AudioClientVtable {
base: IUnknown_Vtbl,
initialize: unsafe extern "system" fn(
*mut c_void,
i32,
u32,
i64,
i64,
*const WAVEFORMATEX,
*const GUID,
) -> HRESULT,
get_buffer_size: unsafe extern "system" fn(*mut c_void, *mut u32) -> HRESULT,
get_stream_latency: unsafe extern "system" fn(*mut c_void, *mut i64) -> HRESULT,
get_current_padding: unsafe extern "system" fn(*mut c_void, *mut u32) -> HRESULT,
is_format_supported: unsafe extern "system" fn(
*mut c_void,
i32,
*const WAVEFORMATEX,
*mut *mut WAVEFORMATEX,
) -> HRESULT,
get_mix_format: unsafe extern "system" fn(*mut c_void, *mut *mut WAVEFORMATEX) -> HRESULT,
get_device_period: unsafe extern "system" fn(*mut c_void, *mut i64, *mut i64) -> HRESULT,
start: unsafe extern "system" fn(*mut c_void) -> HRESULT,
stop: unsafe extern "system" fn(*mut c_void) -> HRESULT,
reset: unsafe extern "system" fn(*mut c_void) -> HRESULT,
set_event_handle: unsafe extern "system" fn(*mut c_void, HANDLE) -> HRESULT,
get_service: unsafe extern "system" fn(*mut c_void, *const GUID, *mut *mut c_void) -> HRESULT,
}
impl ComVtable for AudioClientVtable {
fn unknown(&self) -> &IUnknown_Vtbl {
&self.base
}
}
#[repr(C)]
struct AudioRenderClientVtable {
base: IUnknown_Vtbl,
get_buffer: unsafe extern "system" fn(*mut c_void, u32, *mut *mut u8) -> HRESULT,
release_buffer: unsafe extern "system" fn(*mut c_void, u32, u32) -> HRESULT,
}
impl ComVtable for AudioRenderClientVtable {
fn unknown(&self) -> &IUnknown_Vtbl {
&self.base
}
}
#[repr(C)]
struct AudioCaptureClientVtable {
base: IUnknown_Vtbl,
get_buffer: unsafe extern "system" fn(
*mut c_void,
*mut *mut u8,
*mut u32,
*mut u32,
*mut u64,
*mut u64,
) -> HRESULT,
release_buffer: unsafe extern "system" fn(*mut c_void, u32) -> HRESULT,
get_next_packet_size: unsafe extern "system" fn(*mut c_void, *mut u32) -> HRESULT,
}
impl ComVtable for AudioCaptureClientVtable {
fn unknown(&self) -> &IUnknown_Vtbl {
&self.base
}
}
struct ComPtr<V: ComVtable> {
ptr: *mut RawComInterface<V>,
}
impl<V: ComVtable> ComPtr<V> {
fn new(raw: *mut c_void, what: &str) -> VoiceResult<Self> {
if raw.is_null() {
Err(format!("{what} returned a null COM interface"))
} else {
Ok(Self { ptr: raw.cast() })
}
}
const fn as_void(&self) -> *mut c_void {
self.ptr.cast()
}
fn vtable(&self) -> &V {
// SAFETY: `ptr` was returned as a live COM interface, and each concrete
// interface pointer begins with a pointer to its corresponding vtable.
unsafe { &*(*self.ptr).vtable }
}
}
impl<V: ComVtable> Drop for ComPtr<V> {
fn drop(&mut self) {
let release = self.vtable().unknown().Release;
// SAFETY: this object owns one reference to the live COM interface.
unsafe { release(self.as_void()) };
}
}
#[derive(Clone, Copy)]
struct EventHandle(HANDLE);
// SAFETY: Windows event handles may be used from any thread, and reference
// counting keeps the owning handle open while either side can access it.
unsafe impl Send for EventHandle {}
// SAFETY: `SetEvent` and waits on a Windows event handle are thread-safe.
unsafe impl Sync for EventHandle {}
struct OwnedEvent(EventHandle);
impl OwnedEvent {
fn create() -> VoiceResult<Arc<Self>> {
// SAFETY: null security attributes and name request a private auto-reset,
// initially nonsignaled event.
let raw = unsafe { CreateEventW(null(), 0, 0, null()) };
if raw.is_null() {
Err("CreateEventW failed".to_owned())
} else {
Ok(Arc::new(Self(EventHandle(raw))))
}
}
const fn handle(&self) -> EventHandle {
self.0
}
fn signal(&self) {
// SAFETY: every caller owns an `Arc` that keeps this handle open.
let _ = unsafe { SetEvent(self.0.0) };
}
}
impl Drop for OwnedEvent {
fn drop(&mut self) {
// SAFETY: this is the sole owner and closes the valid event handle once.
unsafe { CloseHandle(self.0.0) };
}
}
struct ComApartment;
impl ComApartment {
fn initialize() -> VoiceResult<Self> {
// SAFETY: this dedicated worker has not initialized COM yet; the reserved
// pointer is required to be null.
let hr = unsafe { CoInitializeEx(null(), COINIT_MULTITHREADED as u32) };
check_hresult(hr, "CoInitializeEx")?;
Ok(Self)
}
}
impl Drop for ComApartment {
fn drop(&mut self) {
// SAFETY: paired with the successful `CoInitializeEx` on this same thread.
unsafe { CoUninitialize() };
}
}
struct BaseStream {
client: ComPtr<AudioClientVtable>,
device: ComPtr<MmDeviceVtable>,
enumerator: ComPtr<MmDeviceEnumeratorVtable>,
event: Arc<OwnedEvent>,
buffer_size: u32,
_apartment: ComApartment,
}
impl BaseStream {
fn open(
config: DeviceConfig,
data_flow: i32,
event: Option<Arc<OwnedEvent>>,
) -> VoiceResult<Self> {
let apartment = ComApartment::initialize()?;
let mut enumerator_raw = null_mut();
// SAFETY: all pointers are valid for the call, and `enumerator_raw` is an
// out parameter for the requested interface.
let hr = unsafe {
CoCreateInstance(
&CLSID_MMDEVICE_ENUMERATOR,
null_mut(),
CLSCTX_ALL,
&IID_IMMDEVICE_ENUMERATOR,
&mut enumerator_raw,
)
};
check_hresult(hr, "CoCreateInstance(MMDeviceEnumerator)")?;
let enumerator: ComPtr<MmDeviceEnumeratorVtable> =
ComPtr::new(enumerator_raw, "CoCreateInstance(MMDeviceEnumerator)")?;
let mut device_raw = null_mut();
// SAFETY: the enumerator is live and the output pointer is writable.
let hr = unsafe {
(enumerator.vtable().get_default_audio_endpoint)(
enumerator.as_void(),
data_flow,
eConsole,
&mut device_raw,
)
};
check_hresult(hr, "IMMDeviceEnumerator::GetDefaultAudioEndpoint")?;
let device: ComPtr<MmDeviceVtable> =
ComPtr::new(device_raw, "IMMDeviceEnumerator::GetDefaultAudioEndpoint")?;
let mut client_raw = null_mut();
// SAFETY: the device is live, activation parameters are optional and null,
// and `client_raw` receives the requested interface.
let hr = unsafe {
(device.vtable().activate)(
device.as_void(),
&IID_IAUDIO_CLIENT,
CLSCTX_ALL,
null(),
&mut client_raw,
)
};
check_hresult(hr, "IMMDevice::Activate(IAudioClient)")?;
let client: ComPtr<AudioClientVtable> =
ComPtr::new(client_raw, "IMMDevice::Activate(IAudioClient)")?;
let bytes_per_second = config
.sample_rate
.checked_mul(4)
.ok_or_else(|| "WASAPI sample rate is too large".to_owned())?;
let format = WAVEFORMATEX {
wFormatTag: WAVE_FORMAT_IEEE_FLOAT as u16,
nChannels: 1,
nSamplesPerSec: config.sample_rate,
nAvgBytesPerSec: bytes_per_second,
nBlockAlign: 4,
wBitsPerSample: 32,
cbSize: 0,
};
let stream_flags = AUDCLNT_STREAMFLAGS_EVENTCALLBACK
| AUDCLNT_STREAMFLAGS_AUTOCONVERTPCM
| AUDCLNT_STREAMFLAGS_SRC_DEFAULT_QUALITY;
let buffer_duration = i64::from(config.period_ms) * 3 * 10_000;
// SAFETY: the client is live and `format` remains valid for the duration
// of this synchronous initialization call.
let hr = unsafe {
(client.vtable().initialize)(
client.as_void(),
AUDCLNT_SHAREMODE_SHARED,
stream_flags,
buffer_duration,
0,
&format,
null(),
)
};
check_hresult(hr, "IAudioClient::Initialize")?;
let event = match event {
Some(event) => event,
None => OwnedEvent::create()?,
};
let raw_event = event.handle().0;
// SAFETY: the client and event are live for the remainder of the stream.
let hr = unsafe { (client.vtable().set_event_handle)(client.as_void(), raw_event) };
if let Err(error) = check_hresult(hr, "IAudioClient::SetEventHandle") {
drop(client);
drop(device);
drop(enumerator);
drop(event);
drop(apartment);
return Err(error);
}
let mut buffer_size = 0;
// SAFETY: the initialized client is live and the output pointer is valid.
let hr = unsafe { (client.vtable().get_buffer_size)(client.as_void(), &mut buffer_size) };
if let Err(error) = check_hresult(hr, "IAudioClient::GetBufferSize") {
drop(client);
drop(device);
drop(enumerator);
drop(event);
drop(apartment);
return Err(error);
}
if buffer_size == 0 {
drop(client);
drop(device);
drop(enumerator);
drop(event);
drop(apartment);
return Err("IAudioClient::GetBufferSize returned zero frames".to_owned());
}
Ok(Self { client, device, enumerator, event, buffer_size, _apartment: apartment })
}
fn event_handle(&self) -> EventHandle {
self.event.handle()
}
fn start(&self) -> VoiceResult<()> {
// SAFETY: the client is fully initialized and its event handle is set.
let hr = unsafe { (self.client.vtable().start)(self.client.as_void()) };
check_hresult(hr, "IAudioClient::Start")
}
fn stop(&self) {
// SAFETY: the client remains live and may be stopped during teardown.
let _ = unsafe { (self.client.vtable().stop)(self.client.as_void()) };
}
}
struct PlaybackStream {
render: ComPtr<AudioRenderClientVtable>,
base: BaseStream,
period_frames: u32,
started: bool,
}
impl PlaybackStream {
fn open(config: DeviceConfig, event: Option<Arc<OwnedEvent>>) -> VoiceResult<Self> {
let base = BaseStream::open(config, eRender, event)?;
let period_frames = u32::try_from(config.period_samples())
.map_err(|_| "WASAPI playback period is too large".to_owned())?;
if period_frames > base.buffer_size {
return Err(format!(
"WASAPI endpoint buffer ({} frames) is smaller than one playback period \
({period_frames} frames)",
base.buffer_size
));
}
let mut render_raw = null_mut();
// SAFETY: the initialized client is live and `render_raw` receives the
// requested service interface.
let hr = unsafe {
(base.client.vtable().get_service)(
base.client.as_void(),
&IID_IAUDIO_RENDER_CLIENT,
&mut render_raw,
)
};
check_hresult(hr, "IAudioClient::GetService(IAudioRenderClient)")?;
let render = ComPtr::new(render_raw, "IAudioClient::GetService(IAudioRenderClient)")?;
let mut stream = Self { render, base, period_frames, started: false };
stream.base.start()?;
stream.started = true;
Ok(stream)
}
}
impl Drop for PlaybackStream {
fn drop(&mut self) {
if self.started {
self.base.stop();
}
}
}
struct CaptureStream {
capture: ComPtr<AudioCaptureClientVtable>,
base: BaseStream,
started: bool,
}
impl CaptureStream {
fn open(config: DeviceConfig, event: Option<Arc<OwnedEvent>>) -> VoiceResult<Self> {
let base = BaseStream::open(config, eCapture, event)?;
let mut capture_raw = null_mut();
// SAFETY: the initialized client is live and `capture_raw` receives the
// requested service interface.
let hr = unsafe {
(base.client.vtable().get_service)(
base.client.as_void(),
&IID_IAUDIO_CAPTURE_CLIENT,
&mut capture_raw,
)
};
check_hresult(hr, "IAudioClient::GetService(IAudioCaptureClient)")?;
let capture = ComPtr::new(capture_raw, "IAudioClient::GetService(IAudioCaptureClient)")?;
let mut stream = Self { capture, base, started: false };
stream.base.start()?;
stream.started = true;
Ok(stream)
}
}
impl Drop for CaptureStream {
fn drop(&mut self) {
if self.started {
self.base.stop();
}
}
}
pub(crate) struct PlaybackDevice {
stop: Arc<AtomicBool>,
event: Arc<OwnedEvent>,
thread: Option<JoinHandle<VoiceResult<()>>>,
}
impl PlaybackDevice {
/// Open and start shared-mode playback on the default console endpoint.
pub fn start(config: DeviceConfig, fill: PlaybackFill) -> VoiceResult<Self> {
let stop = Arc::new(AtomicBool::new(false));
let worker_stop = Arc::clone(&stop);
let (startup_tx, startup_rx) = mpsc::channel();
let thread = thread::Builder::new()
.name("pi-voice-wasapi-playback".to_owned())
.spawn(move || playback_thread(config, fill, worker_stop, startup_tx))
.map_err(|error| format!("failed to spawn WASAPI playback thread: {error}"))?;
match startup_rx.recv() {
Ok(Ok(event)) => Ok(Self { stop, event, thread: Some(thread) }),
Ok(Err(error)) => {
let _ = thread.join();
Err(error)
},
Err(_) => match thread.join() {
Ok(Err(error)) => Err(error),
Ok(Ok(())) => Err("WASAPI playback thread exited during startup".to_owned()),
Err(_) => Err("WASAPI playback thread panicked during startup".to_owned()),
},
}
}
/// Stop playback and wait until its worker can no longer invoke `fill`.
pub fn stop(&mut self) -> VoiceResult<()> {
stop_worker(&self.stop, &self.event, &mut self.thread, "playback")
}
}
impl Drop for PlaybackDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
pub(crate) struct CaptureDevice {
stop: Arc<AtomicBool>,
event: Arc<OwnedEvent>,
thread: Option<JoinHandle<VoiceResult<()>>>,
}
impl CaptureDevice {
/// Open and start shared-mode capture on the default console endpoint.
pub fn start(config: DeviceConfig, sink: CaptureSink) -> VoiceResult<Self> {
let stop = Arc::new(AtomicBool::new(false));
let worker_stop = Arc::clone(&stop);
let (startup_tx, startup_rx) = mpsc::channel();
let thread = thread::Builder::new()
.name("pi-voice-wasapi-capture".to_owned())
.spawn(move || capture_thread(config, sink, worker_stop, startup_tx))
.map_err(|error| format!("failed to spawn WASAPI capture thread: {error}"))?;
match startup_rx.recv() {
Ok(Ok(event)) => Ok(Self { stop, event, thread: Some(thread) }),
Ok(Err(error)) => {
let _ = thread.join();
Err(error)
},
Err(_) => match thread.join() {
Ok(Err(error)) => Err(error),
Ok(Ok(())) => Err("WASAPI capture thread exited during startup".to_owned()),
Err(_) => Err("WASAPI capture thread panicked during startup".to_owned()),
},
}
}
/// Stop capture and wait until its worker can no longer invoke `sink`.
pub fn stop(&mut self) -> VoiceResult<()> {
stop_worker(&self.stop, &self.event, &mut self.thread, "capture")
}
}
impl Drop for CaptureDevice {
fn drop(&mut self) {
let _ = self.stop();
}
}
enum RunError {
DeviceInvalidated,
Other(String),
}
fn playback_thread(
config: DeviceConfig,
mut fill: PlaybackFill,
stop: Arc<AtomicBool>,
startup: Sender<VoiceResult<Arc<OwnedEvent>>>,
) -> VoiceResult<()> {
let mut stream = match PlaybackStream::open(config, None) {
Ok(stream) => stream,
Err(error) => {
let _ = startup.send(Err(error.clone()));
return Err(error);
},
};
let event = Arc::clone(&stream.base.event);
startup
.send(Ok(Arc::clone(&event)))
.map_err(|_| "WASAPI playback startup receiver was dropped".to_owned())?;
loop {
match run_playback(&stream, &stop, &mut fill) {
Ok(()) => return Ok(()),
Err(RunError::Other(error)) => return Err(error),
Err(RunError::DeviceInvalidated) => {
drop(stream);
let Some(reopened) = reopen_playback(config, &event, &stop)? else {
return Ok(());
};
stream = reopened;
},
}
}
}
fn run_playback(
stream: &PlaybackStream,
stop: &AtomicBool,
fill: &mut PlaybackFill,
) -> Result<(), RunError> {
loop {
if stop.load(Ordering::Acquire) {
return Ok(());
}
let event_signaled = wait_for_event(stream.base.event_handle()).map_err(RunError::Other)?;
if stop.load(Ordering::Acquire) {
return Ok(());
}
let mut padding = 0;
// SAFETY: the client is started and the output pointer is valid.
let hr = unsafe {
(stream.base.client.vtable().get_current_padding)(
stream.base.client.as_void(),
&mut padding,
)
};
check_run_hresult(hr, "IAudioClient::GetCurrentPadding")?;
if !event_signaled {
continue;
}
let max_padding = stream
.period_frames
.saturating_mul(2)
.min(stream.base.buffer_size);
if padding
.checked_add(stream.period_frames)
.is_none_or(|queued| queued > max_padding)
{
continue;
}
let frames = stream.period_frames;
let mut data = null_mut();
// SAFETY: the render client is live and `data` receives one logical
// period of writable frames.
let hr =
unsafe { (stream.render.vtable().get_buffer)(stream.render.as_void(), frames, &mut data) };
check_run_hresult(hr, "IAudioRenderClient::GetBuffer")?;
if data.is_null() {
// SAFETY: release balances the successful buffer acquisition.
let _ = unsafe {
(stream.render.vtable().release_buffer)(
stream.render.as_void(),
frames,
AUDCLNT_BUFFERFLAGS_SILENT as u32,
)
};
return Err(RunError::Other("IAudioRenderClient::GetBuffer returned null".to_owned()));
}
if stop.load(Ordering::Acquire) {
// SAFETY: release balances the successful acquisition; silent
// prevents uninitialized samples from being rendered.
let hr = unsafe {
(stream.render.vtable().release_buffer)(
stream.render.as_void(),
frames,
AUDCLNT_BUFFERFLAGS_SILENT as u32,
)
};
check_run_hresult(hr, "IAudioRenderClient::ReleaseBuffer")?;
return Ok(());
}
// SAFETY: WASAPI returned exactly `frames` writable mono IEEE-float
// samples for the format used to initialize this client.
let output = unsafe { slice::from_raw_parts_mut(data.cast::<f32>(), frames as usize) };
fill(output);
// SAFETY: release balances the successful acquisition after `fill`
// initialized the complete fixed-quantum buffer.
let hr =
unsafe { (stream.render.vtable().release_buffer)(stream.render.as_void(), frames, 0) };
check_run_hresult(hr, "IAudioRenderClient::ReleaseBuffer")?;
}
}
fn capture_thread(
config: DeviceConfig,
mut sink: CaptureSink,
stop: Arc<AtomicBool>,
startup: Sender<VoiceResult<Arc<OwnedEvent>>>,
) -> VoiceResult<()> {
let mut stream = match CaptureStream::open(config, None) {
Ok(stream) => stream,
Err(error) => {
let _ = startup.send(Err(error.clone()));
return Err(error);
},
};
let event = Arc::clone(&stream.base.event);
startup
.send(Ok(Arc::clone(&event)))
.map_err(|_| "WASAPI capture startup receiver was dropped".to_owned())?;
loop {
match run_capture(&stream, &stop, &mut sink) {
Ok(()) => return Ok(()),
Err(RunError::Other(error)) => return Err(error),
Err(RunError::DeviceInvalidated) => {
drop(stream);
let Some(reopened) = reopen_capture(config, &event, &stop)? else {
return Ok(());
};
stream = reopened;
},
}
}
}
fn run_capture(
stream: &CaptureStream,
stop: &AtomicBool,
sink: &mut CaptureSink,
) -> Result<(), RunError> {
let silent = vec![0.0_f32; stream.base.buffer_size as usize];
'events: loop {
if stop.load(Ordering::Acquire) {
return Ok(());
}
wait_for_event(stream.base.event_handle()).map_err(RunError::Other)?;
if stop.load(Ordering::Acquire) {
return Ok(());
}
loop {
if stop.load(Ordering::Acquire) {
break 'events;
}
let mut packet_size = 0;
// SAFETY: the capture service is live and the output pointer is valid.
let hr = unsafe {
(stream.capture.vtable().get_next_packet_size)(
stream.capture.as_void(),
&mut packet_size,
)
};
check_run_hresult(hr, "IAudioCaptureClient::GetNextPacketSize")?;
if packet_size == 0 {
break;
}
let mut data = null_mut();
let mut frames = 0;
let mut flags = 0;
// SAFETY: all outputs are valid and device/QPC positions are optional.
let hr = unsafe {
(stream.capture.vtable().get_buffer)(
stream.capture.as_void(),
&mut data,
&mut frames,
&mut flags,
null_mut(),
null_mut(),
)
};
check_run_hresult(hr, "IAudioCaptureClient::GetBuffer")?;
if stop.load(Ordering::Acquire) {
// SAFETY: release balances the successful buffer acquisition.
let hr = unsafe {
(stream.capture.vtable().release_buffer)(stream.capture.as_void(), frames)
};
check_run_hresult(hr, "IAudioCaptureClient::ReleaseBuffer")?;
break 'events;
}
if frames != 0 {
if flags & AUDCLNT_BUFFERFLAGS_SILENT as u32 != 0 {
let Some(samples) = silent.get(..frames as usize) else {
// SAFETY: release balances the successful buffer acquisition.
let _ = unsafe {
(stream.capture.vtable().release_buffer)(stream.capture.as_void(), frames)
};
return Err(RunError::Other(
"WASAPI capture packet exceeds the endpoint buffer".to_owned(),
));
};
sink(samples);
} else {
if data.is_null() {
// SAFETY: release balances the successful buffer acquisition.
let _ = unsafe {
(stream.capture.vtable().release_buffer)(stream.capture.as_void(), frames)
};
return Err(RunError::Other(
"IAudioCaptureClient::GetBuffer returned null".to_owned(),
));
}
// SAFETY: WASAPI returned `frames` readable mono IEEE-float samples
// for the format used to initialize this client.
let samples = unsafe { slice::from_raw_parts(data.cast::<f32>(), frames as usize) };
sink(samples);
}
}
// SAFETY: release balances the successful buffer acquisition.
let hr =
unsafe { (stream.capture.vtable().release_buffer)(stream.capture.as_void(), frames) };
check_run_hresult(hr, "IAudioCaptureClient::ReleaseBuffer")?;
}
}
Ok(())
}
// We deliberately omit `IMMNotificationClient`: a live endpoint stays selected
// across default-device changes. Device invalidation is the unambiguous point
// at which these retries reopen whichever endpoint is currently the default.
fn reopen_playback(
config: DeviceConfig,
event: &Arc<OwnedEvent>,
stop: &AtomicBool,
) -> VoiceResult<Option<PlaybackStream>> {
let mut last_error = "default endpoint remained unavailable".to_owned();
for attempt in 0..REOPEN_ATTEMPTS {
if stop.load(Ordering::Acquire) {
return Ok(None);
}
if attempt != 0 {
thread::sleep(REOPEN_BACKOFF);
if stop.load(Ordering::Acquire) {
return Ok(None);
}
}
match PlaybackStream::open(config, Some(Arc::clone(event))) {
Ok(stream) => return Ok(Some(stream)),
Err(error) => last_error = error,
}
}
Err(format!(
"WASAPI playback endpoint recovery failed after {REOPEN_ATTEMPTS} attempts: {last_error}"
))
}
fn reopen_capture(
config: DeviceConfig,
event: &Arc<OwnedEvent>,
stop: &AtomicBool,
) -> VoiceResult<Option<CaptureStream>> {
let mut last_error = "default endpoint remained unavailable".to_owned();
for attempt in 0..REOPEN_ATTEMPTS {
if stop.load(Ordering::Acquire) {
return Ok(None);
}
if attempt != 0 {
thread::sleep(REOPEN_BACKOFF);
if stop.load(Ordering::Acquire) {
return Ok(None);
}
}
match CaptureStream::open(config, Some(Arc::clone(event))) {
Ok(stream) => return Ok(Some(stream)),
Err(error) => last_error = error,
}
}
Err(format!(
"WASAPI capture endpoint recovery failed after {REOPEN_ATTEMPTS} attempts: {last_error}"
))
}
fn wait_for_event(event: EventHandle) -> VoiceResult<bool> {
// SAFETY: a reference-counted owner keeps this handle live throughout waits.
let result = unsafe { WaitForSingleObject(event.0, WAIT_TIMEOUT_MS) };
match result {
WAIT_OBJECT_0 => Ok(true),
WAIT_TIMEOUT => Ok(false),
code => Err(format!("WaitForSingleObject failed (result 0x{code:08X})")),
}
}
fn stop_worker(
stop: &AtomicBool,
event: &OwnedEvent,
thread: &mut Option<JoinHandle<VoiceResult<()>>>,
direction: &str,
) -> VoiceResult<()> {
stop.store(true, Ordering::Release);
event.signal();
if thread
.as_ref()
.is_some_and(|worker| worker.thread().id() == thread::current().id())
{
// Callback-thread stop cannot wait for its own frame to unwind. The
// callback-thread carve-out accepts that; the flag prevents another
// delivery, and retaining the handle still permits a later external join.
return Ok(());
}
let Some(thread) = thread.take() else {
return Ok(());
};
match thread.join() {
Ok(result) => result,
Err(_) => Err(format!("WASAPI {direction} thread panicked")),
}
}
fn check_hresult(hr: HRESULT, what: &str) -> VoiceResult<()> {
if hr < 0 {
Err(format!("{what} failed (HRESULT 0x{:08X})", hr as u32))
} else {
Ok(())
}
}
fn check_run_hresult(hr: HRESULT, what: &str) -> Result<(), RunError> {
if hr == AUDCLNT_E_DEVICE_INVALIDATED {
Err(RunError::DeviceInvalidated)
} else {
check_hresult(hr, what).map_err(RunError::Other)
}
}
+5 -4
View File
@@ -2,14 +2,14 @@
//! live-conversation peer.
//!
//! napi-free by design. `pi-natives` wraps these types in thin `#[napi]`
//! adapters, so the webrtc/opus/miniaudio dependency graph compiles once into
//! this rlib and rebuilds only when voice code changes — never when the
//! N-API surface elsewhere in the addon does.
//! adapters, so the webrtc/opus dependency graph compiles once into this
//! rlib and rebuilds only when voice code changes — never when the N-API
//! surface elsewhere in the addon does.
//!
//! # Architecture
//! ```text
//! JS (packages/natives) -> #[napi] adapters (pi-natives audio.rs / live.rs)
//! -> pi_voice::audio (miniaudio capture/playback engine)
//! -> pi_voice::audio (capture/playback engine over in-house OS backends)
//! -> pi_voice::live (WebRTC peer + Opus media, feeds audio playback)
//! ```
//!
@@ -19,6 +19,7 @@
//! worker threads and must not block.
pub mod audio;
pub(crate) mod device;
pub mod live;
/// Engine-level result: a human-readable failure message the N-API layer maps
+4
View File
@@ -6,6 +6,10 @@
- Added support for Windows hosts in `bun run build`, enabling local N-API builds against VS Build Tools without requiring a pre-configured vcvars prompt.
### Changed
- Replaced the miniaudio (`maudio`) dependency with in-house platform audio backends for `AudioCapture`/`AudioPlayback`: CoreAudio AudioQueue on macOS, shared-mode WASAPI on Windows, and PulseAudio (ALSA fallback) loaded via `dlopen` on Linux. Removes the bindgen/libclang requirement and the Windows rustc-ICE workaround from the native build.
### Fixed
- Fixed CPU feature detection (AVX2) on Windows hosts, resolving an issue where the native addon loader and local builds incorrectly fell back to the baseline variant, while improving startup performance by ~270ms.
+4 -2
View File
@@ -702,8 +702,10 @@ export interface DesktopSessionOptions {
/** One capturable top-level window in global logical desktop coordinates. */
export interface DesktopWindow {
/**
* Stable numeric window id, valid as a capture target while the window
* lives.
* Backend-defined opaque window id, valid as a capture target while the
* window lives. Numeric on X11/Win32/macOS; a composite AT-SPI string on
* Wayland (e.g. `atspi::1.31:/org/a11y/atspi/accessible/1`). Never parse
* it.
*/
id: string
/** Window title; may be empty for untitled windows. */