use std::collections::{HashMap, VecDeque};
use std::os::raw::{c_char, c_void};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::{Duration, Instant};
use super::bindings::*;
const AMPLITUDE: f64 = 0.25;
#[derive(Clone)]
enum GenSpec {
Silence,
Tone(u32),
File(Arc<Vec<i16>>, u32),
Stream(u32),
}
struct Entry {
version: u64,
spec: GenSpec,
}
static REGISTRY: OnceLock<Mutex<HashMap<String, Entry>>> = OnceLock::new();
fn registry() -> &'static Mutex<HashMap<String, Entry>> {
REGISTRY.get_or_init(|| Mutex::new(HashMap::new()))
}
static AUSRC: OnceLock<usize> = OnceLock::new();
#[derive(Debug, Clone)]
pub struct AudioFrame {
pub samples: Vec<i16>,
pub rate: u32,
}
type StreamQueue = Arc<Mutex<VecDeque<i16>>>;
static STREAM_IN: OnceLock<Mutex<HashMap<String, StreamQueue>>> = OnceLock::new();
fn stream_in() -> &'static Mutex<HashMap<String, StreamQueue>> {
STREAM_IN.get_or_init(|| Mutex::new(HashMap::new()))
}
fn stream_queue(key: &str) -> Option<StreamQueue> {
stream_in()
.lock()
.unwrap_or_else(|e| e.into_inner())
.get(key)
.map(Arc::clone)
}
static RX_TAPS: OnceLock<Mutex<HashMap<String, std::sync::mpsc::Sender<AudioFrame>>>> =
OnceLock::new();
fn rx_taps() -> &'static Mutex<HashMap<String, std::sync::mpsc::Sender<AudioFrame>>> {
RX_TAPS.get_or_init(|| Mutex::new(HashMap::new()))
}
const STREAM_IN_CAP_SAMPLES: usize = 48_000;
pub(super) fn set_generator(key: &str, spec: &str) {
let parsed = parse_spec(spec);
let mut map = registry().lock().unwrap_or_else(|e| e.into_inner());
let version = map.get(key).map(|e| e.version + 1).unwrap_or(0);
map.insert(
key.to_string(),
Entry {
version,
spec: parsed,
},
);
}
pub(super) fn start_audio_stream(key: &str, rate: u32) {
stream_in()
.lock()
.unwrap_or_else(|e| e.into_inner())
.entry(key.to_string())
.or_default();
let mut map = registry().lock().unwrap_or_else(|e| e.into_inner());
let version = map.get(key).map(|e| e.version + 1).unwrap_or(0);
map.insert(
key.to_string(),
Entry {
version,
spec: GenSpec::Stream(rate.max(1)),
},
);
}
pub(super) fn push_audio(key: &str, samples: &[i16]) {
let Some(q) = stream_queue(key) else {
return;
};
let mut q = q.lock().unwrap_or_else(|e| e.into_inner());
q.extend(samples.iter().copied());
while q.len() > STREAM_IN_CAP_SAMPLES {
q.pop_front();
}
}
pub(super) fn subscribe_received_audio(key: &str) -> std::sync::mpsc::Receiver<AudioFrame> {
let (tx, rx) = std::sync::mpsc::channel();
rx_taps()
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(key.to_string(), tx);
rx
}
pub(super) fn init_generator(key: &str) {
let mut map = registry().lock().unwrap_or_else(|e| e.into_inner());
map.entry(key.to_string()).or_insert(Entry {
version: 0,
spec: GenSpec::Silence,
});
}
pub(super) fn remove_generator(key: &str) {
registry()
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(key);
rx_buffers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(key);
tx_buffers()
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(key);
stream_in()
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(key);
rx_taps()
.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(key);
}
fn spec_name(spec: &GenSpec) -> String {
match spec {
GenSpec::Silence => "silence".into(),
GenSpec::Tone(f) => format!("tone({f}Hz)"),
GenSpec::File(s, sr) => format!("file({} samples @{sr}Hz)", s.len()),
GenSpec::Stream(sr) => format!("stream(@{sr}Hz)"),
}
}
fn parse_spec(spec: &str) -> GenSpec {
let (driver, device) = spec.split_once(',').unwrap_or((spec, ""));
match driver {
"ausine" => match device.split(',').next().unwrap_or("").parse::<u32>() {
Ok(f) if (10..=20000).contains(&f) => GenSpec::Tone(f),
_ => {
crate::rlog!(
Warn,
"ringo ausrc: invalid tone freq in '{spec}', using silence"
);
GenSpec::Silence
}
},
"aufile" => match load_wav_mono(device) {
Some((samples, srate)) => GenSpec::File(Arc::new(samples), srate),
None => {
crate::rlog!(
Warn,
"ringo ausrc: failed to load '{device}', using silence"
);
GenSpec::Silence
}
},
_ => GenSpec::Silence,
}
}
fn load_wav_mono(path: &str) -> Option<(Vec<i16>, u32)> {
let data = std::fs::read(path).ok()?;
let pcm = super::sounds::parse_wav(&data)?;
let i16s: Vec<i16> = pcm
.samples
.chunks_exact(2)
.map(|b| i16::from_le_bytes([b[0], b[1]]))
.collect();
let mono = if pcm.channels >= 2 {
i16s.chunks_exact(pcm.channels as usize)
.map(|frame| {
let sum: i32 = frame.iter().take(2).map(|&s| s as i32).sum();
(sum / 2) as i16
})
.collect()
} else {
i16s
};
if mono.is_empty() {
return None;
}
Some((mono, pcm.srate))
}
struct SrcState {
key: String,
run: Arc<AtomicBool>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Drop for SrcState {
fn drop(&mut self) {
crate::rlog!(Info, "ringo ausrc: FREE key={}", self.key);
self.run.store(false, Ordering::Release);
if let Some(t) = self.thread.take() {
let _ = t.join();
}
}
}
struct ReadCb {
rh: ausrc_read_h,
arg: usize,
}
unsafe impl Send for ReadCb {}
unsafe extern "C" fn destructor(arg: *mut c_void) {
let cell = arg as *mut *mut SrcState;
let state = unsafe { *cell };
if !state.is_null() {
drop(unsafe { Box::from_raw(state) });
unsafe { *cell = std::ptr::null_mut() };
}
}
unsafe extern "C" fn alloc_handler(
stp: *mut *mut ausrc_st,
_as: *const ausrc,
prm: *mut ausrc_prm,
dev: *const c_char,
rh: ausrc_read_h,
_errh: ausrc_error_h,
arg: *mut c_void,
) -> i32 {
const EINVAL: i32 = 22;
const ENOMEM: i32 = 12;
const ENOTSUP: i32 = 95;
if stp.is_null() || prm.is_null() || dev.is_null() || rh.is_none() {
return EINVAL;
}
let prm = unsafe { &*prm };
let key = unsafe { std::ffi::CStr::from_ptr(dev) }
.to_string_lossy()
.into_owned();
let srate = prm.srate;
let ptime = if prm.ptime == 0 { 20 } else { prm.ptime };
let fmt = prm.fmt;
if srate == 0 || prm.ch == 0 {
crate::rlog!(Warn, "ringo ausrc: invalid prm srate={srate} ch={}", prm.ch);
return EINVAL;
}
let ch = prm.ch;
if fmt != aufmt::AUFMT_S16LE as i32 && fmt != aufmt::AUFMT_FLOAT as i32 {
crate::rlog!(Warn, "ringo ausrc: unsupported sample format {fmt}");
return ENOTSUP;
}
let cell = unsafe { mem_zalloc(std::mem::size_of::<*mut SrcState>(), Some(destructor)) };
if cell.is_null() {
return ENOMEM;
}
crate::rlog!(
Info,
"ringo ausrc: ALLOC key={key} srate={srate} ch={ch} ptime={ptime} fmt={fmt}"
);
let run = Arc::new(AtomicBool::new(true));
let run_thread = run.clone();
let cb = ReadCb {
rh,
arg: arg as usize,
};
reset_buffer(tx_buffers(), &key, srate);
let key_thread = key.clone();
let thread = match std::thread::Builder::new()
.name("ringo-ausrc".into())
.spawn(move || render_loop(key_thread, srate, ch, ptime, fmt, run_thread, cb))
{
Ok(t) => t,
Err(e) => {
crate::rlog!(Error, "ringo ausrc: spawn render thread failed: {e}");
unsafe { mem_deref(cell) };
return ENOMEM;
}
};
let state = Box::new(SrcState {
key,
run,
thread: Some(thread),
});
unsafe {
*(cell as *mut *mut SrcState) = Box::into_raw(state);
*stp = cell as *mut ausrc_st;
}
0
}
fn render_stream(q: Option<&mut VecDeque<i16>>, pos: &mut f64, step: f64, mono: &mut [i16]) {
let Some(q) = q else {
mono.iter_mut().for_each(|s| *s = 0);
return;
};
for s in mono.iter_mut() {
*s = q.get(*pos as usize).copied().unwrap_or(0);
*pos += step;
}
let consumed = (*pos as usize).min(q.len());
for _ in 0..consumed {
q.pop_front();
}
*pos -= consumed as f64;
if q.is_empty() {
*pos = 0.0; }
}
fn render_loop(
key: String,
srate: u32,
ch: u8,
ptime: u32,
fmt: i32,
run: Arc<AtomicBool>,
cb: ReadCb,
) {
let is_float = fmt == aufmt::AUFMT_FLOAT as i32;
let sample_size = if is_float { 4 } else { 2 };
let frames = (srate as usize * ptime as usize / 1000).max(1);
let sampc = frames * ch as usize;
let mut sampv = vec![0u8; sampc * sample_size];
let mut mono = vec![0i16; frames];
let mut cur_version = u64::MAX;
let mut cur_spec = GenSpec::Silence;
let mut phase = 0.0f64; let mut file_pos = 0.0f64; let mut stream_pos = 0.0f64;
let mut start = Instant::now();
let mut frame_idx: u64 = 0;
while run.load(Ordering::Acquire) {
let target = start + Duration::from_millis(frame_idx * ptime as u64);
let now = Instant::now();
if target > now {
std::thread::sleep(target - now);
} else if now - target > Duration::from_millis(ptime as u64 * 4) {
start = now;
frame_idx = 0;
}
if !run.load(Ordering::Acquire) {
break;
}
let present = {
let map = registry().lock().unwrap_or_else(|e| e.into_inner());
map.get(&key).map(|e| (e.version, e.spec.clone()))
};
match &present {
Some((version, spec)) => {
if *version != cur_version {
cur_version = *version;
cur_spec = spec.clone();
phase = 0.0;
file_pos = 0.0;
stream_pos = 0.0;
crate::rlog!(
Info,
"ringo ausrc: key={key} spec={} (v{version})",
spec_name(spec)
);
}
}
None => {
if cur_version != u64::MAX {
crate::rlog!(Warn, "ringo ausrc: key={key} no registry entry, silence");
cur_version = u64::MAX;
cur_spec = GenSpec::Silence;
}
}
}
match &cur_spec {
GenSpec::Silence => mono.iter_mut().for_each(|s| *s = 0),
GenSpec::Tone(freq) => {
let step = std::f64::consts::TAU * (*freq as f64) / (srate as f64);
for s in mono.iter_mut() {
*s = (phase.sin() * AMPLITUDE * 32767.0) as i16;
phase += step;
if phase >= std::f64::consts::TAU {
phase -= std::f64::consts::TAU;
}
}
}
GenSpec::File(samples, file_srate) => {
let step = *file_srate as f64 / srate as f64;
let len = samples.len() as f64;
for s in mono.iter_mut() {
while file_pos >= len {
file_pos -= len;
}
let i = file_pos as usize;
*s = samples.get(i).copied().unwrap_or(0);
file_pos += step;
}
}
GenSpec::Stream(stream_srate) => {
let step = (*stream_srate as f64 / srate as f64).max(f64::MIN_POSITIVE);
match stream_queue(&key) {
Some(q) => {
let mut q = q.lock().unwrap_or_else(|e| e.into_inner());
render_stream(Some(&mut q), &mut stream_pos, step, &mut mono);
}
None => render_stream(None, &mut stream_pos, step, &mut mono),
}
}
}
if FULL_CAPTURE.load(Ordering::Acquire) {
capture_mono(tx_buffers(), &key, &mono);
}
if is_float {
let out =
unsafe { std::slice::from_raw_parts_mut(sampv.as_mut_ptr() as *mut f32, sampc) };
for (f, &m) in out.chunks_mut(ch as usize).zip(mono.iter()) {
let v = m as f32 / 32768.0;
f.iter_mut().for_each(|x| *x = v);
}
} else {
let out =
unsafe { std::slice::from_raw_parts_mut(sampv.as_mut_ptr() as *mut i16, sampc) };
for (f, &m) in out.chunks_mut(ch as usize).zip(mono.iter()) {
f.iter_mut().for_each(|x| *x = m);
}
}
let mut af: auframe = unsafe { std::mem::zeroed() };
unsafe {
auframe_init(
&mut af,
if is_float {
aufmt::AUFMT_FLOAT
} else {
aufmt::AUFMT_S16LE
},
sampv.as_mut_ptr() as *mut c_void,
sampc,
srate,
ch,
);
}
af.timestamp = frame_idx * ptime as u64 * 1000;
if let Some(rh) = cb.rh {
unsafe { rh(&mut af, cb.arg as *mut c_void) };
}
frame_idx += 1;
}
}
static AUPLAY: OnceLock<usize> = OnceLock::new();
struct PlayState {
run: Arc<AtomicBool>,
thread: Option<std::thread::JoinHandle<()>>,
}
impl Drop for PlayState {
fn drop(&mut self) {
self.run.store(false, Ordering::Release);
if let Some(t) = self.thread.take() {
let _ = t.join();
}
}
}
struct WriteCb {
wh: auplay_write_h,
arg: usize,
}
unsafe impl Send for WriteCb {}
struct AudioBuf {
srate: u32,
samples: VecDeque<i16>,
}
const VERIFY_RETAIN_SECS: usize = 3;
const FULL_RETAIN_SECS: usize = 600;
static FULL_CAPTURE: AtomicBool = AtomicBool::new(false);
static RX_BUFFERS: OnceLock<Mutex<HashMap<String, AudioBuf>>> = OnceLock::new();
static TX_BUFFERS: OnceLock<Mutex<HashMap<String, AudioBuf>>> = OnceLock::new();
fn rx_buffers() -> &'static Mutex<HashMap<String, AudioBuf>> {
RX_BUFFERS.get_or_init(|| Mutex::new(HashMap::new()))
}
fn tx_buffers() -> &'static Mutex<HashMap<String, AudioBuf>> {
TX_BUFFERS.get_or_init(|| Mutex::new(HashMap::new()))
}
pub(super) fn set_full_capture(on: bool) {
FULL_CAPTURE.store(on, Ordering::Release);
}
fn retain_secs() -> usize {
if FULL_CAPTURE.load(Ordering::Acquire) {
FULL_RETAIN_SECS
} else {
VERIFY_RETAIN_SECS
}
}
fn reset_buffer(buffers: &Mutex<HashMap<String, AudioBuf>>, key: &str, srate: u32) {
buffers.lock().unwrap_or_else(|e| e.into_inner()).insert(
key.to_string(),
AudioBuf {
srate,
samples: VecDeque::new(),
},
);
}
pub(super) fn received_window(key: &str) -> Option<(Vec<i16>, u32)> {
buffer_window(rx_buffers(), key)
}
pub(super) fn sent_window(key: &str) -> Option<(Vec<i16>, u32)> {
buffer_window(tx_buffers(), key)
}
fn buffer_window(buffers: &Mutex<HashMap<String, AudioBuf>>, key: &str) -> Option<(Vec<i16>, u32)> {
let map = buffers.lock().unwrap_or_else(|e| e.into_inner());
let buf = map.get(key)?;
Some((buf.samples.iter().copied().collect(), buf.srate))
}
fn capture_mono(buffers: &Mutex<HashMap<String, AudioBuf>>, key: &str, samples: &[i16]) {
let mut map = buffers.lock().unwrap_or_else(|e| e.into_inner());
let Some(buf) = map.get_mut(key) else {
return;
};
buf.samples.extend(samples.iter().copied());
let cap = (buf.srate as usize * retain_secs()).max(1);
while buf.samples.len() > cap {
buf.samples.pop_front();
}
}
fn capture_rx(key: &str, sampv: &[u8], is_float: bool, ch: usize, srate: u32) {
let ch = ch.max(1);
let mono: Vec<i16> = if is_float {
let f =
unsafe { std::slice::from_raw_parts(sampv.as_ptr() as *const f32, sampv.len() / 4) };
f.chunks_exact(ch)
.map(|fr| (fr[0] * 32767.0) as i16)
.collect()
} else {
let s =
unsafe { std::slice::from_raw_parts(sampv.as_ptr() as *const i16, sampv.len() / 2) };
s.chunks_exact(ch).map(|fr| fr[0]).collect()
};
capture_mono(rx_buffers(), key, &mono);
let mut taps = rx_taps().lock().unwrap_or_else(|e| e.into_inner());
if let Some(tx) = taps.get(key) {
let frame = AudioFrame {
samples: mono,
rate: srate,
};
if tx.send(frame).is_err() {
taps.remove(key);
}
}
}
unsafe extern "C" fn play_destructor(arg: *mut c_void) {
let cell = arg as *mut *mut PlayState;
let state = unsafe { *cell };
if !state.is_null() {
drop(unsafe { Box::from_raw(state) });
unsafe { *cell = std::ptr::null_mut() };
}
}
unsafe extern "C" fn play_alloc_handler(
stp: *mut *mut auplay_st,
_ap: *const auplay,
prm: *mut auplay_prm,
dev: *const c_char,
wh: auplay_write_h,
arg: *mut c_void,
) -> i32 {
const EINVAL: i32 = 22;
const ENOMEM: i32 = 12;
if stp.is_null() || prm.is_null() || wh.is_none() {
return EINVAL;
}
let prm = unsafe { &*prm };
let srate = prm.srate;
let ptime = if prm.ptime == 0 { 20 } else { prm.ptime };
let fmt = prm.fmt;
if srate == 0 || prm.ch == 0 {
crate::rlog!(
Warn,
"ringo auplay: invalid prm srate={srate} ch={}",
prm.ch
);
return EINVAL;
}
let ch = prm.ch;
let key = if dev.is_null() {
String::new()
} else {
unsafe { std::ffi::CStr::from_ptr(dev) }
.to_string_lossy()
.into_owned()
};
if !key.is_empty() {
reset_buffer(rx_buffers(), &key, srate);
}
let cell = unsafe { mem_zalloc(std::mem::size_of::<*mut PlayState>(), Some(play_destructor)) };
if cell.is_null() {
return ENOMEM;
}
let run = Arc::new(AtomicBool::new(true));
let run_thread = run.clone();
let cb = WriteCb {
wh,
arg: arg as usize,
};
let thread = match std::thread::Builder::new()
.name("ringo-auplay".into())
.spawn(move || play_loop(key, srate, ch, ptime, fmt, run_thread, cb))
{
Ok(t) => t,
Err(e) => {
crate::rlog!(Error, "ringo auplay: spawn play thread failed: {e}");
unsafe { mem_deref(cell) };
return ENOMEM;
}
};
let state = Box::new(PlayState {
run,
thread: Some(thread),
});
unsafe {
*(cell as *mut *mut PlayState) = Box::into_raw(state);
*stp = cell as *mut auplay_st;
}
0
}
fn play_loop(
key: String,
srate: u32,
ch: u8,
ptime: u32,
fmt: i32,
run: Arc<AtomicBool>,
cb: WriteCb,
) {
let is_float = fmt == aufmt::AUFMT_FLOAT as i32;
let sample_size = if is_float { 4 } else { 2 };
let frames = (srate as usize * ptime as usize / 1000).max(1);
let sampc = frames * ch as usize;
let mut sampv = vec![0u8; sampc * sample_size];
let af_fmt = if is_float {
aufmt::AUFMT_FLOAT
} else {
aufmt::AUFMT_S16LE
};
let mut start = Instant::now();
let mut frame_idx: u64 = 0;
while run.load(Ordering::Acquire) {
let target = start + Duration::from_millis(frame_idx * ptime as u64);
let now = Instant::now();
if target > now {
std::thread::sleep(target - now);
} else if now - target > Duration::from_millis(ptime as u64 * 4) {
start = now;
frame_idx = 0;
}
if !run.load(Ordering::Acquire) {
break;
}
let mut af: auframe = unsafe { std::mem::zeroed() };
unsafe {
auframe_init(
&mut af,
af_fmt,
sampv.as_mut_ptr() as *mut c_void,
sampc,
srate,
ch,
);
}
af.timestamp = frame_idx * ptime as u64 * 1000;
if let Some(wh) = cb.wh {
unsafe { wh(&mut af, cb.arg as *mut c_void) };
if !key.is_empty() {
capture_rx(&key, &sampv, is_float, ch as usize, srate);
}
}
frame_idx += 1;
}
}
pub(super) fn register_module() -> Result<(), String> {
let mut asp: *mut ausrc = std::ptr::null_mut();
let rc = unsafe {
ausrc_register(
&mut asp,
baresip_ausrcl(),
c"ringo".as_ptr(),
Some(alloc_handler),
)
};
if rc != 0 {
return Err(format!("ausrc_register(ringo) failed (rc={rc})"));
}
let _ = AUSRC.set(asp as usize);
let mut pp: *mut auplay = std::ptr::null_mut();
let rc = unsafe {
auplay_register(
&mut pp,
baresip_auplayl(),
c"ringo".as_ptr(),
Some(play_alloc_handler),
)
};
if rc != 0 {
return Err(format!("auplay_register(ringo) failed (rc={rc})"));
}
let _ = AUPLAY.set(pp as usize);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn render_stream_same_rate_drains_in_order() {
let mut q: VecDeque<i16> = (1..=6).collect();
let mut pos = 0.0;
let mut out = [0i16; 4];
render_stream(Some(&mut q), &mut pos, 1.0, &mut out);
assert_eq!(out, [1, 2, 3, 4]);
assert_eq!(q.iter().copied().collect::<Vec<_>>(), vec![5, 6]); let mut out2 = [0i16; 2];
render_stream(Some(&mut q), &mut pos, 1.0, &mut out2);
assert_eq!(out2, [5, 6]);
assert!(q.is_empty());
}
#[test]
fn render_stream_underrun_is_silence() {
let mut q: VecDeque<i16> = VecDeque::from(vec![7, 8]);
let mut pos = 0.0;
let mut out = [0i16; 4];
render_stream(Some(&mut q), &mut pos, 1.0, &mut out);
assert_eq!(out, [7, 8, 0, 0]); assert!(q.is_empty());
assert_eq!(pos, 0.0); let mut out2 = [9i16; 3];
render_stream(None, &mut pos, 1.0, &mut out2);
assert_eq!(out2, [0, 0, 0]);
}
#[test]
fn render_stream_downsamples_2to1() {
let mut q: VecDeque<i16> = (0..8).collect();
let mut pos = 0.0;
let mut out = [0i16; 4];
render_stream(Some(&mut q), &mut pos, 2.0, &mut out);
assert_eq!(out, [0, 2, 4, 6]);
assert!(q.is_empty()); }
#[test]
fn push_audio_is_noop_before_start_then_caps() {
let key = "ringo-test-stream-cap"; push_audio(key, &[1, 2, 3]);
assert!(stream_queue(key).is_none());
start_audio_stream(key, 8000);
push_audio(key, &vec![7i16; STREAM_IN_CAP_SAMPLES + 100]);
let q = stream_queue(key).expect("queue exists after start");
assert_eq!(
q.lock().unwrap_or_else(|e| e.into_inner()).len(),
STREAM_IN_CAP_SAMPLES
);
remove_generator(key);
assert!(stream_queue(key).is_none());
}
}