use std::collections::VecDeque;
use std::panic::AssertUnwindSafe;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};
use std::time::Duration;
use flexaudio_core::backend::{CaptureBackend, RawSink};
use flexaudio_core::chunk_ring::{chunk_ring, ChunkConsumer, ChunkProducer};
use flexaudio_core::clock::{monotonic_now_ns, ClockNormalizer};
use flexaudio_core::normalizer::{InnerProcessor, Normalizer};
use flexaudio_core::raw_ring::{raw_ring, RawConsumer};
use flexaudio_core::secondary_ring::{
secondary_chunk_ring, SecondaryChunkConsumer, SecondaryChunkProducer,
};
use flexaudio_core::types::{
AudioChunk, ChunkFlags, Error, Event, OutputFormat, Result, SecondaryChunk, StreamConfig,
};
struct DenoiseInnerProcessor {
denoiser: flexaudio_denoise::Denoiser,
}
impl DenoiseInnerProcessor {
fn new() -> Self {
Self {
denoiser: flexaudio_denoise::Denoiser::new(2)
.expect("stereo denoiser construction is infallible"),
}
}
}
impl InnerProcessor for DenoiseInnerProcessor {
fn process(&mut self, samples: &mut [f32]) {
let _ = self.denoiser.process(samples);
}
fn flush(&mut self) -> Vec<f32> {
self.denoiser.flush()
}
}
fn build_inner_processor(shared: &SharedState) -> Option<Box<dyn InnerProcessor>> {
if shared.denoise_enabled.load(Ordering::SeqCst) {
Some(Box::new(DenoiseInnerProcessor::new()))
} else {
None
}
}
fn build_normalizer(
shared: &SharedState,
rate: u32,
channels: u16,
output: OutputFormat,
secondary_output: Option<OutputFormat>,
) -> Result<Normalizer> {
let mut n = Normalizer::new(rate, channels, output)?;
if let Some(sec) = secondary_output {
n = n.with_secondary(sec)?;
}
if let Some(proc) = build_inner_processor(shared) {
n = n.with_inner_processor(proc);
}
Ok(n)
}
pub(crate) const RAW_RING_SAMPLES: usize = 48_000;
const WATCHDOG_TICK: Duration = Duration::from_millis(250);
const STALL_THRESHOLD: Duration = Duration::from_secs(2);
const BACKOFF_MIN: Duration = Duration::from_millis(250);
const BACKOFF_MAX: Duration = Duration::from_secs(5);
pub struct Stream {
config: StreamConfig,
shared: Arc<SharedState>,
chunk_consumer: ChunkConsumer,
secondary_consumer: Option<SecondaryChunkConsumer>,
events: Arc<Mutex<VecDeque<Event>>>,
worker: Option<JoinHandle<()>>,
watchdog: Option<JoinHandle<()>>,
started: bool,
}
struct SharedState {
backend: Mutex<Box<dyn CaptureBackend>>,
raw_consumer: Mutex<Option<RawConsumer>>,
raw_generation: AtomicU64,
last_sample_ns: AtomicI64,
stopping: AtomicBool,
recovered_pending: AtomicBool,
events: Arc<Mutex<VecDeque<Event>>>,
chunk_producer: Mutex<Option<ChunkProducer>>,
native_format: Mutex<(u32, u16)>,
switching: AtomicBool,
discontinuity_pending: AtomicBool,
paused: AtomicBool,
delivery: Mutex<()>,
resume_generation: AtomicU64,
gain_bits: AtomicU32,
recording_epoch_ns: AtomicI64,
denoise_enabled: AtomicBool,
secondary_producer: Mutex<Option<SecondaryChunkProducer>>,
}
impl SharedState {
fn push_event(&self, ev: Event) {
let mut q = self.events.lock().unwrap_or_else(|e| e.into_inner());
q.push_back(ev);
}
}
fn start_backend_catching(be: &mut Box<dyn CaptureBackend>, sink: RawSink) -> Result<()> {
match std::panic::catch_unwind(AssertUnwindSafe(|| be.start(sink))) {
Ok(res) => res,
Err(_) => Err(Error::Backend("backend panicked during start()".into())),
}
}
#[must_use]
fn stop_backend_catching(be: &mut Box<dyn CaptureBackend>) -> bool {
std::panic::catch_unwind(AssertUnwindSafe(|| be.stop())).is_ok()
}
impl Stream {
pub fn open(config: StreamConfig, backend: Box<dyn CaptureBackend>) -> Result<Stream> {
if config.ring_capacity_chunks == 0 {
return Err(Error::InvalidArg("ring_capacity_chunks must be > 0".into()));
}
if !config.gain.is_finite() || config.gain < 0.0 {
return Err(Error::InvalidArg(format!(
"gain must be finite and >= 0.0, got {}",
config.gain
)));
}
config.output.validate()?;
crate::validate_exclude_pids(&config)?;
if let Some(sec) = config.secondary_output {
sec.validate()?;
}
let native_format = backend.native_format();
if native_format.0 == 0 || native_format.1 == 0 {
return Err(Error::InvalidArg(
"backend native_format must have non-zero rate and channels".into(),
));
}
let (chunk_producer, chunk_consumer) = chunk_ring(config.ring_capacity_chunks);
let (secondary_producer, secondary_consumer) = if config.secondary_output.is_some() {
let (p, c) = secondary_chunk_ring(config.ring_capacity_chunks);
(Some(p), Some(c))
} else {
(None, None)
};
let events = Arc::new(Mutex::new(VecDeque::new()));
let shared = Arc::new(SharedState {
backend: Mutex::new(backend),
raw_consumer: Mutex::new(None),
raw_generation: AtomicU64::new(0),
last_sample_ns: AtomicI64::new(0),
stopping: AtomicBool::new(false),
recovered_pending: AtomicBool::new(false),
events: events.clone(),
chunk_producer: Mutex::new(Some(chunk_producer)),
native_format: Mutex::new(native_format),
switching: AtomicBool::new(false),
discontinuity_pending: AtomicBool::new(false),
paused: AtomicBool::new(false),
delivery: Mutex::new(()),
resume_generation: AtomicU64::new(0),
gain_bits: AtomicU32::new(config.gain.to_bits()),
recording_epoch_ns: AtomicI64::new(i64::MIN),
denoise_enabled: AtomicBool::new(false),
secondary_producer: Mutex::new(secondary_producer),
});
Ok(Stream {
config,
shared,
chunk_consumer,
secondary_consumer,
events,
worker: None,
watchdog: None,
started: false,
})
}
pub fn start(&mut self) -> Result<()> {
if self.started {
return Ok(());
}
self.shared.stopping.store(false, Ordering::SeqCst);
self.shared.paused.store(false, Ordering::SeqCst);
self.shared
.recording_epoch_ns
.store(i64::MIN, Ordering::SeqCst);
Self::open_backend_once(&self.shared)?;
let chunk_producer = self
.shared
.chunk_producer
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()
.ok_or_else(|| Error::InvalidState("chunk producer already taken".into()))?;
let secondary_producer = self
.shared
.secondary_producer
.lock()
.unwrap_or_else(|e| e.into_inner())
.take();
let worker_shared = self.shared.clone();
let initial_native = *self
.shared
.native_format
.lock()
.unwrap_or_else(|e| e.into_inner());
let output = self.config.output;
let secondary_output = self.config.secondary_output;
let worker = thread::Builder::new()
.name("flexaudio-intake".into())
.spawn(move || {
run_intake(
worker_shared,
chunk_producer,
secondary_producer,
initial_native,
output,
secondary_output,
);
})
.map_err(|e| Error::Backend(format!("spawn intake thread: {e}")))?;
self.worker = Some(worker);
let wd_shared = self.shared.clone();
let watchdog = thread::Builder::new()
.name("flexaudio-watchdog".into())
.spawn(move || {
run_watchdog(wd_shared);
})
.map_err(|e| Error::Backend(format!("spawn watchdog thread: {e}")))?;
self.watchdog = Some(watchdog);
self.started = true;
Ok(())
}
pub fn stop(&mut self) {
self.shared.stopping.store(true, Ordering::SeqCst);
{
let mut be = self
.shared
.backend
.lock()
.unwrap_or_else(|e| e.into_inner());
let _ = stop_backend_catching(&mut be);
}
if let Some(h) = self.worker.take() {
let _ = h.join();
}
if let Some(h) = self.watchdog.take() {
let _ = h.join();
}
self.started = false;
}
pub fn pause(&self) {
let _g = self
.shared
.delivery
.lock()
.unwrap_or_else(|e| e.into_inner());
self.shared.paused.store(true, Ordering::SeqCst);
}
pub fn resume(&self) {
let _g = self
.shared
.delivery
.lock()
.unwrap_or_else(|e| e.into_inner());
if self.shared.paused.load(Ordering::SeqCst) {
self.shared.resume_generation.fetch_add(1, Ordering::SeqCst);
self.shared.paused.store(false, Ordering::SeqCst);
}
}
pub fn is_paused(&self) -> bool {
self.shared.paused.load(Ordering::SeqCst)
}
pub fn set_gain(&self, gain: f32) -> Result<()> {
if !gain.is_finite() || gain < 0.0 {
return Err(Error::InvalidArg(format!(
"gain must be finite and >= 0.0, got {gain}"
)));
}
self.shared
.gain_bits
.store(gain.to_bits(), Ordering::Relaxed);
Ok(())
}
pub fn gain(&self) -> f32 {
f32::from_bits(self.shared.gain_bits.load(Ordering::Relaxed))
}
pub fn poll_chunk(&mut self) -> Option<AudioChunk> {
self.chunk_consumer.try_pop()
}
pub fn poll_secondary(&mut self) -> Option<SecondaryChunk> {
self.secondary_consumer.as_mut().and_then(|c| c.try_pop())
}
pub fn set_denoise(&self, enabled: bool) {
self.shared.denoise_enabled.store(enabled, Ordering::SeqCst);
}
pub fn poll_event(&mut self) -> Option<Event> {
self.events.lock().ok().and_then(|mut q| q.pop_front())
}
pub fn dropped_chunks(&self) -> u64 {
self.chunk_consumer.dropped_count()
}
pub fn config(&self) -> &StreamConfig {
&self.config
}
pub fn native_format(&self) -> (u32, u16) {
*self
.shared
.native_format
.lock()
.unwrap_or_else(|e| e.into_inner())
}
fn open_backend_once(shared: &Arc<SharedState>) -> Result<()> {
let (rate, channels) = {
let be = shared.backend.lock().unwrap_or_else(|e| e.into_inner());
be.native_format()
};
{
let mut nf = shared
.native_format
.lock()
.unwrap_or_else(|e| e.into_inner());
*nf = (rate, channels);
}
let (producer, consumer) = raw_ring(RAW_RING_SAMPLES);
let sink = RawSink::new(producer, rate, channels);
{
let mut be = shared.backend.lock().unwrap_or_else(|e| e.into_inner());
start_backend_catching(&mut be, sink)?;
}
{
let mut rc = shared
.raw_consumer
.lock()
.unwrap_or_else(|e| e.into_inner());
*rc = Some(consumer);
}
shared.raw_generation.fetch_add(1, Ordering::SeqCst);
shared
.last_sample_ns
.store(monotonic_now_ns(), Ordering::SeqCst);
Ok(())
}
#[doc(hidden)]
pub fn switch_backend(&mut self, new_backend: Box<dyn CaptureBackend>) -> Result<()> {
if !self.started {
return Err(Error::InvalidState(
"switch_backend is only available on a started stream".into(),
));
}
self.shared.switching.store(true, Ordering::SeqCst);
{
let mut be = self
.shared
.backend
.lock()
.unwrap_or_else(|e| e.into_inner());
let _ = stop_backend_catching(&mut be);
let (rate, channels) = new_backend.native_format();
let (producer, consumer) = raw_ring(RAW_RING_SAMPLES);
let sink = RawSink::new(producer, rate, channels);
let mut new_backend = new_backend;
match start_backend_catching(&mut new_backend, sink) {
Ok(()) => {
{
let mut nf = self
.shared
.native_format
.lock()
.unwrap_or_else(|e| e.into_inner());
*nf = (rate, channels);
}
self.shared
.discontinuity_pending
.store(true, Ordering::SeqCst);
self.shared
.last_sample_ns
.store(monotonic_now_ns(), Ordering::SeqCst);
self.shared.raw_generation.fetch_add(1, Ordering::SeqCst);
*be = new_backend;
{
let mut rc = self
.shared
.raw_consumer
.lock()
.unwrap_or_else(|e| e.into_inner());
*rc = Some(consumer);
}
}
Err(e) => {
drop(be);
let _ = Self::open_backend_once(&self.shared);
self.shared
.discontinuity_pending
.store(true, Ordering::SeqCst);
self.shared.switching.store(false, Ordering::SeqCst);
return Err(e);
}
}
}
self.shared.switching.store(false, Ordering::SeqCst);
Ok(())
}
pub fn switch_source(&mut self, new_config: StreamConfig) -> Result<()> {
if !self.started {
return Err(Error::InvalidState(
"switch_source is only available on a started stream".into(),
));
}
if new_config.output != self.config.output {
return Err(Error::InvalidArg(
"output format cannot change during switch_source".into(),
));
}
if new_config.secondary_output != self.config.secondary_output {
return Err(Error::InvalidArg(
"secondary output format cannot change during switch_source".into(),
));
}
crate::validate_exclude_pids(&new_config)?;
let backend = crate::build_backend(&new_config)?;
self.switch_backend(backend)?;
self.config = StreamConfig {
kind: new_config.kind,
device_id: new_config.device_id,
target_pid: new_config.target_pid,
mode: new_config.mode,
exclude_self: new_config.exclude_self,
exclude_pids: new_config.exclude_pids,
..self.config.clone()
};
Ok(())
}
}
impl Drop for Stream {
fn drop(&mut self) {
if self.started {
self.stop();
}
}
}
fn apply_epoch(shared: &SharedState, raw_pts: i64) -> i64 {
let epoch = shared.recording_epoch_ns.load(Ordering::SeqCst);
if epoch == i64::MIN {
shared.recording_epoch_ns.store(raw_pts, Ordering::SeqCst);
0
} else {
raw_pts - epoch
}
}
fn apply_gain(data: &mut [f32], gain: f32) {
if gain != 1.0 {
for x in data.iter_mut() {
*x = (*x * gain).clamp(-1.0, 1.0);
}
}
}
fn run_intake(
shared: Arc<SharedState>,
mut chunk_producer: ChunkProducer,
mut secondary_producer: Option<SecondaryChunkProducer>,
initial_native: (u32, u16),
output: OutputFormat,
secondary_output: Option<OutputFormat>,
) {
let (mut rate, mut channels) = initial_native;
let mut normalizer = match build_normalizer(&shared, rate, channels, output, secondary_output) {
Ok(n) => n,
Err(e) => {
shared.push_event(Event::Error(format!("normalizer init failed: {e}")));
return;
}
};
let mut clock = ClockNormalizer::new();
let mut seq: u64 = 0; let mut sec_seq: u64 = 0; let mut current_generation = shared.raw_generation.load(Ordering::SeqCst);
let mut primary_resume_generation = 0;
let mut secondary_resume_generation = 0;
let mut overflow_baseline: u64 = 0;
let mut rec_primary = false;
let mut rec_secondary = false;
let mut disc_primary = false;
let mut disc_secondary = false;
let mut last_pts_primary: i64 = 0;
let mut last_pts_secondary: i64 = 0;
let out_channels = output.channels.max(1) as usize;
let sec_channels = secondary_output
.map(|f| f.channels.max(1) as usize)
.unwrap_or(1);
let mut scratch = vec![0.0f32; RAW_RING_SAMPLES];
loop {
let stopping = shared.stopping.load(Ordering::SeqCst);
let gen = shared.raw_generation.load(Ordering::SeqCst);
if gen != current_generation {
current_generation = gen;
let nf = *shared
.native_format
.lock()
.unwrap_or_else(|e| e.into_inner());
rate = nf.0;
channels = nf.1;
normalizer = match build_normalizer(&shared, rate, channels, output, secondary_output) {
Ok(n) => n,
Err(e) => {
shared.push_event(Event::Error(format!(
"normalizer rebuild failed after source change: {e}"
)));
return;
}
};
clock = ClockNormalizer::new();
overflow_baseline = 0;
}
if shared.recovered_pending.swap(false, Ordering::SeqCst) {
rec_primary = true;
rec_secondary = true;
}
if shared.discontinuity_pending.swap(false, Ordering::SeqCst) {
disc_primary = true;
disc_secondary = true;
}
let mut produced_any = false;
let mut push_err: Option<Error> = None;
let mut overflow_now = overflow_baseline;
{
let mut rc_guard = shared
.raw_consumer
.lock()
.unwrap_or_else(|e| e.into_inner());
if let Some(rc) = rc_guard.as_mut() {
let got = rc.pop_slice(&mut scratch);
overflow_now = rc.overflow_count();
if got > 0 {
let samples = &scratch[..got];
let device_pts = monotonic_now_ns();
let norm_pts = clock.normalize(device_pts);
if let Err(e) = normalizer.push(samples, norm_pts) {
push_err = Some(e);
} else {
shared
.last_sample_ns
.store(monotonic_now_ns(), Ordering::SeqCst);
produced_any = true;
}
}
}
}
if let Some(e) = push_err {
shared.push_event(Event::Error(format!("normalizer push failed: {e}")));
return;
}
if overflow_now > overflow_baseline {
disc_primary = true;
disc_secondary = true;
}
overflow_baseline = overflow_now;
if stopping {
normalizer.flush();
}
let gain = f32::from_bits(shared.gain_bits.load(Ordering::Relaxed));
let mut emitted_any = false;
while let Some((mut data, raw_pts)) = normalizer.pop_chunk() {
if shared.paused.load(Ordering::SeqCst) {
continue;
}
let pts_ns = apply_epoch(&shared, raw_pts).max(last_pts_primary);
last_pts_primary = pts_ns;
debug_assert_eq!(data.len() % out_channels, 0);
let frames = data.len() / out_channels;
apply_gain(&mut data, gain);
let (peak, rms) = peak_rms(&data);
let mut flags = ChunkFlags::empty();
if rec_primary {
flags |= ChunkFlags::RECOVERED | ChunkFlags::DISCONTINUITY;
}
if disc_primary {
flags |= ChunkFlags::DISCONTINUITY;
}
let mut chunk = AudioChunk {
data,
frames,
pts_ns,
seq,
flags,
dropped_before: 0, peak,
rms,
};
{
let _g = shared.delivery.lock().unwrap_or_else(|e| e.into_inner());
if shared.paused.load(Ordering::SeqCst) {
continue;
}
let resume_generation = shared.resume_generation.load(Ordering::SeqCst);
if resume_generation != primary_resume_generation {
chunk.flags |= ChunkFlags::DISCONTINUITY;
}
rec_primary = false;
disc_primary = false;
seq += 1;
if let Some(total) = chunk_producer.push(chunk) {
shared.push_event(Event::ChunkDropped { count: total });
}
primary_resume_generation = resume_generation;
emitted_any = true;
}
}
if let Some(sec_prod) = secondary_producer.as_mut() {
while let Some((mut samples, raw_pts)) = normalizer.pop_secondary() {
if shared.paused.load(Ordering::SeqCst) {
continue;
}
let pts_ns = apply_epoch(&shared, raw_pts).max(last_pts_secondary);
last_pts_secondary = pts_ns;
debug_assert_eq!(samples.len() % sec_channels, 0);
let frames = samples.len() / sec_channels;
apply_gain(&mut samples, gain);
let (peak, rms) = peak_rms(&samples);
let mut flags = ChunkFlags::empty();
if rec_secondary {
flags |= ChunkFlags::RECOVERED | ChunkFlags::DISCONTINUITY;
}
if disc_secondary {
flags |= ChunkFlags::DISCONTINUITY;
}
let mut chunk = SecondaryChunk {
samples,
frames,
pts_ns,
seq: sec_seq,
flags,
dropped_before: 0, peak,
rms,
};
{
let _g = shared.delivery.lock().unwrap_or_else(|e| e.into_inner());
if shared.paused.load(Ordering::SeqCst) {
continue;
}
let resume_generation = shared.resume_generation.load(Ordering::SeqCst);
if resume_generation != secondary_resume_generation {
chunk.flags |= ChunkFlags::DISCONTINUITY;
}
rec_secondary = false;
disc_secondary = false;
sec_seq += 1;
let _ = sec_prod.push(chunk);
secondary_resume_generation = resume_generation;
emitted_any = true;
}
}
}
if stopping {
break;
}
if !produced_any && !emitted_any {
thread::sleep(Duration::from_millis(2));
}
}
}
fn run_watchdog(shared: Arc<SharedState>) {
let mut stalled = false;
let mut backoff = BACKOFF_MIN;
loop {
if shared.stopping.load(Ordering::SeqCst) {
break;
}
thread::sleep(WATCHDOG_TICK);
if shared.stopping.load(Ordering::SeqCst) {
break;
}
if shared.switching.load(Ordering::SeqCst) {
continue;
}
let now = monotonic_now_ns();
let last = shared.last_sample_ns.load(Ordering::SeqCst);
let idle_ns = now.saturating_sub(last);
let idle = Duration::from_nanos(idle_ns.max(0) as u64);
if !stalled {
if idle >= STALL_THRESHOLD {
stalled = true;
backoff = BACKOFF_MIN;
shared.push_event(Event::StreamStalled);
}
continue;
}
{
let mut be = shared.backend.lock().unwrap_or_else(|e| e.into_inner());
let _ = stop_backend_catching(&mut be);
}
if shared.stopping.load(Ordering::SeqCst) {
break;
}
let reopened = match Stream::open_backend_once(&shared) {
Ok(()) => true,
Err(e) => {
shared.push_event(Event::Error(format!("reopen failed: {e}")));
false
}
};
if reopened {
shared.recovered_pending.store(true, Ordering::SeqCst);
stalled = false;
shared.push_event(Event::StreamRecovered);
backoff = BACKOFF_MIN;
} else {
let jittered = jittered_backoff(backoff);
sleep_interruptible(&shared, jittered);
backoff = (backoff * 2).min(BACKOFF_MAX);
}
}
}
fn peak_rms(data: &[f32]) -> (f32, f32) {
if data.is_empty() {
return (0.0, 0.0);
}
let mut peak = 0.0f32;
let mut sum_sq = 0.0f64;
for &x in data {
let a = x.abs();
if a > peak {
peak = a;
}
sum_sq += (x as f64) * (x as f64);
}
let rms = (sum_sq / data.len() as f64).sqrt() as f32;
(peak, rms)
}
fn jittered_backoff(base: Duration) -> Duration {
let base_ns = base.as_nanos() as u64;
let entropy = monotonic_now_ns() as u64;
let span = (base_ns / 8).max(1);
let delta = (entropy % (2 * span)) as i64 - span as i64;
let result = base_ns as i64 + delta;
Duration::from_nanos(result.max(0) as u64)
}
fn sleep_interruptible(shared: &Arc<SharedState>, dur: Duration) {
let step = Duration::from_millis(50);
let mut remaining = dur;
while remaining > Duration::ZERO {
if shared.stopping.load(Ordering::SeqCst) {
return;
}
let s = step.min(remaining);
thread::sleep(s);
remaining = remaining.saturating_sub(s);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::mock::{
MockBackend, PanicMode, PanickingMockBackend, StallThenPanicOnReopenBackend,
StallableMockBackend,
};
use flexaudio_core::types::SourceKind;
use std::time::Instant;
fn collect_for(stream: &mut Stream, dur: Duration) -> Vec<AudioChunk> {
let mut chunks = Vec::new();
let start = Instant::now();
while start.elapsed() < dur {
while let Some(c) = stream.poll_chunk() {
chunks.push(c);
}
thread::sleep(Duration::from_millis(5));
}
chunks
}
fn wait_until<F: FnMut() -> bool>(mut cond: F, timeout: Duration) -> bool {
let start = Instant::now();
while start.elapsed() < timeout {
if cond() {
return true;
}
thread::sleep(Duration::from_millis(10));
}
cond()
}
fn open_err(result: Result<Stream>, ctx: &str) -> Error {
match result {
Ok(_) => panic!("{ctx}: expected an error, got Ok"),
Err(e) => e,
}
}
#[test]
fn open_validates_exclusion_pids_for_system_capture() {
for kind in [SourceKind::SystemLoopback, SourceKind::Mix] {
let config = StreamConfig {
kind,
exclude_pids: vec![0],
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
assert!(matches!(
Stream::open(config, backend),
Err(Error::InvalidArg(message))
if message == "exclude_pids: pid 0 is not a valid process id"
));
for exclude_pids in [vec![], vec![42]] {
let config = StreamConfig {
kind,
exclude_pids,
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
assert!(Stream::open(config, backend).is_ok());
}
}
for kind in [SourceKind::Mic, SourceKind::ProcessLoopback] {
let config = StreamConfig {
kind,
exclude_pids: vec![0],
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
assert!(Stream::open(config, backend).is_ok());
}
}
#[test]
fn switch_source_rejects_zero_exclusion_pid_before_replacing_backend() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start mock capture");
for kind in [SourceKind::SystemLoopback, SourceKind::Mix] {
let new_config = StreamConfig {
kind,
exclude_pids: vec![42, 0],
..Default::default()
};
assert!(matches!(
stream.switch_source(new_config),
Err(Error::InvalidArg(message))
if message == "exclude_pids: pid 0 is not a valid process id"
));
assert_eq!(stream.config.kind, SourceKind::Mic);
assert!(stream.config.exclude_pids.is_empty());
}
stream.stop();
}
#[test]
fn open_rejects_zero_ring_capacity() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
ring_capacity_chunks: 0,
..Default::default()
};
let err = open_err(Stream::open(config, backend), "capacity 0");
assert!(
matches!(err, Error::InvalidArg(_)),
"expected InvalidArg: {err:?}"
);
}
#[test]
fn open_rejects_invalid_output_channels() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
output: OutputFormat {
sample_rate: 48_000,
channels: 3,
},
..Default::default()
};
let err = open_err(Stream::open(config, backend), "ch=3");
assert!(
matches!(err, Error::UnsupportedFormat(_)),
"expected UnsupportedFormat: {err:?}"
);
}
#[test]
fn open_rejects_out_of_range_output_rate() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
output: OutputFormat {
sample_rate: 1_000_000,
channels: 2,
},
..Default::default()
};
let err = open_err(Stream::open(config, backend), "extreme rate");
assert!(
matches!(err, Error::UnsupportedFormat(_)),
"expected UnsupportedFormat: {err:?}"
);
}
#[test]
fn open_rejects_zero_native_format() {
struct ZeroFormatBackend;
impl CaptureBackend for ZeroFormatBackend {
fn native_format(&self) -> (u32, u16) {
(0, 0)
}
fn start(&mut self, _sink: RawSink) -> Result<()> {
Ok(())
}
fn stop(&mut self) {}
}
let backend = Box::new(ZeroFormatBackend);
let err = open_err(
Stream::open(StreamConfig::default(), backend),
"native_format 0",
);
assert!(
matches!(err, Error::InvalidArg(_)),
"expected InvalidArg: {err:?}"
);
}
#[test]
fn poll_event_yields_chunk_dropped() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
ring_capacity_chunks: 1,
..Default::default()
};
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let got_drop = wait_until(
|| {
while let Some(ev) = stream.poll_event() {
if matches!(ev, Event::ChunkDropped { .. }) {
return true;
}
}
false
},
Duration::from_secs(3),
);
stream.stop();
assert!(
got_drop,
"expected to retrieve ChunkDropped through poll_event"
);
}
#[test]
fn poll_event_is_none_when_empty() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
assert!(stream.poll_event().is_none());
}
#[test]
fn watchdog_detects_stall_and_flags_recovered() {
let backend = Box::new(StallableMockBackend::new(
48_000,
2,
440.0,
Duration::from_millis(300),
));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let mut chunks: Vec<AudioChunk> = Vec::new();
let mut saw_stalled = false;
let mut saw_recovered = false;
let deadline = Instant::now() + Duration::from_secs(8);
let mut recovered_chunk_seen = false;
while Instant::now() < deadline && !recovered_chunk_seen {
while let Some(c) = stream.poll_chunk() {
if c.flags.contains(ChunkFlags::RECOVERED) {
recovered_chunk_seen = true;
}
chunks.push(c);
}
while let Some(ev) = stream.poll_event() {
match ev {
Event::StreamStalled => saw_stalled = true,
Event::StreamRecovered => saw_recovered = true,
_ => {}
}
}
thread::sleep(Duration::from_millis(20));
}
stream.stop();
while let Some(c) = stream.poll_chunk() {
if c.flags.contains(ChunkFlags::RECOVERED) {
recovered_chunk_seen = true;
}
chunks.push(c);
}
assert!(saw_stalled, "expected Event::StreamStalled to fire");
assert!(saw_recovered, "expected Event::StreamRecovered to fire");
assert!(
recovered_chunk_seen,
"expected RECOVERED on the first post-recovery chunk"
);
let recovered: Vec<&AudioChunk> = chunks
.iter()
.filter(|c| c.flags.contains(ChunkFlags::RECOVERED))
.collect();
assert!(!recovered.is_empty());
for c in &recovered {
assert!(
c.flags.contains(ChunkFlags::DISCONTINUITY),
"expected DISCONTINUITY with RECOVERED: flags={:?}",
c.flags
);
}
for w in chunks.windows(2) {
assert!(
w[1].seq > w[0].seq,
"seq should increase monotonically across recovery: {} -> {}",
w[0].seq,
w[1].seq
);
}
}
#[test]
fn no_recovered_flag_under_steady_feed() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let chunks = collect_for(&mut stream, Duration::from_millis(500));
let mut saw_stalled = false;
while let Some(ev) = stream.poll_event() {
if matches!(ev, Event::StreamStalled) {
saw_stalled = true;
}
}
stream.stop();
assert!(!chunks.is_empty(), "expected chunks with steady input");
assert!(
!saw_stalled,
"steady input should not be reported as stalled"
);
for c in &chunks {
assert!(
!c.flags.contains(ChunkFlags::RECOVERED),
"RECOVERED should not be set with steady input: flags={:?}",
c.flags
);
}
}
#[test]
fn pause_stops_delivering_chunks() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let got_before = wait_until(|| stream.poll_chunk().is_some(), Duration::from_secs(2));
assert!(got_before, "expected a chunk before pause");
stream.pause();
while stream.poll_chunk().is_some() {}
let after = collect_for(&mut stream, Duration::from_millis(300));
stream.stop();
assert!(
after.is_empty(),
"expected no new chunks while paused; received {}",
after.len()
);
}
#[test]
fn long_pause_does_not_trigger_stall() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let got_before = wait_until(|| stream.poll_chunk().is_some(), Duration::from_secs(2));
assert!(got_before, "expected a chunk before pause");
stream.pause();
while stream.poll_chunk().is_some() {}
let mut saw_stalled = false;
let mut saw_recovered = false;
let deadline = Instant::now() + Duration::from_millis(2800);
while Instant::now() < deadline {
while let Some(ev) = stream.poll_event() {
match ev {
Event::StreamStalled => saw_stalled = true,
Event::StreamRecovered => saw_recovered = true,
_ => {}
}
}
assert!(
stream.is_paused(),
"is_paused should remain true during the pause window"
);
thread::sleep(Duration::from_millis(20));
}
assert!(
!saw_stalled,
"StreamStalled should not fire during a long pause"
);
assert!(
!saw_recovered,
"StreamRecovered should not fire because there was no stall"
);
stream.resume();
let resumed = wait_until(|| stream.poll_chunk().is_some(), Duration::from_secs(2));
stream.stop();
assert!(resumed, "chunk delivery should resume after resume");
}
#[test]
fn resume_flags_discontinuity_and_keeps_seq_continuous() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let before = collect_for(&mut stream, Duration::from_millis(200));
assert!(!before.is_empty(), "expected chunks before pause");
let last_seq = before.last().unwrap().seq;
stream.pause();
let mut last_seq = last_seq;
while let Some(c) = stream.poll_chunk() {
last_seq = c.seq;
}
assert!(collect_for(&mut stream, Duration::from_millis(150)).is_empty());
stream.resume();
let mut first_after: Option<AudioChunk> = None;
let got = wait_until(
|| match stream.poll_chunk() {
Some(c) => {
first_after = Some(c);
true
}
None => false,
},
Duration::from_secs(2),
);
stream.stop();
assert!(got, "expected a chunk after resume");
let first = first_after.unwrap();
assert!(
first.flags.contains(ChunkFlags::DISCONTINUITY),
"expected DISCONTINUITY on the first chunk after resume: flags={:?}",
first.flags
);
assert_eq!(
first.seq,
last_seq + 1,
"seq should remain continuous across pause ({last_seq} -> {})",
first.seq
);
assert_eq!(
first.dropped_before, 0,
"pause should not cause dropped chunks"
);
}
#[test]
fn resume_stress_marks_first_chunk_of_each_tap_discontinuous() {
const ROUNDS: usize = 300;
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
ring_capacity_chunks: 200,
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let mut failures = 0usize;
for round in 0..ROUNDS {
stream.pause();
while stream.poll_chunk().is_some() {}
while stream.poll_secondary().is_some() {}
stream.resume();
let mut primary = None;
let got_primary = wait_until(
|| match stream.poll_chunk() {
Some(chunk) => {
primary = Some(chunk);
true
}
None => false,
},
Duration::from_secs(2),
);
let mut secondary = None;
let got_secondary = wait_until(
|| match stream.poll_secondary() {
Some(chunk) => {
secondary = Some(chunk);
true
}
None => false,
},
Duration::from_secs(2),
);
let primary_ok = got_primary
&& primary.is_some_and(|chunk| chunk.flags.contains(ChunkFlags::DISCONTINUITY));
let secondary_ok = got_secondary
&& secondary.is_some_and(|chunk| chunk.flags.contains(ChunkFlags::DISCONTINUITY));
if !primary_ok || !secondary_ok {
failures += 1;
eprintln!("round {round}: primary_ok={primary_ok}, secondary_ok={secondary_ok}");
}
}
stream.stop();
assert_eq!(
failures, 0,
"{failures} / {ROUNDS} resume attempts had no DISCONTINUITY on the first primary or \
secondary chunk"
);
}
#[test]
fn resume_flags_secondary_first_chunk_with_raw_intake_blocked() {
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
ring_capacity_chunks: 200,
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let mut last_secondary_seq = None;
let got_before = wait_until(
|| {
while stream.poll_chunk().is_some() {}
if let Some(chunk) = stream.poll_secondary() {
last_secondary_seq = Some(chunk.seq);
true
} else {
false
}
},
Duration::from_secs(2),
);
assert!(got_before, "expected a secondary chunk before pause");
stream.pause();
let shared = stream.shared.clone();
{
let raw = shared
.raw_consumer
.lock()
.unwrap_or_else(|e| e.into_inner());
while stream.poll_chunk().is_some() {}
while let Some(chunk) = stream.poll_secondary() {
last_secondary_seq = Some(chunk.seq);
}
stream.resume();
drop(raw);
}
let mut first_after = None;
let got_after = wait_until(
|| {
while stream.poll_chunk().is_some() {}
if let Some(chunk) = stream.poll_secondary() {
first_after = Some(chunk);
true
} else {
false
}
},
Duration::from_secs(2),
);
let overflow_count = shared
.raw_consumer
.lock()
.unwrap_or_else(|e| e.into_inner())
.as_ref()
.expect("raw consumer")
.overflow_count();
stream.stop();
assert!(got_after, "expected a secondary chunk after resume");
assert_eq!(overflow_count, 0, "raw overflow must not mask resume flags");
let first = first_after.expect("first secondary chunk after resume");
assert!(
first.flags.contains(ChunkFlags::DISCONTINUITY),
"expected DISCONTINUITY on the first secondary chunk after resume: {:?}",
first.flags
);
assert!(
!first.flags.contains(ChunkFlags::RECOVERED),
"watchdog recovery must not mask resume flags"
);
assert_eq!(
first.seq,
last_secondary_seq.expect("last secondary sequence before pause") + 1,
"secondary sequence should remain continuous across pause"
);
assert_eq!(
first.dropped_before, 0,
"pause must not drop secondary chunks"
);
}
#[test]
fn resume_without_pause_is_noop() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let _ = collect_for(&mut stream, Duration::from_millis(200));
stream.resume();
let after = collect_for(&mut stream, Duration::from_millis(200));
stream.stop();
assert!(!after.is_empty(), "expected chunks to arrive");
for c in &after {
assert!(
!c.flags.contains(ChunkFlags::DISCONTINUITY),
"resume without pause should not set DISCONTINUITY: flags={:?}",
c.flags
);
}
}
#[test]
fn double_pause_then_single_resume_recovers() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let before = collect_for(&mut stream, Duration::from_millis(200));
assert!(!before.is_empty(), "expected chunks before pause");
stream.pause();
stream.pause();
assert!(stream.is_paused());
while stream.poll_chunk().is_some() {}
assert!(collect_for(&mut stream, Duration::from_millis(150)).is_empty());
stream.resume();
assert!(!stream.is_paused());
let got = wait_until(|| stream.poll_chunk().is_some(), Duration::from_secs(2));
stream.stop();
assert!(got, "delivery should resume after one resume call");
}
#[test]
fn gain_scales_samples_and_meters() {
for (gain, lo, hi) in [(2.0f32, 0.95f32, 1.0f32), (0.5, 0.2, 0.3)] {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
gain,
..Default::default()
};
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let chunks = collect_for(&mut stream, Duration::from_millis(300));
stream.stop();
assert!(!chunks.is_empty(), "expected chunks with gain={gain}");
let mut max_peak = 0.0f32;
for c in &chunks {
let recomputed = c.data.iter().fold(0.0f32, |m, &x| m.max(x.abs()));
assert_eq!(
c.peak, recomputed,
"peak should be computed from data after gain is applied"
);
max_peak = max_peak.max(c.peak);
}
assert!(
(lo..=hi).contains(&max_peak),
"expected peak for gain={gain} in {lo}..={hi}: {max_peak}"
);
}
}
#[test]
fn set_gain_takes_effect_mid_stream() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
assert_eq!(stream.gain(), 1.0, "default gain is 1.0");
let got_before = wait_until(|| stream.poll_chunk().is_some(), Duration::from_secs(2));
assert!(got_before, "expected a chunk before set_gain");
stream.set_gain(0.0).expect("set_gain(0.0)");
assert_eq!(stream.gain(), 0.0);
let got_silent = wait_until(
|| matches!(stream.poll_chunk(), Some(c) if c.peak == 0.0),
Duration::from_secs(2),
);
assert!(got_silent, "expected a silent chunk after set_gain(0.0)");
let after = collect_for(&mut stream, Duration::from_millis(300));
stream.stop();
assert!(
!after.is_empty(),
"chunks should continue to flow during silence"
);
for c in &after {
assert!(
c.data.iter().all(|&x| x == 0.0),
"all samples should be 0 at gain 0.0"
);
assert_eq!(c.peak, 0.0);
assert_eq!(c.rms, 0.0);
}
}
#[test]
fn gain_clamps_to_unit_range() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
gain: 100.0,
..Default::default()
};
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let chunks = collect_for(&mut stream, Duration::from_millis(300));
stream.stop();
assert!(!chunks.is_empty(), "expected chunks to arrive");
let mut max_peak = 0.0f32;
for c in &chunks {
assert!(
c.data.iter().all(|&x| (-1.0..=1.0).contains(&x)),
"samples should not exceed ±1.0"
);
max_peak = max_peak.max(c.peak);
}
assert_eq!(max_peak, 1.0, "clamping should make the peak exactly 1.0");
}
#[test]
fn invalid_gain_rejected() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
gain: -1.0,
..Default::default()
};
let err = open_err(Stream::open(config, backend), "gain=-1.0");
assert!(
matches!(err, Error::InvalidArg(_)),
"expected InvalidArg: {err:?}"
);
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let config = StreamConfig {
gain: f32::NAN,
..Default::default()
};
let err = open_err(Stream::open(config, backend), "gain=NaN");
assert!(
matches!(err, Error::InvalidArg(_)),
"expected InvalidArg: {err:?}"
);
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let stream = Stream::open(StreamConfig::default(), backend).expect("open");
assert!(matches!(stream.set_gain(-1.0), Err(Error::InvalidArg(_))));
assert!(matches!(
stream.set_gain(f32::NAN),
Err(Error::InvalidArg(_))
));
assert_eq!(
stream.gain(),
1.0,
"failed set_gain should not change the current value"
);
}
#[test]
fn backend_panic_in_start_returns_err_not_silent_death() {
let backend = Box::new(PanickingMockBackend::new(
48_000,
2,
440.0,
PanicMode::Start,
));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
let result = stream.start();
match result {
Ok(()) => panic!("backend panicked in start, but start() returned Ok"),
Err(Error::Backend(msg)) => {
assert!(
msg.contains("panicked"),
"expected Error::Backend with a message identifying the panic: {msg}"
);
}
Err(other) => panic!("expected Error::Backend, got a different error: {other:?}"),
}
stream.stop();
}
#[test]
fn backend_panic_in_stop_does_not_kill_process() {
let backend = Box::new(PanickingMockBackend::new(48_000, 2, 440.0, PanicMode::Stop));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let chunks = collect_for(&mut stream, Duration::from_millis(300));
assert!(
!chunks.is_empty(),
"chunks should flow normally before stop (happy path unchanged)"
);
stream.stop();
let _ = stream.poll_chunk();
let _ = stream.poll_event();
}
#[test]
fn backend_panic_on_watchdog_reopen_surfaces_event_error() {
let backend = Box::new(StallThenPanicOnReopenBackend::new(
48_000,
2,
440.0,
Duration::from_millis(300),
));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let mut saw_stalled = false;
let mut saw_reopen_error = false;
let deadline = Instant::now() + Duration::from_secs(8);
while Instant::now() < deadline && !saw_reopen_error {
while stream.poll_chunk().is_some() {}
while let Some(ev) = stream.poll_event() {
match ev {
Event::StreamStalled => saw_stalled = true,
Event::Error(msg) if msg.contains("reopen failed") => {
saw_reopen_error = true;
}
_ => {}
}
}
thread::sleep(Duration::from_millis(20));
}
stream.stop();
assert!(
saw_stalled,
"expected stall detection (Event::StreamStalled)"
);
assert!(
saw_reopen_error,
"backend panic during reopen should surface as Event::Error(\"reopen failed: ...\") \
(no silent death)"
);
}
#[test]
fn recording_clock_is_zero_based() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let mut first: Option<AudioChunk> = None;
let got = wait_until(
|| match stream.poll_chunk() {
Some(c) => {
first = Some(c);
true
}
None => false,
},
Duration::from_secs(2),
);
stream.stop();
assert!(got, "expected the first chunk to arrive");
let first = first.unwrap();
assert_eq!(
first.pts_ns, 0,
"the first delivered chunk should start at recording time zero (pts_ns == 0): {}",
first.pts_ns
);
}
#[test]
fn pause_preserves_absolute_clock() {
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(StreamConfig::default(), backend).expect("open");
stream.start().expect("start");
let before = collect_for(&mut stream, Duration::from_millis(250));
assert!(!before.is_empty(), "expected chunks before pause");
let mut last = before.last().cloned().unwrap();
stream.pause();
while let Some(c) = stream.poll_chunk() {
last = c;
}
let d = Duration::from_millis(600);
thread::sleep(d);
stream.resume();
let mut first_after: Option<AudioChunk> = None;
let got = wait_until(
|| match stream.poll_chunk() {
Some(c) => {
first_after = Some(c);
true
}
None => false,
},
Duration::from_secs(2),
);
stream.stop();
assert!(got, "expected a chunk after resume");
let first = first_after.unwrap();
assert!(
first.flags.contains(ChunkFlags::DISCONTINUITY),
"expected DISCONTINUITY on the first chunk after resume: {:?}",
first.flags
);
assert_eq!(
first.seq,
last.seq + 1,
"seq should remain continuous across pause"
);
assert_eq!(first.dropped_before, 0, "pause should not drop chunks");
let delta = first.pts_ns - last.pts_ns;
let d_ns = d.as_nanos() as i64;
assert!(
delta >= d_ns * 4 / 5,
"pts should advance by at least the pause duration (>= {} ns): delta={delta} ns",
d_ns * 4 / 5
);
assert!(
delta <= d_ns + 500_000_000,
"pts should not advance too far (<= D + 500ms): delta={delta} ns"
);
}
#[test]
fn dual_output_delivers_primary_and_secondary() {
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
ring_capacity_chunks: 200,
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let mut primary: Vec<AudioChunk> = Vec::new();
let mut secondary: Vec<SecondaryChunk> = Vec::new();
let deadline = Instant::now() + Duration::from_millis(500);
while Instant::now() < deadline {
while let Some(c) = stream.poll_chunk() {
primary.push(c);
}
while let Some(c) = stream.poll_secondary() {
secondary.push(c);
}
thread::sleep(Duration::from_millis(5));
}
stream.stop();
while let Some(c) = stream.poll_chunk() {
primary.push(c);
}
while let Some(c) = stream.poll_secondary() {
secondary.push(c);
}
assert!(!primary.is_empty(), "expected primary chunks");
assert!(!secondary.is_empty(), "expected secondary chunks");
for c in &primary {
assert_eq!(
c.data.len(),
960 * 2,
"primary is 48k/stereo = 1920 samples"
);
}
for c in &secondary {
assert_eq!(c.samples.len(), 320, "secondary is 16k/mono = 320 samples");
}
assert_eq!(primary[0].pts_ns, 0, "first primary chunk starts at zero");
for w in secondary.windows(2) {
assert!(
w[1].pts_ns >= w[0].pts_ns,
"secondary PTS should not decrease"
);
}
assert!(
secondary[0].pts_ns >= 0,
"secondary PTS should be non-negative (based on primary epoch)"
);
assert_eq!(secondary[0].seq, 0);
for w in secondary.windows(2) {
assert_eq!(
w[1].seq,
w[0].seq + 1,
"secondary seq should increase consecutively"
);
}
}
#[test]
fn start_after_pause_delivers_the_secondary_tap() {
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
ring_capacity_chunks: 200,
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.pause();
stream.start().expect("start");
assert!(
!stream.is_paused(),
"start should clear the pre-start pause"
);
let mut primary = false;
let mut secondary = false;
let got_both = wait_until(
|| {
while stream.poll_chunk().is_some() {
primary = true;
}
while stream.poll_secondary().is_some() {
secondary = true;
}
primary && secondary
},
Duration::from_secs(2),
);
stream.stop();
assert!(primary, "expected primary delivery after a pre-start pause");
assert!(
got_both && secondary,
"expected secondary delivery after a pre-start pause"
);
}
#[test]
fn secondary_output_cannot_change_on_switch() {
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let new_config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 8_000,
channels: 1,
}),
..Default::default()
};
let err = stream.switch_source(new_config);
stream.stop();
assert!(
matches!(err, Err(Error::InvalidArg(_))),
"expected InvalidArg when changing the secondary format: {err:?}"
);
}
struct FloodingMockBackend {
running: Arc<AtomicBool>,
handle: Option<thread::JoinHandle<()>>,
}
impl FloodingMockBackend {
fn new() -> Self {
Self {
running: Arc::new(AtomicBool::new(false)),
handle: None,
}
}
}
impl CaptureBackend for FloodingMockBackend {
fn native_format(&self) -> (u32, u16) {
(48_000, 2)
}
fn start(&mut self, mut sink: RawSink) -> Result<()> {
self.running.store(true, Ordering::SeqCst);
let running = self.running.clone();
let handle = thread::spawn(move || {
let burst = vec![0.1f32; RAW_RING_SAMPLES * 2];
while running.load(Ordering::SeqCst) {
sink.push(&burst, 0);
thread::sleep(Duration::from_millis(5));
}
});
self.handle = Some(handle);
Ok(())
}
fn stop(&mut self) {
self.running.store(false, Ordering::SeqCst);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
#[test]
fn ring_overflow_marks_discontinuity() {
let backend = Box::new(FloodingMockBackend::new());
let config = StreamConfig {
ring_capacity_chunks: 200,
..Default::default()
};
let mut stream = Stream::open(config, backend).expect("open");
stream.start().expect("start");
let mut saw_discontinuity = false;
let deadline = Instant::now() + Duration::from_secs(2);
while Instant::now() < deadline && !saw_discontinuity {
while let Some(c) = stream.poll_chunk() {
if c.flags.contains(ChunkFlags::DISCONTINUITY) {
saw_discontinuity = true;
}
}
thread::sleep(Duration::from_millis(10));
}
stream.stop();
while let Some(c) = stream.poll_chunk() {
if c.flags.contains(ChunkFlags::DISCONTINUITY) {
saw_discontinuity = true;
}
}
assert!(
saw_discontinuity,
"RawRing overflow should set DISCONTINUITY"
);
}
const DENOISE_TAP_TIMEOUT: Duration = Duration::from_secs(15);
#[test]
fn denoise_enabled_still_delivers_both_taps() {
fn drain_taps(stream: &mut Stream, primary: &mut usize, secondary: &mut usize) {
while stream.poll_chunk().is_some() {
*primary += 1;
}
while let Some(chunk) = stream.poll_secondary() {
assert_eq!(
chunk.samples.len(),
320,
"secondary is 16k/mono = 320 samples"
);
*secondary += 1;
}
}
let config = StreamConfig {
secondary_output: Some(OutputFormat {
sample_rate: 16_000,
channels: 1,
}),
..Default::default()
};
let backend = Box::new(MockBackend::new(48_000, 2, 440.0));
let mut stream = Stream::open(config, backend).expect("open");
stream.set_denoise(true);
stream.start().expect("start");
let mut primary = 0usize;
let mut secondary = 0usize;
let primary_ok = wait_until(
|| {
drain_taps(&mut stream, &mut primary, &mut secondary);
primary > 0
},
DENOISE_TAP_TIMEOUT,
);
let secondary_ok = wait_until(
|| {
drain_taps(&mut stream, &mut primary, &mut secondary);
secondary > 0
},
DENOISE_TAP_TIMEOUT,
);
let last_sample_ns = stream.shared.last_sample_ns.load(Ordering::SeqCst);
stream.stop();
assert!(
primary_ok,
"no primary delivery with denoise within {DENOISE_TAP_TIMEOUT:?}: \
primary={primary}, secondary={secondary}, last_sample_ns={last_sample_ns}"
);
assert!(
secondary_ok,
"no secondary delivery with denoise within {DENOISE_TAP_TIMEOUT:?}: \
primary={primary}, secondary={secondary}, last_sample_ns={last_sample_ns}"
);
}
}