use super::types::SampleType;
use ringbuf::traits::{Consumer, Observer, Producer};
use ringbuf::{HeapCons, HeapProd};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread::{self, JoinHandle};
use std::time::Duration;
const IDLE_POLL: Duration = Duration::from_micros(500);
const MAX_SCHEDULE_AHEAD: usize = 1 << 20;
#[derive(Debug, Default)]
pub struct ClearSignal {
epoch: AtomicUsize,
ack: AtomicUsize,
}
impl ClearSignal {
pub fn request(&self) -> usize {
self.epoch.fetch_add(1, Ordering::Release) + 1
}
pub fn epoch(&self) -> usize {
self.epoch.load(Ordering::Acquire)
}
pub fn acked(&self) -> usize {
self.ack.load(Ordering::Acquire)
}
fn ack(&self, epoch: usize) {
self.ack.store(epoch, Ordering::Release);
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct MixCommand<S: SampleType> {
pub index: usize,
pub samples: Vec<S>,
}
pub struct Mixer<S: SampleType> {
commands: HeapCons<MixCommand<S>>,
output: HeapProd<S>,
mix: Vec<S>,
cursor: usize,
late: usize,
dropped: usize,
clear: Arc<ClearSignal>,
cleared: usize,
}
impl<S: SampleType> Mixer<S> {
pub fn new(commands: HeapCons<MixCommand<S>>, output: HeapProd<S>) -> Self {
Mixer {
commands,
output,
mix: Vec::new(),
cursor: 0,
late: 0,
dropped: 0,
clear: Arc::new(ClearSignal::default()),
cleared: 0,
}
}
pub fn clear_signal(&self) -> Arc<ClearSignal> {
Arc::clone(&self.clear)
}
fn clear_pending(&mut self) {
self.commands.clear();
self.mix.clear();
}
fn apply(&mut self, command: MixCommand<S>) {
let MixCommand { index, samples } = command;
let skip = self.cursor.saturating_sub(index);
if skip >= samples.len() {
self.late += samples.len();
return;
}
self.late += skip;
let start = (index + skip) - self.cursor;
let end = start + (samples.len() - skip);
if end > MAX_SCHEDULE_AHEAD {
self.dropped += samples.len() - skip;
return;
}
if self.mix.len() < end {
self.mix.resize(end, S::SILENCE);
}
for (slot, sample) in self.mix[start..end].iter_mut().zip(&samples[skip..]) {
*slot = slot.mix(*sample);
}
}
fn flush(&mut self) -> usize {
let ready = self.mix.len().min(self.output.vacant_len());
if ready == 0 {
return 0;
}
let pushed = self.output.push_slice(&self.mix[..ready]);
self.mix.drain(..pushed);
self.cursor += pushed;
pushed
}
pub fn tick(&mut self) -> bool {
let mut worked = false;
let epoch = self.clear.epoch();
if epoch != self.cleared {
self.clear_pending();
self.cleared = epoch;
self.clear.ack(epoch);
worked = true;
}
while let Some(command) = self.commands.try_pop() {
self.apply(command);
worked = true;
}
self.flush() > 0 || worked
}
pub fn run(mut self, running: Arc<AtomicBool>) {
while running.load(Ordering::Relaxed) {
if !self.tick() {
thread::sleep(IDLE_POLL);
}
}
}
pub fn cursor(&self) -> usize {
self.cursor
}
pub fn late_samples(&self) -> usize {
self.late
}
pub fn dropped_samples(&self) -> usize {
self.dropped
}
}
pub struct MixerHandle {
running: Arc<AtomicBool>,
clear: Arc<ClearSignal>,
thread: Option<JoinHandle<()>>,
}
impl MixerHandle {
pub fn spawn<S: SampleType>(mixer: Mixer<S>) -> Self {
let running = Arc::new(AtomicBool::new(true));
let flag = Arc::clone(&running);
let clear = mixer.clear_signal();
let thread = thread::Builder::new()
.name("atome_mixer".to_string())
.spawn(move || mixer.run(flag))
.expect("failed to spawn mixer thread");
MixerHandle {
running,
clear,
thread: Some(thread),
}
}
pub fn clear_samples(&self) {
self.clear.request();
}
pub fn clear_signal(&self) -> Arc<ClearSignal> {
Arc::clone(&self.clear)
}
}
impl Drop for MixerHandle {
fn drop(&mut self) {
self.running.store(false, Ordering::Relaxed);
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}