aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-audio/src/playback.rs
diff options
context:
space:
mode:
authorelvis <elvis@claros.ar>2026-09-06 19:21:16 -0300
committerelvis <elvis@claros.ar>2026-09-06 19:23:37 -0300
commitf8f98e83481e7376a235fb305a095552433f79ab (patch)
tree6d0ecddb26bef8430a086431de40f8ce9b1039e6 /crates/asist-audio/src/playback.rs
downloadasist-p-f8f98e83481e7376a235fb305a095552433f79ab.tar.gz
asist-p-f8f98e83481e7376a235fb305a095552433f79ab.zip
Local voice assistant on top of Canary, llama.cpp and qwentts
Pipeline en Rust de hilos y canales que une los tres motores: el micrófono alimenta un segmentador con VAD, las intervenciones cerradas van al reconocedor, la transcripción al modelo y cada frase que este cierra sale hacia el sintetizador sin esperar al resto de la respuesta. Dos reglas sostienen el diseño: ninguna etapa bloquea a la anterior —quien va sobrado descarta trabajo en lugar de acumular retraso— y todo lo que viaja por los canales lleva el turno al que pertenece, así que interrumpir es subir el contador y levantar dos banderas de cancelación. Seis crates: core (configuración, eventos, HTTP, telemetría, herramientas), audio (cpal, VAD, anillo de reproducción), asr, llm, tts y app (supervisor de procesos y orquestador). Los motores van como submódulos fijados a un commit, con los cambios locales en vendor/patches. Midiendo el pipeline aparecieron tres cuellos de botella de configuración que valieron más que cualquier cambio de código, todos documentados en docs/RENDIMIENTO.md: - tts-server decodificaba el audio en bloques de 24 s, de modo que el modo «streaming» llegaba de una pieza: 4948 ms -> 585 ms hasta el primer audio. - La plantilla de chat del modelo abre <think> y no lo cierra nunca, sin variable que lo apague: 8630 ms -> 413 ms hasta el primer token, con una copia de la plantilla que deja el bloque cerrado de entrada. - Cualquier indicación de estilo junto a la guía de herramientas hace que este modelo de 2B deje de llamarlas y se invente el dato (8/8 aciertos con la guía sola, 0/8 con la persona de asistente de voz). El turno alterna ahora entre dos instrucciones de sistema. La ejecución de órdenes del sistema queda implementada y apagada, tras cuatro barreras: lista blanca sobre el ejecutable, rutas rechazadas, sin shell que interprete metacaracteres y plazo máximo. 68 pruebas unitarias sin modelos, más seis de integración que se saltan solas si no hay servidores y se turnan la GPU: en paralelo, los dos servidores no caben en 4 GB y miden contención en vez de latencia. Claude-Session: https://claude.ai/code/session_01FNxz5cSdQSscJH9H7b8uGU
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}"))),
+ }
+}