aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-audio/src/playback.rs
diff options
context:
space:
mode:
Diffstat (limited to 'crates/asist-audio/src/playback.rs')
-rw-r--r--crates/asist-audio/src/playback.rs287
1 files changed, 287 insertions, 0 deletions
diff --git a/crates/asist-audio/src/playback.rs b/crates/asist-audio/src/playback.rs
new file mode 100644
index 0000000..fcdfe7d
--- /dev/null
+++ b/crates/asist-audio/src/playback.rs
@@ -0,0 +1,287 @@
+//! 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<std::collections::VecDeque<f32>>,
+ /// 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<f32>,
+}
+
+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<f32>> {
+ 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<Ring>,
+ 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<Self> {
+ 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<cpal::Device> {
+ 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<cpal::SupportedStreamConfig> {
+ // 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}"))),
+ }
+}