use std::collections::VecDeque;
use std::panic::AssertUnwindSafe;
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32, AtomicU64, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
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, NormalizedChunk, Normalizer};
use flexaudio_core::raw_ring::{raw_ring, RawConsumer};
use flexaudio_core::secondary_ring::{
secondary_chunk_ring, SecondaryChunkConsumer, SecondaryChunkProducer,
};
use flexaudio_core::types::{
AudioChunk, AudioLoss, ChunkFlags, Error, ErrorContext, ErrorGroup, ErrorKind, Event,
Operation, OutputFormat, OutputTap, Permission, Result, SecondaryChunk, ShutdownReport,
StreamConfig,
};
use flexaudio_core::CaptureDiagnostics;
use std::num::NonZeroU64;
mod control;
mod delivery;
mod intake;
mod lifecycle;
mod processor;
mod shared;
mod source;
mod tap_drain;
#[cfg(test)]
use control::{drain_backend_events, MailboxDrain};
use control::{run_watchdog, stop_backend_reconciling};
use intake::run_intake;
use processor::build_normalizer;
mod terminal;
use terminal::TerminalFailure;
#[cfg(test)]
mod contract_tests;
#[cfg(test)]
mod permission_tests;
#[cfg(test)]
mod recovery_tests;
pub(crate) const RAW_RING_SAMPLES: usize = 48_000;
const WATCHDOG_TICK: Duration = Duration::from_millis(250);
const MAX_BACKEND_EVENTS_PER_TICK: usize = 64;
const MAX_FINAL_EVENT_BATCHES: usize = 64;
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,
capture_consumer: Option<ChunkConsumer>,
secondary_consumer: Option<SecondaryChunkConsumer>,
events: Arc<Mutex<VecDeque<Event>>>,
worker: Option<JoinHandle<()>>,
watchdog: Option<JoinHandle<()>>,
started: bool,
shutdown: Option<ShutdownReport>,
}
struct SharedState {
cleanup: Mutex<Vec<Error>>,
backend_stopped: AtomicBool,
raw_diagnostics: Mutex<Option<CaptureDiagnostics>>,
raw_overflow_reported: Mutex<u64>,
primary_frame_index: AtomicU64,
secondary_frame_index: AtomicU64,
capture_frame_index: AtomicU64,
capture_enabled: AtomicBool,
capture_producer: Mutex<Option<ChunkProducer>>,
backend: Mutex<Box<dyn CaptureBackend>>,
raw_consumer: Mutex<Option<RawConsumer>>,
raw_generation: AtomicU64,
last_sample_ns: AtomicI64,
stopping: AtomicBool,
terminal: TerminalFailure,
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: Arc<AtomicU32>,
recording_epoch_ns: AtomicI64,
denoise_enabled: AtomicBool,
secondary_producer: Mutex<Option<SecondaryChunkProducer>>,
}
struct RawSnapshot {
generation: u64,
denoise_enabled: bool,
native_format: (u32, u16),
samples: usize,
overflows: u64,
losses: Result<Vec<AudioLoss>>,
recovered: bool,
discontinuity: bool,
}
#[derive(Clone, Copy)]
enum GenerationChange {
Initial,
Recovery,
Switch,
}
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())),
}
}
fn stop_backend_catching(be: &mut Box<dyn CaptureBackend>) -> Result<()> {
match std::panic::catch_unwind(AssertUnwindSafe(|| be.stop_checked())) {
Ok(result) => result,
Err(_) => Err(Error::Backend("backend panicked during stop".into())),
}
.map_err(|error| error.with_context(ErrorContext::new(Operation::Stop)))
}
impl Stream {
pub fn open(config: StreamConfig, backend: Box<dyn CaptureBackend>) -> Result<Stream> {
validate_chunk_ms(config.chunk_ms)?;
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(),
));
}
Normalizer::new(native_format.0, native_format.1, config.output)?;
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 {
cleanup: Mutex::new(Vec::new()),
backend_stopped: AtomicBool::new(false),
raw_diagnostics: Mutex::new(None),
raw_overflow_reported: Mutex::new(0),
primary_frame_index: AtomicU64::new(0),
secondary_frame_index: AtomicU64::new(0),
capture_frame_index: AtomicU64::new(0),
capture_enabled: AtomicBool::new(false),
capture_producer: Mutex::new(None),
backend: Mutex::new(backend),
raw_consumer: Mutex::new(None),
raw_generation: AtomicU64::new(0),
last_sample_ns: AtomicI64::new(0),
stopping: AtomicBool::new(false),
terminal: TerminalFailure::default(),
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: Arc::new(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,
capture_consumer: None,
secondary_consumer,
events,
worker: None,
watchdog: None,
started: false,
shutdown: None,
})
}
pub fn enable_capture_tap(&mut self) -> Result<()> {
if self.started {
return Err(Error::InvalidArg(
"capture tap must be enabled before start".into(),
));
}
if self.capture_consumer.is_none() {
let (producer, consumer) = chunk_ring(self.config.ring_capacity_chunks);
*self
.shared
.capture_producer
.lock()
.unwrap_or_else(|e| e.into_inner()) = Some(producer);
self.capture_consumer = Some(consumer);
self.shared.capture_enabled.store(true, Ordering::SeqCst);
}
Ok(())
}
pub fn poll_capture(&mut self) -> Option<AudioChunk> {
let _delivery = self
.shared
.delivery
.lock()
.unwrap_or_else(|e| e.into_inner());
if self.shared.terminal.is_failed() {
return None;
}
self.capture_consumer.as_mut()?.try_pop()
}
}
impl Drop for Stream {
fn drop(&mut self) {
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) -> bool {
let mut clipped = false;
if gain != 1.0 {
for x in data.iter_mut() {
let scaled = *x * gain;
let clamped = scaled.clamp(-1.0, 1.0);
clipped |= !(-1.0..=1.0).contains(&scaled) && !scaled.is_nan();
*x = clamped;
}
}
clipped
}
fn advance_frame_index(counter: &AtomicU64, frames: u64, rate: u32) -> Result<u64> {
if rate == 0 {
return Err(Error::InvalidArg(
"frame index sample rate must be nonzero".into(),
));
}
let previous = counter
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |index| {
let next = index.checked_add(frames)?;
u64::try_from(u128::from(next) * 48_000 / u128::from(rate)).ok()?;
Some(next)
})
.map_err(|_| Error::InvalidState("canonical frame timeline exhausted".into()))?;
u64::try_from(u128::from(previous) * 48_000 / u128::from(rate))
.map_err(|_| Error::InvalidState("canonical frame timeline exhausted".into()))
}
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;
#[cfg(test)]
mod repro_tests;
#[cfg(test)]
#[path = "stream_frame_tests.rs"]
mod frame_tests;
pub(crate) fn validate_chunk_ms(chunk_ms: u32) -> Result<()> {
if chunk_ms == 20 {
Ok(())
} else {
Err(Error::InvalidArg("chunk_ms must be 20".into()))
}
}
fn with_cleanup(primary: Error, cleanup: Option<Error>) -> Error {
match cleanup {
Some(error) => Error::Multiple(ErrorGroup::new(primary, error, Vec::new())),
None => primary,
}
}
fn cleanup_error_since(shared: &SharedState, first: usize) -> Option<Error> {
let mut errors = shared
.cleanup
.lock()
.unwrap_or_else(|e| e.into_inner())
.clone()
.into_iter()
.skip(first);
match (errors.next(), errors.next()) {
(None, _) => None,
(Some(error), None) => Some(error),
(Some(primary), Some(first)) => Some(Error::Multiple(ErrorGroup::new(
primary,
first,
errors.collect(),
))),
}
}
fn is_terminal_kind(error: &Error) -> bool {
matches!(
error.kind(),
ErrorKind::PermissionDenied | ErrorKind::NativeFormatChanged
)
}