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-llm/src/client.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-llm/src/client.rs')
| -rw-r--r-- | crates/asist-llm/src/client.rs | 342 |
1 files changed, 342 insertions, 0 deletions
diff --git a/crates/asist-llm/src/client.rs b/crates/asist-llm/src/client.rs new file mode 100644 index 0000000..1dca2af --- /dev/null +++ b/crates/asist-llm/src/client.rs @@ -0,0 +1,342 @@ +//! Transporte hacia llama-server. + +use std::time::{Duration, Instant}; + +use serde_json::{json, Value}; + +use asist_core::config::LlmConfig; +use asist_core::error::{Error, Result}; +use asist_core::http::{Cancel, HttpClient}; +use asist_core::tools::{ToolCall, ToolRegistry}; + +use crate::chat::Conversation; + +/// Un trozo de respuesta según sale del modelo. +#[derive(Debug, Clone)] +pub enum Delta { + /// Texto para el usuario. + Text(String), + /// El modelo ha decidido llamar a herramientas; no habrá más texto. + ToolCalls(Vec<ToolCall>), +} + +/// Cómo terminó una respuesta en streaming. +#[derive(Debug, Clone, Default)] +pub struct StreamOutcome { + pub text: String, + pub tool_calls: Vec<ToolCall>, + /// Tiempo hasta el primer fragmento con contenido. + pub ttft: Option<Duration>, + pub stopped_early: bool, +} + +pub struct LlmClient { + http: HttpClient, + config: LlmConfig, +} + +impl LlmClient { + pub fn new(authority: String, config: &LlmConfig) -> 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") + } + + /// Espera a que el servidor termine de cargar el modelo. + 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::Llm(format!( + "{} no respondió a /health en {} s", + self.http.authority, + timeout.as_secs() + ))) + } + + /// Comprueba que la plantilla de chat no arranca un bloque de + /// razonamiento sin cerrar. + /// + /// Merece la pena avisar: con la plantilla original de Qwen3.5 el modelo + /// se pasa entre 7 y 9 segundos «pensando» antes de la primera palabra + /// audible, y desde fuera parece que el asistente se ha colgado. + pub fn warn_if_thinking_template(&self) { + let Ok(body) = self.http.get("/props") else { + return; + }; + let Ok(props) = serde_json::from_slice::<Value>(&body) else { + return; + }; + let Some(template) = props.get("chat_template").and_then(Value::as_str) else { + return; + }; + + let opens = template.matches("<think>").count(); + let closes = template.matches("</think>").count(); + if opens > closes { + tracing::warn!( + target: "llm", + "la plantilla de chat abre <think> sin cerrarlo: el modelo razonará \ + varios segundos antes de contestar. Arranca llama-server con \ + --chat-template-file config/qwen35-no-think.jinja" + ); + } + } + + fn request_body(&self, chat: &Conversation, tools: Option<&ToolRegistry>) -> Value { + let mut body = json!({ + "messages": chat.to_json(), + "stream": true, + "temperature": self.config.temperature, + "top_p": self.config.top_p, + "top_k": self.config.top_k, + "min_p": self.config.min_p, + "repeat_penalty": self.config.repeat_penalty, + "max_tokens": self.config.max_tokens, + }); + if !self.config.model.is_empty() { + body["model"] = json!(self.config.model); + } + if let Some(tools) = tools.filter(|t| !t.is_empty()) { + body["tools"] = tools.schema(); + body["tool_choice"] = json!("auto"); + } + body + } + + /// Envía la conversación y entrega la respuesta a trozos. + /// + /// `on_delta` devuelve `false` para cortar —lo que hace el orquestador + /// cuando el usuario interrumpe—; `cancel` hace lo mismo desde otro hilo. + pub fn stream( + &self, + chat: &Conversation, + tools: Option<&ToolRegistry>, + cancel: &Cancel, + mut on_delta: impl FnMut(&Delta) -> bool, + ) -> Result<StreamOutcome> { + let body = self.request_body(chat, tools); + let started = Instant::now(); + let mut outcome = StreamOutcome::default(); + let mut assembler = ToolCallAssembler::default(); + let mut finished = false; + + let result = self + .http + .post_json_lines("/v1/chat/completions", &body, cancel, |line| { + let Some(payload) = line.strip_prefix("data: ") else { + return true; + }; + if payload.trim() == "[DONE]" { + finished = true; + return false; + } + let Ok(event) = serde_json::from_str::<Value>(payload) else { + tracing::debug!(target: "llm", %payload, "fragmento SSE ilegible"); + return true; + }; + let Some(choice) = event.get("choices").and_then(|c| c.get(0)) else { + return true; + }; + let delta = choice.get("delta").unwrap_or(&Value::Null); + + if let Some(calls) = delta.get("tool_calls").and_then(Value::as_array) { + assembler.absorb(calls); + } + if let Some(text) = delta.get("content").and_then(Value::as_str) { + if !text.is_empty() { + if outcome.ttft.is_none() { + outcome.ttft = Some(started.elapsed()); + } + outcome.text.push_str(text); + if !on_delta(&Delta::Text(text.to_string())) { + outcome.stopped_early = true; + return false; + } + } + } + // Un `finish_reason` cierra la respuesta aunque el servidor no + // llegue a mandar el [DONE] (pasa al cortar la conexión). + if choice.get("finish_reason").is_some_and(|r| !r.is_null()) { + finished = true; + return false; + } + true + }); + + match result { + Ok(()) => {} + // Cortar a propósito no es un fallo: `on_delta` devolvió false o + // se activó la cancelación. + Err(Error::Cancelled) if finished || outcome.stopped_early || cancel.is_cancelled() => { + } + Err(err) => return Err(err), + } + + outcome.tool_calls = assembler.finish(); + if !outcome.tool_calls.is_empty() + && !on_delta(&Delta::ToolCalls(outcome.tool_calls.clone())) + { + outcome.stopped_early = true; + } + Ok(outcome) + } +} + +/// Reensambla las llamadas a herramientas, que llegan repartidas en fragmentos. +/// +/// El streaming manda el nombre en un fragmento y los argumentos en varios +/// más, identificados sólo por su `index`; hasta que no termina el flujo no +/// hay una llamada completa que ejecutar. +#[derive(Default)] +struct ToolCallAssembler { + partial: Vec<PartialCall>, +} + +#[derive(Default, Clone)] +struct PartialCall { + id: String, + name: String, + arguments: String, +} + +impl ToolCallAssembler { + fn absorb(&mut self, calls: &[Value]) { + for call in calls { + let index = call.get("index").and_then(Value::as_u64).unwrap_or(0) as usize; + if self.partial.len() <= index { + self.partial.resize(index + 1, PartialCall::default()); + } + let slot = &mut self.partial[index]; + if let Some(id) = call.get("id").and_then(Value::as_str) { + if !id.is_empty() { + slot.id = id.to_string(); + } + } + let Some(function) = call.get("function") else { + continue; + }; + if let Some(name) = function.get("name").and_then(Value::as_str) { + if !name.is_empty() { + slot.name = name.to_string(); + } + } + if let Some(args) = function.get("arguments").and_then(Value::as_str) { + slot.arguments.push_str(args); + } + } + } + + fn finish(self) -> Vec<ToolCall> { + self.partial + .into_iter() + .enumerate() + .filter(|(_, call)| !call.name.is_empty()) + .map(|(i, call)| ToolCall { + id: if call.id.is_empty() { + format!("call_{i}") + } else { + call.id + }, + name: call.name, + arguments: if call.arguments.is_empty() { + "{}".into() + } else { + call.arguments + }, + }) + .collect() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn absorb(assembler: &mut ToolCallAssembler, raw: &str) { + let value: Value = serde_json::from_str(raw).unwrap(); + assembler.absorb(value.as_array().unwrap()); + } + + #[test] + fn los_argumentos_troceados_se_reensamblan() { + let mut assembler = ToolCallAssembler::default(); + absorb( + &mut assembler, + r#"[{"index":0,"id":"c1","function":{"name":"eco","arguments":"{\"a\":"}}]"#, + ); + absorb( + &mut assembler, + r#"[{"index":0,"function":{"arguments":"1}"}}]"#, + ); + let calls = assembler.finish(); + assert_eq!(calls.len(), 1); + assert_eq!(calls[0].id, "c1"); + assert_eq!(calls[0].name, "eco"); + assert_eq!(calls[0].arguments, r#"{"a":1}"#); + } + + #[test] + fn se_reensamblan_varias_llamadas_en_paralelo() { + let mut assembler = ToolCallAssembler::default(); + absorb( + &mut assembler, + r#"[{"index":0,"id":"a","function":{"name":"uno","arguments":"{}"}}]"#, + ); + absorb( + &mut assembler, + r#"[{"index":1,"id":"b","function":{"name":"dos","arguments":"{}"}}]"#, + ); + let calls = assembler.finish(); + assert_eq!(calls.len(), 2); + assert_eq!(calls[0].name, "uno"); + assert_eq!(calls[1].name, "dos"); + } + + #[test] + fn sin_argumentos_se_asume_el_objeto_vacio() { + let mut assembler = ToolCallAssembler::default(); + absorb( + &mut assembler, + r#"[{"index":0,"id":"c","function":{"name":"hora_actual"}}]"#, + ); + assert_eq!(assembler.finish()[0].arguments, "{}"); + } + + #[test] + fn un_hueco_sin_nombre_no_produce_una_llamada() { + // El servidor puede numerar el índice 1 sin haber mandado nunca el 0. + let mut assembler = ToolCallAssembler::default(); + absorb( + &mut assembler, + r#"[{"index":1,"id":"b","function":{"name":"dos","arguments":"{}"}}]"#, + ); + let calls = assembler.finish(); + assert_eq!( + calls.len(), + 1, + "el hueco no debe convertirse en una llamada vacía" + ); + assert_eq!(calls[0].name, "dos"); + } + + #[test] + fn sin_id_se_genera_uno_estable() { + let mut assembler = ToolCallAssembler::default(); + absorb( + &mut assembler, + r#"[{"index":0,"function":{"name":"x","arguments":"{}"}}]"#, + ); + assert_eq!(assembler.finish()[0].id, "call_0"); + } +} |