//! Reproducción del audio sintetizado. //! //! El TTS produce a ráfagas —un bloque de codec cada vez— y el altavoz consume //! a ritmo constante, así que entre ambos hay un anillo. La retrollamada sólo //! vacía el anillo; quien sintetiza sólo lo llena. Cortar la respuesta es //! entonces una operación trivial y sin condiciones de carrera: se vacía el //! anillo y la voz calla en el siguiente bloque de audio. use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Condvar, Mutex}; use std::time::Duration; use cpal::traits::{DeviceTrait, HostTrait, StreamTrait}; use cpal::SampleFormat; use asist_core::config::AudioConfig; use asist_core::error::{Error, Result}; use crate::{describe, to_mono_at, InputFormat, TTS_SAMPLE_RATE}; /// Estado compartido entre quien sintetiza y la retrollamada del altavoz. struct Ring { samples: Mutex>, /// Despierta a quien espera a que se vacíe la cola. drained: Condvar, /// Muestras servidas al altavoz desde que arrancó. Sirve para saber si ya /// ha sonado algo de verdad, que es la latencia que percibe el usuario. played: AtomicU64, /// Hay audio en curso (queda cola o se está alimentando). active: AtomicBool, gain: Mutex, } impl Ring { fn new(gain: f32) -> Self { Self { samples: Mutex::new(std::collections::VecDeque::new()), drained: Condvar::new(), played: AtomicU64::new(0), active: AtomicBool::new(false), gain: Mutex::new(gain), } } fn lock(&self) -> std::sync::MutexGuard<'_, std::collections::VecDeque> { self.samples.lock().unwrap_or_else(|e| e.into_inner()) } } /// Mando a distancia del altavoz: se puede clonar y repartir por los hilos. #[derive(Clone)] pub struct PlaybackHandle { ring: Arc, sample_rate: u32, } impl PlaybackHandle { /// Encola muestras mono ya a la frecuencia del dispositivo. pub fn push(&self, samples: &[f32]) { if samples.is_empty() { return; } let mut queue = self.ring.lock(); queue.extend(samples.iter().copied()); self.ring.active.store(true, Ordering::SeqCst); } /// Encola audio del TTS (24 kHz mono), remuestreando si hace falta. pub fn push_tts(&self, samples: &[f32]) { if self.sample_rate == TTS_SAMPLE_RATE { self.push(samples); return; } let format = InputFormat { sample_rate: TTS_SAMPLE_RATE as usize, channels: 1, }; self.push(&to_mono_at(samples, format, self.sample_rate)); } /// Calla ahora mismo y tira lo que quedaba por sonar. pub fn stop(&self) { let mut queue = self.ring.lock(); queue.clear(); self.ring.active.store(false, Ordering::SeqCst); self.ring.drained.notify_all(); } /// Segundos de audio pendientes de sonar. pub fn queued_secs(&self) -> f32 { self.ring.lock().len() as f32 / self.sample_rate as f32 } /// `true` si ya ha salido audio por el altavoz. pub fn has_played(&self) -> bool { self.ring.played.load(Ordering::SeqCst) > 0 } pub fn reset_played(&self) { self.ring.played.store(0, Ordering::SeqCst); } pub fn is_active(&self) -> bool { self.ring.active.load(Ordering::SeqCst) } pub fn set_gain(&self, gain: f32) { *self.ring.gain.lock().unwrap_or_else(|e| e.into_inner()) = gain; } /// Espera a que se vacíe la cola, o a que venza el plazo. /// /// Devuelve `true` si terminó de sonar todo y `false` si se agotó el /// tiempo, para que quien llama distinga «ya está» de «sigue sonando». pub fn wait_drained(&self, timeout: Duration) -> bool { let deadline = std::time::Instant::now() + timeout; let mut queue = self.ring.lock(); while !queue.is_empty() { let left = deadline.saturating_duration_since(std::time::Instant::now()); if left.is_zero() { return false; } let (guard, result) = self .ring .drained .wait_timeout(queue, left.min(Duration::from_millis(50))) .unwrap_or_else(|e| e.into_inner()); queue = guard; if result.timed_out() && queue.is_empty() { break; } } self.ring.active.store(false, Ordering::SeqCst); true } } pub struct Playback { stream: cpal::Stream, handle: PlaybackHandle, pub device_name: String, pub sample_rate: u32, } impl Playback { pub fn open(config: &AudioConfig) -> Result { let host = cpal::default_host(); let device = select_device(&host, &config.output_device)?; let device_name = describe(&device); let supported = preferred_config(&device)?; let sample_rate = supported.sample_rate(); let channels = supported.channels() as usize; let stream_config: cpal::StreamConfig = supported.clone().into(); let ring = Arc::new(Ring::new(config.output_gain)); let on_error = |err| tracing::error!(target: "audio", %err, "flujo de salida"); // La retrollamada: coger lo que haya, rellenar con silencio lo que // falte y salir. Nunca espera a que llegue audio, porque quedarse // esperando aquí se oye como un corte. macro_rules! build { ($sample:ty, $silence:expr, $convert:expr) => {{ let ring = Arc::clone(&ring); device .build_output_stream( &stream_config, move |data: &mut [$sample], _: &_| { let gain = *ring.gain.lock().unwrap_or_else(|e| e.into_inner()); let mut queue = ring.lock(); let mut served = 0u64; for frame in data.chunks_mut(channels) { match queue.pop_front() { Some(sample) => { let value = (sample * gain).clamp(-1.0, 1.0); // Mono a N canales: la misma muestra en todos. for out in frame.iter_mut() { *out = $convert(value); } served += 1; } None => { for out in frame.iter_mut() { *out = $silence; } } } } if served > 0 { ring.played.fetch_add(served, Ordering::SeqCst); } if queue.is_empty() { ring.drained.notify_all(); } }, on_error, None, ) .map_err(|e| Error::Audio(format!("no se pudo abrir la salida: {e}")))? }}; } let stream = match supported.sample_format() { SampleFormat::F32 => build!(f32, 0.0f32, |v: f32| v), SampleFormat::I16 => build!(i16, 0i16, |v: f32| (v * 32767.0) as i16), SampleFormat::U16 => { build!(u16, u16::MAX / 2, |v: f32| ((v * 32767.0) as i32 + 32768) as u16) } other => { return Err(Error::Audio(format!( "formato de salida no soportado: {other:?}" ))) } }; stream .play() .map_err(|e| Error::Audio(format!("no se pudo arrancar la salida: {e}")))?; tracing::info!( target: "audio", dispositivo = %device_name, hz = sample_rate, canales = channels, "salida abierta" ); Ok(Self { stream, handle: PlaybackHandle { ring, sample_rate }, device_name, sample_rate, }) } pub fn handle(&self) -> PlaybackHandle { self.handle.clone() } pub fn stop(self) { self.handle.stop(); drop(self.stream); } } fn select_device(host: &cpal::Host, wanted: &str) -> Result { if wanted.is_empty() { return host .default_output_device() .ok_or_else(|| Error::Audio("no hay dispositivo de salida".into())); } let wanted_lower = wanted.to_lowercase(); let devices = host .output_devices() .map_err(|e| Error::Audio(format!("no se pudieron listar las salidas: {e}")))?; let mut seen = Vec::new(); for device in devices { let name = describe(&device); if name.to_lowercase().contains(&wanted_lower) { return Ok(device); } seen.push(name); } Err(Error::Audio(format!( "ninguna salida coincide con «{wanted}». Disponibles: {}", seen.join(", ") ))) } fn preferred_config(device: &cpal::Device) -> Result { // 24 kHz nativo evita remuestrear la síntesis; si no, se coge lo de fábrica. let native = device .supported_output_configs() .map_err(|e| Error::Audio(format!("no se pudo consultar la salida: {e}")))? .filter(|range| { range.min_sample_rate() <= TTS_SAMPLE_RATE && TTS_SAMPLE_RATE <= range.max_sample_rate() }) .find(|range| range.sample_format() == SampleFormat::F32) .map(|range| range.with_sample_rate(TTS_SAMPLE_RATE)); match native { Some(config) => Ok(config), None => device .default_output_config() .map_err(|e| Error::Audio(format!("sin configuración de salida: {e}"))), } }