#[cfg(feature = "audio")]
use anyhow::Context;
use anyhow::{anyhow, Result};
#[cfg(feature = "audio")]
use cpal::traits::{DeviceTrait, HostTrait, StreamTrait};
#[cfg(feature = "audio")]
use std::sync::atomic::AtomicU64;
use std::sync::{
atomic::{AtomicBool, Ordering},
mpsc, Arc,
};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
use super::codec::SAMPLES_PER_FRAME;
#[cfg(feature = "audio")]
use super::codec::SAMPLE_RATE;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum AudioSource {
#[default]
Mic,
System,
Both,
None_,
}
impl AudioSource {
pub fn parse(s: &str) -> Result<Self> {
match s.to_ascii_lowercase().as_str() {
"mic" | "microphone" => Ok(Self::Mic),
"system" | "loopback" | "output" => Ok(Self::System),
"both" | "mic+system" | "mix" => Ok(Self::Both),
"none" | "off" | "silence" => Ok(Self::None_),
other => anyhow::bail!("--audio-source must be mic|system|both|none, got '{other}'"),
}
}
}
#[cfg(feature = "audio")]
fn is_loopback_name(name: &str) -> bool {
let n = name.to_ascii_lowercase();
[
"monitor",
"loopback",
"stereo mix",
"stereomix",
"what u hear",
"wave out",
"blackhole",
"vb-cable",
"vb cable",
"vbcable",
"background music",
"zoom",
"hue sync",
"huesync",
"serato virtual",
"parrot",
"solstice",
"virtual audio",
"virtual mic",
"teams audio",
"ecamm",
]
.iter()
.any(|k| n.contains(k))
}
#[cfg(feature = "audio")]
pub fn find_loopback() -> Option<cpal::Device> {
let host = cpal::default_host();
let devices = host.input_devices().ok()?;
devices
.filter_map(|d| d.name().ok().map(|n| (n, d)))
.find(|(n, _)| is_loopback_name(n))
.map(|(_, d)| d)
}
#[cfg(not(feature = "audio"))]
pub fn list_input_devices() -> Vec<(String, bool)> {
Vec::new()
}
#[cfg(feature = "audio")]
pub fn list_input_devices() -> Vec<(String, bool)> {
let host = cpal::default_host();
let Ok(devices) = host.input_devices() else {
return Vec::new();
};
devices
.filter_map(|d| d.name().ok())
.map(|n| {
let loopback = is_loopback_name(&n);
(n, loopback)
})
.collect()
}
pub fn open_source(which: AudioSource) -> Result<Box<dyn SystemAudio>> {
match which {
AudioSource::None_ => Ok(Box::new(NullSource::new())),
#[cfg(not(feature = "audio"))]
_ => {
let _ = which;
Ok(Box::new(NullSource::new()))
}
#[cfg(feature = "audio")]
AudioSource::Mic => match cpal::default_host().default_input_device() {
Some(device) => Ok(Box::new(MicrophoneSource::from_device(device)?)),
None => Ok(Box::new(NullSource::new())),
},
#[cfg(feature = "audio")]
AudioSource::System => match find_loopback() {
Some(device) => Ok(Box::new(MicrophoneSource::from_device(device)?)),
None => {
tracing::warn!(
"no system-audio loopback device found; capturing the microphone instead (WASAPI loopback / PipeWire monitor / ScreenCaptureKit backends not yet implemented)"
);
open_source(AudioSource::Mic)
}
},
#[cfg(feature = "audio")]
AudioSource::Both => match find_loopback() {
Some(device) => Ok(Box::new(MixedSource::new(device)?)),
None => {
tracing::warn!(
"no system-audio loopback device found; --audio-source both captures the microphone alone"
);
open_source(AudioSource::Mic)
}
},
}
}
pub fn default_source() -> Result<Box<dyn SystemAudio>> {
open_source(AudioSource::Mic)
}
#[cfg(feature = "audio")]
pub struct MixedSource {
mic_name: Option<String>,
system_name: String,
running: Option<Arc<AtomicBool>>,
worker: Option<JoinHandle<()>>,
}
#[cfg(feature = "audio")]
impl MixedSource {
pub fn new(system_device: cpal::Device) -> Result<Self> {
let system_name = system_device
.name()
.context("reading loopback device name")?;
let mic_name = cpal::default_host()
.default_input_device()
.and_then(|d| d.name().ok());
Ok(Self {
mic_name,
system_name,
running: None,
worker: None,
})
}
}
#[cfg(feature = "audio")]
fn open_named(name: &str) -> Result<MicrophoneSource> {
let host = cpal::default_host();
let devices = host
.input_devices()
.context("enumerating audio input devices")?;
for device in devices {
if device.name().as_deref().unwrap_or("") == name {
return MicrophoneSource::from_device(device);
}
}
anyhow::bail!("audio device '{name}' disappeared")
}
pub fn mix_frames(a: &[f32], b: &[f32]) -> Vec<f32> {
a.iter()
.zip(b.iter())
.map(|(x, y)| (x + y).clamp(-1.0, 1.0))
.collect()
}
#[cfg(feature = "audio")]
impl SystemAudio for MixedSource {
fn device_name(&self) -> &str {
"mic+system mix"
}
fn start(&mut self) -> Result<mpsc::Receiver<AudioFrame>> {
if self.running.is_some() {
return Err(anyhow!("mixed source is already running"));
}
if self.running.is_some() {
return Err(anyhow!("mixed source is already running"));
}
let mic_name = self.mic_name.clone();
let system_name = self.system_name.clone();
let (sender, receiver) = mpsc::sync_channel(CAPTURE_QUEUE_FRAMES);
let running = Arc::new(AtomicBool::new(true));
let worker_running = Arc::clone(&running);
let worker = thread::spawn(move || {
let Ok(mut system) = open_named(&system_name) else {
return;
};
let mut mic = mic_name.as_deref().and_then(|n| open_named(n).ok());
let Ok(system_frames) = system.start() else {
return;
};
let mic_frames = match mic.as_mut() {
Some(m) => m.start().ok(),
None => None,
};
while worker_running.load(Ordering::Relaxed) {
let mic_frame = match &mic_frames {
Some(rx) => rx.recv_timeout(Duration::from_millis(5)).ok(),
None => {
thread::sleep(Duration::from_millis(5));
None
}
};
let sys_frame = system_frames.recv_timeout(Duration::from_millis(5)).ok();
match (mic_frame, sys_frame) {
(None, None) => continue,
(m, s) => {
let pcm = match (m, s) {
(Some(m), Some(s)) => mix_frames(&m.pcm, &s.pcm),
(Some(m), None) => m.pcm,
(None, Some(s)) => s.pcm,
(None, None) => unreachable!(),
};
if matches!(
sender.try_send(AudioFrame {
pcm,
capture_time: Instant::now(),
}),
Err(mpsc::TrySendError::Disconnected(_))
) {
break;
}
}
}
}
});
self.running = Some(running);
self.worker = Some(worker);
Ok(receiver)
}
fn stop(&mut self) -> Result<()> {
if let Some(running) = self.running.take() {
running.store(false, Ordering::Relaxed);
}
if let Some(worker) = self.worker.take() {
worker
.join()
.map_err(|_| anyhow!("mixed source worker panicked"))?;
}
Ok(())
}
}
#[cfg(feature = "audio")]
impl Drop for MixedSource {
fn drop(&mut self) {
let _ = self.stop();
}
}
#[cfg(not(feature = "audio"))]
pub fn default_device_name() -> Option<String> {
None
}
#[cfg(feature = "audio")]
pub fn default_device_name() -> Option<String> {
let host = cpal::default_host();
host.default_input_device()
.and_then(|d| d.name().ok())
.map(|n| n.to_string())
}
#[derive(Debug)]
pub struct AudioFrame {
pub pcm: Vec<f32>,
pub capture_time: Instant,
}
pub trait SystemAudio {
fn device_name(&self) -> &str;
fn start(&mut self) -> Result<mpsc::Receiver<AudioFrame>>;
fn stop(&mut self) -> Result<()>;
}
const CAPTURE_QUEUE_FRAMES: usize = 8;
#[cfg(feature = "audio")]
pub struct MicrophoneSource {
device_name: String,
device: cpal::Device,
config: cpal::StreamConfig,
sample_format: cpal::SampleFormat,
stream: Option<cpal::Stream>,
dropped_frames: Arc<AtomicU64>,
}
#[cfg(feature = "audio")]
impl MicrophoneSource {
pub fn new() -> Result<Self> {
let device = cpal::default_host()
.default_input_device()
.ok_or_else(|| anyhow!("no default audio input device is available"))?;
Self::from_device(device)
}
pub fn from_device(device: cpal::Device) -> Result<Self> {
let device_name = device.name().context("reading audio input device name")?;
let config = device
.supported_input_configs()
.with_context(|| format!("reading input configurations for {device_name}"))?
.find(|range| {
range.min_sample_rate().0 <= SAMPLE_RATE
&& SAMPLE_RATE <= range.max_sample_rate().0
&& range.channels() > 0
&& matches!(
range.sample_format(),
cpal::SampleFormat::F32 | cpal::SampleFormat::I16 | cpal::SampleFormat::U16
)
})
.ok_or_else(|| {
anyhow!(
"audio input device {device_name} has no supported 48 kHz F32, I16, or U16 configuration"
)
})?;
let sample_format = config.sample_format();
let config = config
.with_sample_rate(cpal::SampleRate(SAMPLE_RATE))
.config();
Ok(Self {
device_name,
device,
config,
sample_format,
stream: None,
dropped_frames: Arc::new(AtomicU64::new(0)),
})
}
pub fn dropped_frames(&self) -> u64 {
self.dropped_frames.load(Ordering::Relaxed)
}
fn build_stream(&self, sender: mpsc::SyncSender<AudioFrame>) -> Result<cpal::Stream> {
let channels = usize::from(self.config.channels);
let dropped_frames = Arc::clone(&self.dropped_frames);
let error_callback = |error| tracing::warn!(%error, "audio input stream error");
match self.sample_format {
cpal::SampleFormat::F32 => {
let mut assembler = FrameAssembler::new(channels);
self.device
.build_input_stream(
&self.config,
move |data: &[f32], _| {
assembler.push(data, |sample| sample, &sender, &dropped_frames);
},
error_callback,
None,
)
.context("building F32 audio input stream")
}
cpal::SampleFormat::I16 => {
let mut assembler = FrameAssembler::new(channels);
self.device
.build_input_stream(
&self.config,
move |data: &[i16], _| {
assembler.push(
data,
|sample| f32::from(sample) / 32_768.0,
&sender,
&dropped_frames,
);
},
error_callback,
None,
)
.context("building I16 audio input stream")
}
cpal::SampleFormat::U16 => {
let mut assembler = FrameAssembler::new(channels);
self.device
.build_input_stream(
&self.config,
move |data: &[u16], _| {
assembler.push(
data,
|sample| (f32::from(sample) - 32_768.0) / 32_768.0,
&sender,
&dropped_frames,
);
},
error_callback,
None,
)
.context("building U16 audio input stream")
}
format => Err(anyhow!("unsupported audio input sample format {format}")),
}
}
}
#[cfg(feature = "audio")]
impl SystemAudio for MicrophoneSource {
fn device_name(&self) -> &str {
&self.device_name
}
fn start(&mut self) -> Result<mpsc::Receiver<AudioFrame>> {
if self.stream.is_some() {
return Err(anyhow!(
"audio input {0} is already running",
self.device_name
));
}
let (sender, receiver) = mpsc::sync_channel(CAPTURE_QUEUE_FRAMES);
let stream = self.build_stream(sender)?;
stream
.play()
.with_context(|| format!("starting audio input {}", self.device_name))?;
self.stream = Some(stream);
Ok(receiver)
}
fn stop(&mut self) -> Result<()> {
self.stream.take();
Ok(())
}
}
#[cfg(feature = "audio")]
struct FrameAssembler {
channels: usize,
pcm: Vec<f32>,
}
#[cfg(feature = "audio")]
impl FrameAssembler {
fn new(channels: usize) -> Self {
Self {
channels,
pcm: Vec::with_capacity(SAMPLES_PER_FRAME),
}
}
fn push<T: Copy>(
&mut self,
data: &[T],
mut convert: impl FnMut(T) -> f32,
sender: &mpsc::SyncSender<AudioFrame>,
dropped_frames: &AtomicU64,
) {
for input_frame in data.chunks_exact(self.channels) {
let left = convert(input_frame[0]);
let right = if self.channels == 1 {
left
} else {
convert(input_frame[1])
};
self.pcm.extend([left, right]);
if self.pcm.len() == SAMPLES_PER_FRAME {
let pcm = std::mem::replace(&mut self.pcm, Vec::with_capacity(SAMPLES_PER_FRAME));
let frame = AudioFrame {
pcm,
capture_time: Instant::now(),
};
if matches!(sender.try_send(frame), Err(mpsc::TrySendError::Full(_))) {
dropped_frames.fetch_add(1, Ordering::Relaxed);
}
}
}
}
}
pub struct NullSource {
running: Option<Arc<AtomicBool>>,
worker: Option<JoinHandle<()>>,
}
impl NullSource {
pub fn new() -> Self {
Self {
running: None,
worker: None,
}
}
}
impl Default for NullSource {
fn default() -> Self {
Self::new()
}
}
impl SystemAudio for NullSource {
fn device_name(&self) -> &str {
"silence"
}
fn start(&mut self) -> Result<mpsc::Receiver<AudioFrame>> {
if self.running.is_some() {
return Err(anyhow!("silence source is already running"));
}
let (sender, receiver) = mpsc::sync_channel(CAPTURE_QUEUE_FRAMES);
let running = Arc::new(AtomicBool::new(true));
let worker_running = Arc::clone(&running);
let worker = thread::spawn(move || {
while worker_running.load(Ordering::Relaxed) {
let frame = AudioFrame {
pcm: vec![0.0; SAMPLES_PER_FRAME],
capture_time: Instant::now(),
};
if matches!(
sender.try_send(frame),
Err(mpsc::TrySendError::Disconnected(_))
) {
break;
}
thread::sleep(Duration::from_millis(20));
}
});
self.running = Some(running);
self.worker = Some(worker);
Ok(receiver)
}
fn stop(&mut self) -> Result<()> {
if let Some(running) = self.running.take() {
running.store(false, Ordering::Relaxed);
}
if let Some(worker) = self.worker.take() {
worker
.join()
.map_err(|_| anyhow!("silence source worker panicked"))?;
}
Ok(())
}
}
impl Drop for NullSource {
fn drop(&mut self) {
let _ = self.stop();
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
#[test]
fn null_source_emits_stereo_frames_with_monotonic_capture_times() {
let mut source = NullSource::new();
let receiver = source.start().expect("starts silence source");
let first = receiver
.recv_timeout(Duration::from_millis(100))
.expect("first silence frame");
let second = receiver
.recv_timeout(Duration::from_millis(100))
.expect("second silence frame");
assert_eq!(first.pcm.len(), SAMPLES_PER_FRAME);
assert_eq!(second.pcm.len(), SAMPLES_PER_FRAME);
assert!(second.capture_time >= first.capture_time);
source.stop().expect("stops silence source");
}
#[test]
fn audio_source_parses_and_rejects_garbage() {
use AudioSource::*;
assert_eq!(AudioSource::parse("mic").unwrap(), Mic);
assert_eq!(AudioSource::parse("MICROPHONE").unwrap(), Mic);
assert_eq!(AudioSource::parse("system").unwrap(), System);
assert_eq!(AudioSource::parse("loopback").unwrap(), System);
assert_eq!(AudioSource::parse("both").unwrap(), Both);
assert_eq!(AudioSource::parse("mix").unwrap(), Both);
assert_eq!(AudioSource::parse("none").unwrap(), None_);
assert_eq!(AudioSource::parse("off").unwrap(), None_);
assert!(AudioSource::parse("surround").is_err());
assert!(AudioSource::parse("").is_err());
}
#[cfg(feature = "audio")]
#[test]
fn loopback_names_match_monitors_and_mixes() {
for yes in [
"Monitor of Built-in Audio",
"Stereo Mix (Realtek)",
"loopback PCM",
"What U Hear",
] {
assert!(is_loopback_name(yes), "{yes:?} should count as loopback");
}
for yes in [
"BlackHole 2ch",
"VB-Cable",
"Background Music",
"Background Music (UI Sounds)",
"ZoomAudioDevice",
"Hue Sync Audio",
"Serato Virtual Audio",
"Microsoft Teams Audio",
] {
assert!(is_loopback_name(yes), "{yes:?} should count as loopback");
}
for no in [
"Built-in Microphone",
"USB Headset",
"MacBook Pro Speakers",
"MacBook Pro Microphone",
] {
assert!(!is_loopback_name(no), "{no:?} should not count as loopback");
}
}
#[test]
fn mix_clips_rather_than_wrapping() {
assert_eq!(mix_frames(&[0.5, -0.5], &[0.5, -0.5]), vec![1.0, -1.0]);
assert_eq!(mix_frames(&[0.9, 0.9], &[0.9, 0.9]), vec![1.0, 1.0]);
assert_eq!(mix_frames(&[0.1], &[0.2, 0.3]), vec![0.3]);
}
}