From f8f98e83481e7376a235fb305a095552433f79ab Mon Sep 17 00:00:00 2001 From: elvis Date: Sun, 6 Sep 2026 19:21:16 -0300 Subject: Local voice assistant on top of Canary, llama.cpp and qwentts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 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 --- crates/asist-audio/src/playback.rs | 287 +++++++++++++++++++++++++++++++++++++ 1 file changed, 287 insertions(+) create mode 100644 crates/asist-audio/src/playback.rs (limited to 'crates/asist-audio/src/playback.rs') 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>, + /// 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}"))), + } +} -- cgit v1.2.3