use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc;
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use flexaudio_core::backend::{CaptureBackend, RawSink};
use flexaudio_core::types::{Error, ProcessMode, Result};
use crate::common::{translate_pid_to_object, FALLBACK_FORMAT};
use crate::system::run_tap_thread;
use crate::tap::TapKind;
pub struct MacProcessBackend {
target_pid: u32,
mode: ProcessMode,
stop_flag: Arc<AtomicBool>,
handle: Option<JoinHandle<()>>,
native: (u32, u16),
}
impl MacProcessBackend {
pub fn new(target_pid: u32, mode: ProcessMode) -> Self {
Self {
target_pid,
mode,
stop_flag: Arc::new(AtomicBool::new(false)),
handle: None,
native: FALLBACK_FORMAT,
}
}
pub fn target_pid(&self) -> u32 {
self.target_pid
}
pub fn mode(&self) -> ProcessMode {
self.mode
}
}
impl CaptureBackend for MacProcessBackend {
fn native_format(&self) -> (u32, u16) {
self.native
}
fn start(&mut self, sink: RawSink) -> Result<()> {
if self.handle.is_some() {
return Ok(());
}
crate::version::ensure_process_tap_supported()?;
self.stop_flag.store(false, Ordering::SeqCst);
let stop_flag = self.stop_flag.clone();
let (ready_tx, ready_rx) = mpsc::channel::<Result<()>>();
let target_pid = self.target_pid;
let mode = self.mode;
let handle = thread::Builder::new()
.name("flexaudio-macos-process".into())
.spawn(move || {
let kind = match translate_pid_to_object(target_pid as i32) {
Ok(0) => {
let _ = ready_tx.send(Err(Error::DeviceNotFound));
return;
}
Ok(object_id) => match mode {
ProcessMode::Include => TapKind::IncludeProcesses(vec![object_id]),
ProcessMode::Exclude => TapKind::ExcludeProcesses(vec![object_id]),
},
Err(e) => {
let _ = ready_tx.send(Err(e));
return;
}
};
run_tap_thread(kind, sink, stop_flag, ready_tx);
})
.map_err(|e| Error::Backend(format!("spawn macos process thread: {e}")))?;
match ready_rx.recv() {
Ok(Ok(())) => {
self.handle = Some(handle);
Ok(())
}
Ok(Err(e)) => {
self.stop_flag.store(false, Ordering::SeqCst);
let _ = handle.join();
Err(e)
}
Err(_) => {
self.stop_flag.store(false, Ordering::SeqCst);
let _ = handle.join();
Err(Error::Backend(
"macos process thread exited before reporting readiness".into(),
))
}
}
}
fn stop(&mut self) {
self.stop_flag.store(true, Ordering::SeqCst);
if let Some(h) = self.handle.take() {
let _ = h.join();
}
}
}
impl Drop for MacProcessBackend {
fn drop(&mut self) {
self.stop();
}
}
#[cfg(test)]
mod tests {
use super::*;
use flexaudio_core::raw_ring;
#[test]
fn new_and_native_format_do_not_panic() {
let backend = MacProcessBackend::new(1234, ProcessMode::Include);
let (rate, channels) = backend.native_format();
assert!(rate > 0);
assert!(channels > 0);
assert_eq!(backend.target_pid(), 1234);
assert_eq!(backend.mode(), ProcessMode::Include);
}
#[test]
fn start_then_stop_tolerates_missing_target() {
let mut backend = MacProcessBackend::new(0xFFFF_FFFE, ProcessMode::Include);
let (rate, channels) = backend.native_format();
let cap = (rate as usize * channels as usize).max(1);
let (prod, _cons) = raw_ring(cap);
let sink = RawSink::new(prod, rate, channels);
match backend.start(sink) {
Ok(()) => {
backend.stop();
backend.stop();
}
Err(_e) => { }
}
}
}