aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-app/src/session.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-app/src/session.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-app/src/session.rs')
-rw-r--r--crates/asist-app/src/session.rs160
1 files changed, 160 insertions, 0 deletions
diff --git a/crates/asist-app/src/session.rs b/crates/asist-app/src/session.rs
new file mode 100644
index 0000000..b74a4f3
--- /dev/null
+++ b/crates/asist-app/src/session.rs
@@ -0,0 +1,160 @@
+//! Estado compartido entre los hilos del pipeline.
+//!
+//! Sólo hay dos cosas que de verdad tienen que compartirse: qué turno es el
+//! vigente, y las banderas que permiten cortar lo que está en marcha. Todo lo
+//! demás viaja por canales. Mantener esta superficie pequeña es lo que hace
+//! que la interrupción sea razonable de seguir: cortar es subir el turno y
+//! levantar dos banderas.
+
+use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
+use std::sync::Arc;
+
+use asist_core::event::TurnId;
+use asist_core::http::Cancel;
+
+#[derive(Clone)]
+pub struct Session {
+ inner: Arc<Inner>,
+}
+
+struct Inner {
+ /// Turno que se está atendiendo. Cualquier trabajo de un turno anterior
+ /// que aparezca por un canal es basura y se tira.
+ current: AtomicU64,
+ /// El asistente está hablando (o a punto de hacerlo).
+ speaking: AtomicBool,
+ /// Cierre solicitado.
+ stopping: AtomicBool,
+ /// Corta la respuesta del modelo en curso.
+ llm: Cancel,
+ /// Corta la síntesis en curso.
+ tts: Cancel,
+}
+
+impl Default for Session {
+ fn default() -> Self {
+ Self::new()
+ }
+}
+
+impl Session {
+ pub fn new() -> Self {
+ Self {
+ inner: Arc::new(Inner {
+ current: AtomicU64::new(0),
+ speaking: AtomicBool::new(false),
+ stopping: AtomicBool::new(false),
+ llm: Cancel::new(),
+ tts: Cancel::new(),
+ }),
+ }
+ }
+
+ pub fn current(&self) -> TurnId {
+ TurnId(self.inner.current.load(Ordering::SeqCst))
+ }
+
+ /// Abre un turno nuevo y devuelve su identificador.
+ pub fn begin_turn(&self) -> TurnId {
+ let id = TurnId(self.inner.current.fetch_add(1, Ordering::SeqCst) + 1);
+ self.inner.llm.reset();
+ self.inner.tts.reset();
+ id
+ }
+
+ /// ¿Sigue siendo `turn` el turno vigente?
+ ///
+ /// Lo consultan los hilos antes de gastar trabajo: una frase que
+ /// pertenece a un turno ya superado no debe sintetizarse ni oírse.
+ pub fn is_current(&self, turn: TurnId) -> bool {
+ turn == self.current()
+ }
+
+ /// Corta todo lo que esté en marcha para el turno actual.
+ pub fn interrupt(&self) {
+ self.inner.llm.cancel();
+ self.inner.tts.cancel();
+ self.inner.speaking.store(false, Ordering::SeqCst);
+ }
+
+ pub fn llm_cancel(&self) -> &Cancel {
+ &self.inner.llm
+ }
+
+ pub fn tts_cancel(&self) -> &Cancel {
+ &self.inner.tts
+ }
+
+ pub fn set_speaking(&self, speaking: bool) {
+ self.inner.speaking.store(speaking, Ordering::SeqCst);
+ }
+
+ pub fn is_speaking(&self) -> bool {
+ self.inner.speaking.load(Ordering::SeqCst)
+ }
+
+ pub fn request_stop(&self) {
+ self.inner.stopping.store(true, Ordering::SeqCst);
+ self.interrupt();
+ }
+
+ pub fn is_stopping(&self) -> bool {
+ self.inner.stopping.load(Ordering::SeqCst)
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn cada_turno_recibe_un_identificador_creciente() {
+ let session = Session::new();
+ assert_eq!(session.begin_turn(), TurnId(1));
+ assert_eq!(session.begin_turn(), TurnId(2));
+ assert_eq!(session.current(), TurnId(2));
+ }
+
+ #[test]
+ fn el_trabajo_de_un_turno_viejo_deja_de_ser_vigente() {
+ let session = Session::new();
+ let primero = session.begin_turn();
+ session.begin_turn();
+ assert!(
+ !session.is_current(primero),
+ "el turno viejo debe descartarse"
+ );
+ }
+
+ #[test]
+ fn abrir_turno_rearma_las_cancelaciones() {
+ let session = Session::new();
+ session.begin_turn();
+ session.interrupt();
+ assert!(session.llm_cancel().is_cancelled());
+
+ session.begin_turn();
+ assert!(
+ !session.llm_cancel().is_cancelled(),
+ "un turno nuevo no puede heredar la cancelación del anterior"
+ );
+ }
+
+ #[test]
+ fn interrumpir_calla_y_cancela_las_dos_etapas() {
+ let session = Session::new();
+ session.begin_turn();
+ session.set_speaking(true);
+ session.interrupt();
+ assert!(session.tts_cancel().is_cancelled());
+ assert!(!session.is_speaking());
+ }
+
+ #[test]
+ fn pedir_el_cierre_tambien_interrumpe() {
+ let session = Session::new();
+ session.request_stop();
+ assert!(session.is_stopping());
+ assert!(session.llm_cancel().is_cancelled());
+ }
+}