Windows port Phase 1: real cpal/WASAPI audio backend
Replace the Phase 0 no-op CpalBackend stub with a working cpal backend (WASAPI on Windows), preserving the exact PipeWire AudioBackend contract so the mixer/encoder/jitter pipeline is unchanged. - Capture: input stream -> downmix to mono -> 960-sample (20ms) i16 frames -> tx, matching the encoder/jitter frame size. - Playback: 200ms stereo ring prefilled to PLAYBACK_TARGET_SAMPLES; the output callback drains it (silence on underrun) while the owning thread feeds it from rx. ring_fill is the exact delta-maintained occupancy counter (fetch_add on push, fetch_sub on pop), preserving the clock-paced production design (not ringbuf's stale occupied_len). - cpal::Stream is !Send, but AudioBackend is Send+Sync and shared via Arc, so each stream lives on its own owning thread (built/played/dropped there); the struct holds only the running flag + JoinHandle. stop() flips the flag and joins. - Generic over F32/I16/U16 sample formats; device selected by name else default; requires a native 48kHz config (clear error otherwise, no resampling yet). Mirrors the PipeWire drain_loop and playout-health line. - Cargo.toml: add cpal 0.15 under cfg(windows). Verified by temporarily compiling cpal_impl against real cpal on Linux/ALSA: build + clippy clean, 6/6 cpal_impl unit tests pass. Reverted to windows-only gating; shipped Linux state green (316/316). Runtime/WASAPI end-to-end is unverified and pending a Windows host (plan M2). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+567
-32
@@ -1,31 +1,80 @@
|
||||
//! Windows audio backend (cpal/WASAPI) — **Phase 0 stub**.
|
||||
//! Windows audio backend — cpal / WASAPI (Phase 1).
|
||||
//!
|
||||
//! This is a compile-and-run placeholder so the Windows build links and the app
|
||||
//! starts up (networking, UI, and text chat all functional) while the real
|
||||
//! capture/playback implementation lands in Phase 1. Every method satisfies the
|
||||
//! [`AudioBackend`] contract as a no-op: no microphone is captured and nothing is
|
||||
//! played. It deliberately pulls in no extra dependency — `cpal` is added only
|
||||
//! when the real implementation arrives.
|
||||
//! 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:
|
||||
//!
|
||||
//! Phase 1 will replace this with cpal streams on the WASAPI host, mapping:
|
||||
//! - `start_capture` → input stream, f32→i16, mono 48 kHz, into `tx`;
|
||||
//! - `start_playback` → output stream draining a `ringbuf`, keeping `ring_fill`
|
||||
//! updated so the existing hardware-clock pacing in the mixer keeps working;
|
||||
//! - `stop` → drop the streams.
|
||||
//! - **Capture**: mono, 48 kHz, S16 PCM, emitted as `Vec<i16>` frames of
|
||||
//! [`CAPTURE_FRAME`] (960 = 20 ms) samples — matching the encoder/jitter frame.
|
||||
//! - **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 from the network mixer.
|
||||
//!
|
||||
//! ## Sample rate
|
||||
//!
|
||||
//! The whole pipeline assumes 48 kHz (Opus + the 960-sample frame). Phase 1 only
|
||||
//! selects a native-48 kHz device config; if the device can't do 48 kHz we return
|
||||
//! a clear error rather than silently producing pitch-shifted audio. Arbitrary
|
||||
//! sample-rate support (resampling) is a Phase 1.1 follow-up.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::AtomicUsize;
|
||||
use std::sync::mpsc::{Receiver, Sender};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
|
||||
use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::thread::{self, JoinHandle};
|
||||
use std::time::Duration;
|
||||
|
||||
use super::{AudioBackend, AudioError};
|
||||
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
|
||||
use cpal::{Device, FromSample, Sample, SampleFormat, SampleRate, SizedSample, Stream, StreamConfig};
|
||||
use ringbuf::{
|
||||
traits::{Consumer, Producer, Split},
|
||||
HeapRb,
|
||||
};
|
||||
|
||||
/// No-op Windows audio backend (Phase 0). See module docs.
|
||||
pub struct CpalBackend;
|
||||
use super::{AudioBackend, 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;
|
||||
/// 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);
|
||||
|
||||
/// Windows audio backend. See module docs.
|
||||
pub struct CpalBackend {
|
||||
capture: Mutex<Option<StreamWorker>>,
|
||||
playback: Mutex<Option<StreamWorker>>,
|
||||
}
|
||||
|
||||
/// A spawned owning thread plus the flag that tells it to drop its stream and exit.
|
||||
struct StreamWorker {
|
||||
running: Arc<AtomicBool>,
|
||||
thread: JoinHandle<()>,
|
||||
}
|
||||
|
||||
impl CpalBackend {
|
||||
pub fn new() -> Self {
|
||||
crate::log_msg("CpalBackend: Phase 0 stub active (no audio I/O yet)");
|
||||
CpalBackend
|
||||
Self {
|
||||
capture: Mutex::new(None),
|
||||
playback: Mutex::new(None),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -38,30 +87,516 @@ impl Default for CpalBackend {
|
||||
impl AudioBackend for CpalBackend {
|
||||
fn start_capture(
|
||||
&self,
|
||||
_tx: Sender<Vec<i16>>,
|
||||
_target_node: Option<String>,
|
||||
tx: Sender<Vec<i16>>,
|
||||
target_node: Option<String>,
|
||||
) -> Result<(), AudioError> {
|
||||
// No capture stream yet: dropping `_tx` simply means no samples are ever
|
||||
// produced (silent mic), which is the intended Phase 0 behaviour.
|
||||
crate::log_msg("CpalBackend::start_capture: not yet implemented (Phase 1) — capturing silence");
|
||||
let mut guard = self.capture.lock().unwrap();
|
||||
if guard.is_some() {
|
||||
return Err(AudioError::Stream("Capture already started".to_string()));
|
||||
}
|
||||
let running = Arc::new(AtomicBool::new(true));
|
||||
let running_thread = running.clone();
|
||||
let thread = thread::Builder::new()
|
||||
.name("peerspeak-cpal-capture".to_string())
|
||||
.spawn(move || {
|
||||
if let Err(e) = run_capture(tx, target_node, running_thread) {
|
||||
crate::log_msg(&format!("cpal capture error: {e}"));
|
||||
}
|
||||
})
|
||||
.map_err(|e| AudioError::Init(e.to_string()))?;
|
||||
*guard = Some(StreamWorker { running, thread });
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn start_playback(
|
||||
&self,
|
||||
rx: Receiver<Vec<i16>>,
|
||||
_target_node: Option<String>,
|
||||
_ring_fill: Arc<AtomicUsize>,
|
||||
target_node: Option<String>,
|
||||
ring_fill: Arc<AtomicUsize>,
|
||||
) -> Result<(), AudioError> {
|
||||
// Drain and discard incoming audio on a detached thread so the mixer's
|
||||
// producer never blocks or sees a closed channel. This keeps the rest of
|
||||
// the pipeline running normally while output is silent.
|
||||
std::thread::spawn(move || while rx.recv().is_ok() {});
|
||||
crate::log_msg("CpalBackend::start_playback: not yet implemented (Phase 1) — discarding output");
|
||||
let mut guard = self.playback.lock().unwrap();
|
||||
if guard.is_some() {
|
||||
return Err(AudioError::Stream("Playback already started".to_string()));
|
||||
}
|
||||
let running = Arc::new(AtomicBool::new(true));
|
||||
let running_thread = running.clone();
|
||||
let thread = thread::Builder::new()
|
||||
.name("peerspeak-cpal-playback".to_string())
|
||||
.spawn(move || {
|
||||
if let Err(e) = run_playback(rx, target_node, ring_fill, running_thread) {
|
||||
crate::log_msg(&format!("cpal playback error: {e}"));
|
||||
}
|
||||
})
|
||||
.map_err(|e| AudioError::Init(e.to_string()))?;
|
||||
*guard = Some(StreamWorker { running, thread });
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn stop(&self) -> Result<(), AudioError> {
|
||||
for slot in [&self.capture, &self.playback] {
|
||||
if let Some(worker) = slot.lock().unwrap().take() {
|
||||
worker.running.store(false, Ordering::Relaxed);
|
||||
let _ = worker.thread.join();
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Device / config selection
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Resolve a device (by `target` name, else the system default) and a stream
|
||||
/// config running natively at [`SAMPLE_RATE`].
|
||||
///
|
||||
/// For output we require [`PLAYBACK_CHANNELS`] (stereo) so the interleaved ring
|
||||
/// maps 1:1 to the device buffer; for input we prefer mono but accept any channel
|
||||
/// count and downmix. A device with no 48 kHz config is a hard error (no
|
||||
/// resampling yet — see module docs).
|
||||
fn resolve(
|
||||
output: bool,
|
||||
target: Option<String>,
|
||||
) -> 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 {
|
||||
Some(name) => find_device_by_name(&host, output, &name).or_else(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();
|
||||
Ok((device, config, sample_format))
|
||||
}
|
||||
|
||||
fn find_device_by_name(host: &cpal::Host, output: bool, name: &str) -> Option<Device> {
|
||||
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))
|
||||
}
|
||||
|
||||
/// Pick a supported config at exactly [`SAMPLE_RATE`]. Output must be stereo;
|
||||
/// input prefers mono, then any channel count (downmixed later).
|
||||
fn choose_config(
|
||||
device: &Device,
|
||||
output: bool,
|
||||
) -> Result<cpal::SupportedStreamConfig, AudioError> {
|
||||
let ranges: Vec<cpal::SupportedStreamConfigRange> = 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<u16>| {
|
||||
ranges
|
||||
.iter()
|
||||
.find(|r| supports_48k(r) && channels.is_none_or(|c| r.channels() == c))
|
||||
.cloned()
|
||||
};
|
||||
|
||||
let chosen = if output {
|
||||
pick(Some(PLAYBACK_CHANNELS as u16))
|
||||
} else {
|
||||
pick(Some(1)).or_else(|| pick(None))
|
||||
};
|
||||
|
||||
chosen
|
||||
.map(|r| r.with_sample_rate(SampleRate(SAMPLE_RATE)))
|
||||
.ok_or_else(|| {
|
||||
AudioError::Device(format!(
|
||||
"device '{}' has no {SAMPLE_RATE} Hz {} config; resampling not yet implemented (Phase 1.1)",
|
||||
device.name().unwrap_or_else(|_| "<unknown>".to_string()),
|
||||
if output { "stereo output" } else { "input" },
|
||||
))
|
||||
})
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Capture
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn run_capture(
|
||||
tx: Sender<Vec<i16>>,
|
||||
target: Option<String>,
|
||||
running: Arc<AtomicBool>,
|
||||
) -> Result<(), AudioError> {
|
||||
let (device, config, sample_format) = resolve(false, target)?;
|
||||
let channels = config.channels as usize;
|
||||
|
||||
let stream = match sample_format {
|
||||
SampleFormat::F32 => build_input::<f32>(&device, &config, tx, channels),
|
||||
SampleFormat::I16 => build_input::<i16>(&device, &config, tx, channels),
|
||||
SampleFormat::U16 => build_input::<u16>(&device, &config, tx, channels),
|
||||
other => Err(AudioError::Stream(format!(
|
||||
"unsupported capture sample format: {other:?}"
|
||||
))),
|
||||
}?;
|
||||
|
||||
stream.play().map_err(|e| AudioError::Stream(e.to_string()))?;
|
||||
|
||||
// The RT callback does the work; this thread just keeps `stream` alive until
|
||||
// `stop()` flips the flag, at which point the stream is dropped (= stopped).
|
||||
while running.load(Ordering::Relaxed) {
|
||||
thread::sleep(WORKER_POLL);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn build_input<T>(
|
||||
device: &Device,
|
||||
config: &StreamConfig,
|
||||
tx: Sender<Vec<i16>>,
|
||||
channels: usize,
|
||||
) -> Result<Stream, AudioError>
|
||||
where
|
||||
T: SizedSample + Send + 'static,
|
||||
i16: FromSample<T>,
|
||||
{
|
||||
let mut acc = FrameAccumulator::new(CAPTURE_FRAME);
|
||||
let err_fn = |e| crate::log_msg(&format!("cpal capture stream error: {e}"));
|
||||
device
|
||||
.build_input_stream::<T, _, _>(
|
||||
config,
|
||||
move |data: &[T], _| {
|
||||
for frame in data.chunks_exact(channels) {
|
||||
let mono = downmix_to_mono(frame);
|
||||
if let Some(full) = acc.push(mono) {
|
||||
// Consumer gone (call ended) → stop feeding; the owning
|
||||
// thread will drop the stream on `stop()`.
|
||||
if tx.send(full).is_err() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
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<T>(frame: &[T]) -> i16
|
||||
where
|
||||
T: Copy,
|
||||
i16: FromSample<T>,
|
||||
{
|
||||
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
|
||||
}
|
||||
|
||||
/// 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<i16>,
|
||||
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<Vec<i16>> {
|
||||
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<Vec<i16>>,
|
||||
target: Option<String>,
|
||||
ring_fill: Arc<AtomicUsize>,
|
||||
running: Arc<AtomicBool>,
|
||||
) -> Result<(), AudioError> {
|
||||
let (device, config, sample_format) = resolve(true, target)?;
|
||||
|
||||
let rb = HeapRb::<i16>::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));
|
||||
|
||||
let stream = match sample_format {
|
||||
SampleFormat::F32 => {
|
||||
build_output::<f32, _>(&device, &config, consumer, ring_fill.clone(), underrun.clone())
|
||||
}
|
||||
SampleFormat::I16 => {
|
||||
build_output::<i16, _>(&device, &config, consumer, ring_fill.clone(), underrun.clone())
|
||||
}
|
||||
SampleFormat::U16 => {
|
||||
build_output::<u16, _>(&device, &config, consumer, ring_fill.clone(), underrun.clone())
|
||||
}
|
||||
other => Err(AudioError::Stream(format!(
|
||||
"unsupported playback sample format: {other:?}"
|
||||
))),
|
||||
}?;
|
||||
|
||||
stream.play().map_err(|e| AudioError::Stream(e.to_string()))?;
|
||||
|
||||
let logger = spawn_health_logger(
|
||||
running.clone(),
|
||||
ring_fill.clone(),
|
||||
underrun.clone(),
|
||||
dropped.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;
|
||||
}
|
||||
for &sample in &frame {
|
||||
let _ = producer.try_push(sample);
|
||||
}
|
||||
ring_fill.fetch_add(frame.len(), 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);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn build_output<T, C>(
|
||||
device: &Device,
|
||||
config: &StreamConfig,
|
||||
mut consumer: C,
|
||||
ring_fill: Arc<AtomicUsize>,
|
||||
underrun: Arc<AtomicU64>,
|
||||
) -> Result<Stream, AudioError>
|
||||
where
|
||||
T: SizedSample + FromSample<i16> + Send + 'static,
|
||||
C: Consumer<Item = i16> + Send + 'static,
|
||||
{
|
||||
let err_fn = |e| crate::log_msg(&format!("cpal playback stream error: {e}"));
|
||||
device
|
||||
.build_output_stream::<T, _, _>(
|
||||
config,
|
||||
move |data: &mut [T], _| {
|
||||
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()))
|
||||
}
|
||||
|
||||
/// 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<T, C>(consumer: &mut C, out: &mut [T]) -> (usize, u64)
|
||||
where
|
||||
T: Sample + FromSample<i16>,
|
||||
C: Consumer<Item = i16>,
|
||||
{
|
||||
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)
|
||||
}
|
||||
|
||||
/// 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<AtomicBool>,
|
||||
ring_fill: Arc<AtomicUsize>,
|
||||
underrun: Arc<AtomicU64>,
|
||||
dropped: Arc<AtomicU64>,
|
||||
) -> JoinHandle<()> {
|
||||
let verbose = std::env::var_os("PEERSPEAK_AUDIO_VERBOSE").is_some();
|
||||
thread::spawn(move || {
|
||||
let (mut last_u, mut last_d) = (0u64, 0u64);
|
||||
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),
|
||||
));
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// 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<Vec<i16>>, running: &AtomicBool, mut on_frame: impl FnMut(Vec<i16>)) {
|
||||
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::<i16>(&[100, 100]), 100);
|
||||
assert_eq!(downmix_to_mono::<i16>(&[100, -100]), 0);
|
||||
assert_eq!(downmix_to_mono::<i16>(&[50]), 50);
|
||||
assert_eq!(downmix_to_mono::<i16>(&[]), 0);
|
||||
// 4-channel average rounds toward zero (integer division).
|
||||
assert_eq!(downmix_to_mono::<i16>(&[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::<i16>::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 drain_loop_exits_when_running_flips_even_with_sender_alive() {
|
||||
let (tx, rx) = mpsc::channel::<Vec<i16>>();
|
||||
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::<Vec<i16>>();
|
||||
let running = Arc::new(AtomicBool::new(true));
|
||||
drop(tx);
|
||||
drain_loop(&rx, &running, |_| panic!("no frame should arrive"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn drain_loop_delivers_frames() {
|
||||
let (tx, rx) = mpsc::channel::<Vec<i16>>();
|
||||
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]]);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user