use crate::audio::ring_buffer::RingBuffer;
use std::collections::VecDeque;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use tokio::sync::mpsc;
const LOOKBACK_MS: u64 = 300;
pub fn spawn_audio_tee(
mut input: mpsc::Receiver<Vec<i16>>,
ring: Arc<Mutex<RingBuffer>>,
pause: Arc<AtomicBool>,
pause_forwarding: bool,
) -> mpsc::Receiver<Vec<i16>> {
let (tx, rx) = mpsc::channel(super::CHANNEL_CAPACITY);
let chunk_duration_ms = super::CHUNK_DURATION_MS.max(1);
let max_lookback_chunks = (LOOKBACK_MS / chunk_duration_ms) as usize;
tokio::spawn(async move {
let mut lookback: VecDeque<Vec<i16>> = VecDeque::new();
let mut was_paused = false;
while let Some(chunk) = input.recv().await {
let f32_samples: Vec<f32> = chunk.iter().map(|&s| s as f32 / 32768.0).collect();
if let Ok(mut guard) = ring.lock() {
guard.push(&f32_samples);
}
let paused = pause_forwarding && pause.load(Ordering::Relaxed);
if paused {
lookback.push_back(chunk);
while lookback.len() > max_lookback_chunks {
lookback.pop_front();
}
was_paused = true;
} else {
if was_paused {
for buffered in lookback.drain(..) {
if tx.send(buffered).await.is_err() {
return;
}
}
was_paused = false;
}
if tx.send(chunk).await.is_err() {
break; }
}
}
});
rx
}