diff options
| author | elvis <elvis@claros.ar> | 2026-09-06 19:21:16 -0300 |
|---|---|---|
| committer | elvis <elvis@claros.ar> | 2026-09-06 19:23:37 -0300 |
| commit | f8f98e83481e7376a235fb305a095552433f79ab (patch) | |
| tree | 6d0ecddb26bef8430a086431de40f8ce9b1039e6 /crates/asist-tts/src/lib.rs | |
| download | asist-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-tts/src/lib.rs')
| -rw-r--r-- | crates/asist-tts/src/lib.rs | 256 |
1 files changed, 256 insertions, 0 deletions
diff --git a/crates/asist-tts/src/lib.rs b/crates/asist-tts/src/lib.rs new file mode 100644 index 0000000..b250116 --- /dev/null +++ b/crates/asist-tts/src/lib.rs @@ -0,0 +1,256 @@ +//! Síntesis de voz sobre qwentts.cpp. +//! +//! Se habla con `tts-server` por HTTP en lugar de invocar el binario +//! `qwen-tts` una vez por frase. La diferencia no es de estilo: el proceso +//! carga 2,1 GB de modelo hablante más 291 MB de codec, y pagar esa carga en +//! cada frase pondría varios segundos delante de cada respuesta. El servidor +//! los mantiene residentes en la GPU y responde en formato `pcm`, que llega +//! troceado según se genera. + +use std::time::{Duration, Instant}; + +use base64::Engine; +use serde_json::json; + +use asist_core::config::{ReferenceVoice, TtsConfig}; +use asist_core::error::{Error, Result}; +use asist_core::http::{Cancel, HttpClient}; + +/// Frecuencia a la que sintetiza el modelo. +pub const SAMPLE_RATE: u32 = 24_000; + +/// Cómo terminó una síntesis. +#[derive(Debug, Clone, Default)] +pub struct SpeechOutcome { + /// Tiempo hasta el primer bloque de audio. + pub ttfb: Option<Duration>, + pub total: Duration, + pub samples: usize, + pub cancelled: bool, +} + +impl SpeechOutcome { + pub fn audio_secs(&self) -> f32 { + self.samples as f32 / SAMPLE_RATE as f32 + } + + /// Múltiplo del tiempo real. Por encima de 1,0 el sintetizador va más + /// lento que el habla y la cola de reproducción acabará vaciándose. + pub fn rtf(&self) -> f32 { + let audio = self.audio_secs(); + if audio <= 0.0 { + return f32::INFINITY; + } + self.total.as_secs_f32() / audio + } +} + +pub struct TtsClient { + http: HttpClient, + config: TtsConfig, +} + +impl TtsClient { + pub fn new(authority: String, config: &TtsConfig) -> Self { + Self { + http: HttpClient::new(authority) + .with_read_timeout(Duration::from_secs(config.request_timeout_secs)), + config: config.clone(), + } + } + + pub fn healthy(&self) -> bool { + self.http.healthy("/health") + } + + pub fn wait_ready(&self, timeout: Duration) -> Result<()> { + let deadline = Instant::now() + timeout; + while Instant::now() < deadline { + if self.healthy() { + return Ok(()); + } + std::thread::sleep(Duration::from_millis(500)); + } + Err(Error::Tts(format!( + "{} no respondió a /health en {} s", + self.http.authority, + timeout.as_secs() + ))) + } + + /// Registra una voz clonada a partir de los latentes de `qwen-codec`. + /// + /// Se mandan los `.spk` y `.rvq` ya extraídos en lugar del WAV original: + /// el servidor los toma tal cual y se ahorra volver a analizar la + /// referencia en cada arranque. + pub fn register_voice(&self, voice: &ReferenceVoice) -> Result<()> { + let read = |path: &std::path::Path, what: &str| -> Result<Vec<u8>> { + std::fs::read(path).map_err(|e| { + Error::Tts(format!("no se pudo leer {what} ({}): {e}", path.display())) + }) + }; + let b64 = base64::engine::general_purpose::STANDARD; + let speaker = b64.encode(read(&voice.speaker, "el embedding del hablante")?); + let codes = b64.encode(read(&voice.codes, "los códigos de referencia")?); + let transcript = String::from_utf8_lossy(&read(&voice.transcript, "la transcripción")?) + .trim() + .to_string(); + + if transcript.is_empty() { + return Err(Error::Tts(format!( + "la transcripción de referencia {} está vacía; sin ella no se activa el clonado", + voice.transcript.display() + ))); + } + + let body = json!({ + "name": voice.name, + "ref_text": transcript, + "spk_b64": speaker, + "rvq_b64": codes, + }); + self.http.post_json("/v1/audio/voices", &body)?; + tracing::info!(target: "tts", voz = %voice.name, "voz clonada registrada"); + Ok(()) + } + + /// Voces que el servidor conoce ahora mismo. + pub fn voices(&self) -> Result<Vec<String>> { + let body = self.http.get("/v1/audio/voices")?; + let parsed: serde_json::Value = serde_json::from_slice(&body)?; + Ok(parsed + .get("voices") + .and_then(|v| v.as_array()) + .map(|voices| { + voices + .iter() + .filter_map(|v| { + v.get("name") + .and_then(|n| n.as_str()) + .or_else(|| v.as_str()) + .map(String::from) + }) + .collect() + }) + .unwrap_or_default()) + } + + /// Sintetiza `text` y entrega el audio a trozos según se genera. + /// + /// `on_audio` recibe muestras mono f32 a 24 kHz y devuelve `false` para + /// cortar; `cancel` hace lo propio desde otro hilo. Ese corte es lo que + /// permite callar de golpe cuando el usuario interrumpe, sin esperar a que + /// termine de generarse una frase que ya no se va a oír. + pub fn speak( + &self, + text: &str, + cancel: &Cancel, + mut on_audio: impl FnMut(&[f32]) -> bool, + ) -> Result<SpeechOutcome> { + let body = json!({ + "input": text, + "voice": self.config.voice, + "response_format": "pcm", + "temperature": self.config.temperature, + "top_k": self.config.top_k, + "top_p": self.config.top_p, + "repetition_penalty": self.config.repetition_penalty, + "max_new_tokens": self.config.max_new_tokens, + }); + + let started = Instant::now(); + let mut outcome = SpeechOutcome::default(); + // El servidor no garantiza que cada bloque traiga un número par de + // bytes, así que el byte suelto de un bloque espera al siguiente en + // vez de desplazar todas las muestras posteriores. + let mut odd: Option<u8> = None; + let mut samples = Vec::with_capacity(8192); + let mut stopped = false; + + let result = self + .http + .post_json_streaming("/v1/audio/speech", &body, cancel, |chunk| { + if outcome.ttfb.is_none() { + outcome.ttfb = Some(started.elapsed()); + } + samples.clear(); + let mut bytes = chunk; + if let Some(pending) = odd.take() { + if let Some((first, rest)) = bytes.split_first() { + samples.push(i16::from_le_bytes([pending, *first]) as f32 / 32768.0); + bytes = rest; + } else { + odd = Some(pending); + } + } + if bytes.len() % 2 == 1 { + odd = Some(bytes[bytes.len() - 1]); + bytes = &bytes[..bytes.len() - 1]; + } + for pair in bytes.chunks_exact(2) { + samples.push(i16::from_le_bytes([pair[0], pair[1]]) as f32 / 32768.0); + } + outcome.samples += samples.len(); + if !on_audio(&samples) { + stopped = true; + return false; + } + true + }); + + outcome.total = started.elapsed(); + match result { + Ok(()) => {} + Err(Error::Cancelled) => outcome.cancelled = true, + Err(err) => return Err(err), + } + outcome.cancelled |= stopped; + + tracing::debug!( + target: "tts", + caracteres = text.chars().count(), + ttfb_ms = outcome.ttfb.map(|d| d.as_millis()).unwrap_or(0), + audio_s = outcome.audio_secs(), + rtf = outcome.rtf(), + cancelada = outcome.cancelled, + "síntesis" + ); + Ok(outcome) + } + + /// Sintetiza una frase corta para pagar por adelantado la construcción de + /// los grafos: la primera petición cuesta unos 3,5 s más que las + /// siguientes, y no conviene gastarlos en el primer turno de verdad. + pub fn warmup(&self) -> Result<()> { + if !self.config.warmup { + return Ok(()); + } + let started = Instant::now(); + let cancel = Cancel::new(); + // El audio se tira: sólo interesa el efecto secundario. + self.speak("Listo.", &cancel, |_| true)?; + tracing::info!(target: "tts", ms = started.elapsed().as_millis(), "precalentado"); + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn el_rtf_relaciona_tiempo_de_calculo_y_audio() { + let outcome = SpeechOutcome { + total: Duration::from_secs(2), + samples: SAMPLE_RATE as usize * 4, + ..Default::default() + }; + assert_eq!(outcome.audio_secs(), 4.0); + assert_eq!(outcome.rtf(), 0.5); + } + + #[test] + fn sin_audio_el_rtf_es_infinito_en_vez_de_dividir_por_cero() { + assert!(SpeechOutcome::default().rtf().is_infinite()); + } +} |