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-core/src | |
| 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-core/src')
| -rw-r--r-- | crates/asist-core/src/config.rs | 508 | ||||
| -rw-r--r-- | crates/asist-core/src/error.rs | 82 | ||||
| -rw-r--r-- | crates/asist-core/src/event.rs | 127 | ||||
| -rw-r--r-- | crates/asist-core/src/http.rs | 351 | ||||
| -rw-r--r-- | crates/asist-core/src/lib.rs | 17 | ||||
| -rw-r--r-- | crates/asist-core/src/telemetry.rs | 255 | ||||
| -rw-r--r-- | crates/asist-core/src/text.rs | 330 | ||||
| -rw-r--r-- | crates/asist-core/src/tools.rs | 563 |
8 files changed, 2233 insertions, 0 deletions
diff --git a/crates/asist-core/src/config.rs b/crates/asist-core/src/config.rs new file mode 100644 index 0000000..4ca8d8c --- /dev/null +++ b/crates/asist-core/src/config.rs @@ -0,0 +1,508 @@ +//! Configuración del asistente: un TOML con valores por defecto que ya +//! incorporan lo aprendido midiendo el pipeline (ver `docs/RENDIMIENTO.md`). + +use std::path::{Path, PathBuf}; + +use serde::{Deserialize, Serialize}; + +use crate::error::{Error, Result}; + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct Config { + pub general: General, + pub audio: AudioConfig, + pub vad: VadConfig, + pub asr: AsrConfig, + pub llm: LlmConfig, + pub tts: TtsConfig, + pub tools: ToolsConfig, + pub supervisor: SupervisorConfig, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct General { + /// Idioma de la conversación; se propaga a ASR y TTS. + pub language: String, + /// Instrucción de sistema. Pide frases cortas y sin markdown a propósito: + /// el TTS lee literalmente los asteriscos y las viñetas. + pub system_prompt: String, + /// Instrucción de sistema de la pasada en que el modelo decide si usar + /// una herramienta. + /// + /// Va aparte, y **sustituye** a `system_prompt` en esa pasada, por una + /// razón medida y no por gusto: con Qwen3.5-2B, añadir cualquier + /// indicación de estilo a esta guía —en cualquier posición, incluso dos + /// palabras— hace que el modelo deje de llamar a las herramientas y se + /// invente el dato. Medido sobre 8 intentos: la guía sola acierta 8/8; + /// con «Responde breve.» detrás, 1/8; con la persona de asistente de voz, + /// 0/8. Ver docs/RENDIMIENTO.md. + pub tools_prompt: String, + /// Turnos de historial que se envían al modelo (0 = sin memoria). + pub history_turns: usize, + /// Imprime un resumen de latencias por turno al terminar cada respuesta. + pub report_latency: bool, +} + +impl Default for General { + fn default() -> Self { + Self { + language: "es".into(), + system_prompt: concat!( + "Eres un asistente de voz en español. Tus respuestas se leen en voz alta, ", + "así que responde en una o dos frases cortas, en texto plano corrido. ", + "No uses markdown, ni listas, ni asteriscos, ni emojis, ni encabezados. ", + "No escribas URLs ni código salvo que te lo pidan explícitamente. ", + "Si no sabes algo, dilo en una frase." + ) + .into(), + tools_prompt: concat!( + "Antes de responder, comprueba si alguna de tus herramientas te da el dato. ", + "Si es así, llámala primero y espera su resultado; no contestes de memoria. ", + "Sólo cuando tengas el resultado, resúmelo en una frase." + ) + .into(), + history_turns: 8, + report_latency: true, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct AudioConfig { + /// Nombre (o subcadena) del dispositivo de entrada; vacío = el de por defecto. + pub input_device: String, + /// Ídem para la salida. + pub output_device: String, + /// Segundos de audio que la reproducción mantiene en cola antes de + /// arrancar. Amortigua los baches del TTS sin añadir latencia perceptible. + pub playback_prebuffer: f32, + /// Ganancia aplicada a la reproducción. + pub output_gain: f32, +} + +impl Default for AudioConfig { + fn default() -> Self { + Self { + input_device: String::new(), + output_device: String::new(), + playback_prebuffer: 0.20, + output_gain: 1.0, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct VadConfig { + /// Ventana de audio por decisión del detector. + pub frame_seconds: f32, + /// Audio conservado antes del inicio de voz para no cortar la primera sílaba. + pub preroll_seconds: f32, + /// Silencio que cierra una intervención. + pub silence_hold: f32, + /// Intervención más corta que merezca una transcripción final. + pub min_utterance: f32, + /// Corte forzado, que acota el coste de la decodificación final. + pub max_utterance: f32, + /// Múltiplo del suelo de ruido a partir del cual se considera voz. + pub threshold_factor: f32, + pub min_threshold: f32, + pub max_threshold: f32, + /// Permite hablar encima del asistente para cortarlo. + /// + /// Desactivado por defecto: con altavoces abiertos el micrófono se oye a sí + /// mismo y el asistente se interrumpe solo. Actívalo con auriculares o con + /// cancelación de eco del sistema. + pub barge_in: bool, + /// Con barge-in activo, cuánto más fuerte que el umbral normal debe sonar + /// la voz para cortar. Sube el listón frente al eco del altavoz. + pub barge_in_factor: f32, +} + +impl Default for VadConfig { + fn default() -> Self { + Self { + frame_seconds: 0.1, + preroll_seconds: 0.3, + silence_hold: 0.8, + min_utterance: 0.3, + max_utterance: 20.0, + threshold_factor: 3.0, + min_threshold: 0.0008, + max_threshold: 0.02, + barge_in: false, + barge_in_factor: 4.0, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct AsrConfig { + /// Carpeta del modelo Canary en ONNX. + pub model_dir: PathBuf, + pub source_lang: String, + pub target_lang: String, + /// Ventana deslizante de las transcripciones provisionales. + pub window: f32, + /// Avance de la ventana entre decodificaciones. + pub step: f32, + /// Ventanas que deben coincidir para dar una palabra por estable. + pub stability: usize, + /// Emitir transcripciones provisionales. Cuestan CPU y sólo sirven para + /// verlas en pantalla: el turno se dispara con la final. + pub partials: bool, + /// Carga una segunda instancia del modelo para que las decodificaciones + /// finales no bloqueen a las provisionales. Duplica la memoria. + pub dedicated_final_model: bool, + /// Proveedor de ejecución de ONNX Runtime: `cpu`, `cuda`, ... + pub execution_provider: String, + pub inter_threads: usize, + pub intra_threads: usize, +} + +impl Default for AsrConfig { + fn default() -> Self { + Self { + model_dir: PathBuf::from("vendor/canary-rs/models/canary-180m-flash-onnx"), + source_lang: "es".into(), + target_lang: "es".into(), + window: 6.0, + step: 0.4, + stability: 2, + partials: true, + dedicated_final_model: false, + // CPU a propósito: la GPU de 4 GB está ocupada por el hablante del + // TTS, y disputársela sale más caro que decodificar en CPU. + execution_provider: "cpu".into(), + inter_threads: 2, + intra_threads: 4, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct LlmConfig { + pub host: String, + pub port: u16, + pub model: String, + pub temperature: f32, + pub top_p: f32, + pub top_k: i32, + pub min_p: f32, + pub repeat_penalty: f32, + pub max_tokens: u32, + /// Vueltas máximas del bucle de herramientas antes de rendirse. + pub max_tool_rounds: usize, + pub request_timeout_secs: u64, +} + +impl Default for LlmConfig { + fn default() -> Self { + Self { + host: "127.0.0.1".into(), + port: 8012, + model: String::new(), + temperature: 0.7, + top_p: 0.9, + top_k: 40, + min_p: 0.1, + repeat_penalty: 1.1, + // Una respuesta hablada larga cansa; el recorte también acota el + // coste de la síntesis, que es la etapa lenta. + max_tokens: 300, + max_tool_rounds: 4, + request_timeout_secs: 120, + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct TtsConfig { + pub host: String, + pub port: u16, + /// Voz registrada en el servidor, o una del modelo. + pub voice: String, + /// Voz clonada que se registra al arrancar, si se define. + pub reference: Option<ReferenceVoice>, + pub language: String, + pub temperature: f32, + pub top_k: i32, + pub top_p: f32, + pub repetition_penalty: f32, + pub max_new_tokens: u32, + /// Sintetiza una frase corta al arrancar. La primera petición paga la + /// construcción de los grafos: ~3,5 s que conviene no gastar en el + /// primer turno real. + pub warmup: bool, + pub request_timeout_secs: u64, +} + +impl Default for TtsConfig { + fn default() -> Self { + Self { + host: "127.0.0.1".into(), + port: 8013, + voice: "asistente".into(), + reference: None, + language: "spanish".into(), + temperature: 0.9, + top_k: 50, + top_p: 1.0, + repetition_penalty: 1.05, + max_new_tokens: 2048, + warmup: true, + request_timeout_secs: 180, + } + } +} + +/// Latentes de una voz clonada, tal y como los produce `qwen-codec --talker`. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ReferenceVoice { + /// Nombre con el que se registra en el servidor. + pub name: String, + /// Embedding del hablante (`.spk`). + pub speaker: PathBuf, + /// Códigos de referencia (`.rvq`), que activan el clonado ICL. + pub codes: PathBuf, + /// Transcripción de la referencia (`.txt`). + pub transcript: PathBuf, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct ToolsConfig { + /// Deja que el modelo llame a herramientas. + pub enabled: bool, + /// Usa `general.tools_prompt` en solitario para la pasada en que el + /// modelo decide si llamar a una herramienta, y `general.system_prompt` + /// sólo para redactar la respuesta hablada. + /// + /// Es lo único que hace que las herramientas funcionen de verdad con este + /// modelo (ver `tools_prompt`), a cambio de que las respuestas que no + /// usan herramienta pierdan la guía de estilo. Se compensa en parte + /// porque el limpiador de texto quita el markdown antes de hablar. + /// Ponlo a `false` para priorizar el estilo sobre las herramientas. + pub dedicated_prompt: bool, + /// Habilita la herramienta de ejecución de órdenes del sistema. + /// + /// Apagada por defecto a conciencia: darle una shell a un modelo que + /// obedece a lo que oye por el micrófono es un cambio de postura de + /// seguridad, no una opción de comodidad. + pub shell: bool, + /// Órdenes admitidas, comparadas contra el ejecutable (argv[0]). + /// Una lista vacía deniega todo aunque `shell` esté activo. + pub shell_allowlist: Vec<String>, + /// Segundos que puede durar una orden antes de que se la mate. + pub shell_timeout_secs: u64, + /// Registra la orden pero no la ejecuta. Útil para estrenar la lista blanca. + pub shell_dry_run: bool, + /// Directorio de trabajo de las órdenes; vacío = el del proceso. + pub shell_working_dir: String, +} + +impl Default for ToolsConfig { + fn default() -> Self { + Self { + enabled: true, + dedicated_prompt: true, + shell: false, + shell_allowlist: vec![ + "date".into(), + "uptime".into(), + "free".into(), + "df".into(), + "ls".into(), + ], + shell_timeout_secs: 10, + shell_dry_run: false, + shell_working_dir: String::new(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct SupervisorConfig { + /// Lanza los servidores en lugar de suponer que ya están arriba. + pub manage: bool, + /// Segundos de espera a que un servidor conteste a `/health`. + pub startup_timeout_secs: u64, + pub llama: LlamaProcess, + pub tts: TtsProcess, +} + +impl Default for SupervisorConfig { + fn default() -> Self { + Self { + manage: true, + startup_timeout_secs: 180, + llama: LlamaProcess::default(), + tts: TtsProcess::default(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct LlamaProcess { + pub binary: PathBuf, + pub model: PathBuf, + pub mmproj: PathBuf, + /// Plantilla de chat que se pasa con `--chat-template-file`. + /// + /// La del modelo abre `<think>` sin cerrarlo nunca, y el asistente se pasa + /// entonces varios segundos razonando antes de decir la primera palabra. + /// Esta copia deja el bloque cerrado de entrada. + pub chat_template: PathBuf, + pub extra_args: Vec<String>, +} + +impl Default for LlamaProcess { + fn default() -> Self { + Self { + binary: PathBuf::from("vendor/llama.cpp/build/bin/llama-server"), + model: PathBuf::from("models/Qwen3.5-2B.Q5_K_M.gguf"), + mmproj: PathBuf::from("models/mmproj-BF16.gguf"), + chat_template: PathBuf::from("config/qwen35-no-think.jinja"), + extra_args: [ + "--threads", + "10", + "--threads-batch", + "10", + "--batch-size", + "512", + "--ubatch-size", + "256", + "--gpu-layers", + "10", + "--split-mode", + "layer", + "--tensor-split", + "1", + "--main-gpu", + "0", + "--no-mmap", + "--ctx-size", + "8192", + "--parallel", + "2", + "--cache-ram", + "6144", + "--rope-freq-base", + "1000000", + "--rope-freq-scale", + "0.25", + "--jinja", + ] + .iter() + .map(|s| s.to_string()) + .collect(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct TtsProcess { + pub binary: PathBuf, + pub model: PathBuf, + pub codec: PathBuf, + /// Segundos de audio que el codec acumula antes de decodificar un bloque. + /// + /// El valor de fábrica son 24 s, que en la práctica significa «no + /// devuelvas nada hasta terminar la frase entera»: medido, baja el primer + /// audio de 4,9 s a 0,6 s. Es el ajuste con más efecto de todo el sistema. + pub codec_chunk_dur: f32, + pub extra_args: Vec<String>, +} + +impl Default for TtsProcess { + fn default() -> Self { + Self { + binary: PathBuf::from("vendor/qwentts.cpp/build/tts-server"), + model: PathBuf::from("models/qwen-talker-1.7b-base-Q8_0.gguf"), + codec: PathBuf::from("models/qwen-tokenizer-12hz-Q8_0.gguf"), + codec_chunk_dur: 1.0, + extra_args: Vec::new(), + } + } +} + +impl Config { + pub fn load(path: impl AsRef<Path>) -> Result<Self> { + let path = path.as_ref(); + let raw = std::fs::read_to_string(path) + .map_err(|e| Error::Config(format!("no se pudo leer {}: {e}", path.display())))?; + let mut config: Config = + toml::from_str(&raw).map_err(|e| Error::Config(format!("{}: {e}", path.display())))?; + // Las rutas relativas se resuelven contra la carpeta del TOML, no + // contra el directorio desde el que se lanza el binario. + if let Some(base) = path.parent().filter(|p| !p.as_os_str().is_empty()) { + config.rebase(base); + } + config.validate()?; + Ok(config) + } + + /// Reinterpreta las rutas relativas respecto de `base`. + pub fn rebase(&mut self, base: &Path) { + let fix = |p: &mut PathBuf| { + if p.is_relative() { + *p = base.join(&*p); + } + }; + fix(&mut self.asr.model_dir); + fix(&mut self.supervisor.llama.binary); + fix(&mut self.supervisor.llama.model); + fix(&mut self.supervisor.llama.mmproj); + fix(&mut self.supervisor.llama.chat_template); + fix(&mut self.supervisor.tts.binary); + fix(&mut self.supervisor.tts.model); + fix(&mut self.supervisor.tts.codec); + if let Some(reference) = self.tts.reference.as_mut() { + fix(&mut reference.speaker); + fix(&mut reference.codes); + fix(&mut reference.transcript); + } + } + + fn validate(&self) -> Result<()> { + if self.vad.silence_hold <= 0.0 { + return Err(Error::Config("vad.silence_hold debe ser > 0".into())); + } + if self.vad.min_utterance >= self.vad.max_utterance { + return Err(Error::Config( + "vad.min_utterance debe ser menor que vad.max_utterance".into(), + )); + } + if self.asr.step <= 0.0 || self.asr.window <= self.asr.step { + return Err(Error::Config( + "asr.window debe ser mayor que asr.step, y ambos > 0".into(), + )); + } + if self.tools.shell && self.tools.shell_allowlist.is_empty() { + return Err(Error::Config( + "tools.shell está activo pero tools.shell_allowlist está vacía: \ + declara las órdenes permitidas o desactiva tools.shell" + .into(), + )); + } + Ok(()) + } + + pub fn llm_authority(&self) -> String { + format!("{}:{}", self.llm.host, self.llm.port) + } + + pub fn tts_authority(&self) -> String { + format!("{}:{}", self.tts.host, self.tts.port) + } +} diff --git a/crates/asist-core/src/error.rs b/crates/asist-core/src/error.rs new file mode 100644 index 0000000..b559bee --- /dev/null +++ b/crates/asist-core/src/error.rs @@ -0,0 +1,82 @@ +use std::fmt; + +pub type Result<T> = std::result::Result<T, Error>; + +/// Un fallo de una etapa del pipeline. El orquestador decide si es fatal o si +/// basta con abortar el turno en curso, así que el error lleva esa distinción +/// en lugar de dejarla al criterio de quien lo recibe. +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("configuración inválida: {0}")] + Config(String), + + #[error("audio: {0}")] + Audio(String), + + #[error("ASR: {0}")] + Asr(String), + + #[error("LLM: {0}")] + Llm(String), + + #[error("TTS: {0}")] + Tts(String), + + #[error("HTTP {status} en {url}: {body}")] + Http { + status: u16, + url: String, + body: String, + }, + + #[error("transporte hacia {url}: {source}")] + Transport { + url: String, + #[source] + source: std::io::Error, + }, + + #[error("herramienta «{tool}»: {message}")] + Tool { tool: String, message: String }, + + #[error("operación cancelada")] + Cancelled, + + #[error(transparent)] + Json(#[from] serde_json::Error), + + #[error(transparent)] + Io(#[from] std::io::Error), +} + +impl Error { + /// Un error fatal tumba el asistente; el resto sólo aborta el turno. + /// + /// Los servidores locales se caen y se reinician, y el usuario prefiere + /// oír «no te he entendido» a que el proceso muera, así que sólo la + /// configuración y el audio (que no se pueden reintentar) son fatales. + pub fn is_fatal(&self) -> bool { + matches!(self, Error::Config(_) | Error::Audio(_)) + } + + /// `true` cuando reintentar tiene sentido: el servidor aún está cargando + /// el modelo, o se ha caído la conexión a mitad de una petición. + pub fn is_retryable(&self) -> bool { + match self { + Error::Transport { .. } => true, + Error::Http { status, .. } => *status == 503 || *status == 429 || *status >= 500, + _ => false, + } + } +} + +/// Contexto legible para los `Result` que cruzan una frontera de crate. +pub trait Context<T> { + fn ctx(self, f: impl FnOnce() -> String) -> Result<T>; +} + +impl<T, E: fmt::Display> Context<T> for std::result::Result<T, E> { + fn ctx(self, f: impl FnOnce() -> String) -> Result<T> { + self.map_err(|e| Error::Config(format!("{}: {}", f(), e))) + } +} diff --git a/crates/asist-core/src/event.rs b/crates/asist-core/src/event.rs new file mode 100644 index 0000000..549bec0 --- /dev/null +++ b/crates/asist-core/src/event.rs @@ -0,0 +1,127 @@ +use std::time::{Duration, Instant}; + +/// Identifica un turno de conversación (una intervención del usuario y la +/// respuesta que provoca). Todo lo que viaja por el bus lo lleva, para que un +/// resultado que llega tarde no se confunda con el turno que ya está en curso +/// —el caso típico cuando el usuario interrumpe al asistente. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Default)] +pub struct TurnId(pub u64); + +impl std::fmt::Display for TurnId { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!(f, "#{}", self.0) + } +} + +/// Motivo por el que se corta la reproducción en curso. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum InterruptReason { + /// El usuario ha empezado a hablar encima del asistente (barge-in). + UserSpoke, + /// Petición explícita: tecla, señal o cierre. + Requested, +} + +/// Todo lo que las etapas se cuentan entre sí. Un único enum mantiene el bus +/// observable: el renderizador y la telemetría ven exactamente los mismos +/// eventos que el orquestador, sin canales paralelos que se desincronicen. +#[derive(Debug, Clone)] +pub enum Event { + /// El VAD ha detectado el arranque de una intervención. + SpeechStarted { turn: TurnId, at: Instant }, + + /// Transcripción provisional de la ventana deslizante: `committed` ya es + /// estable, `volatile` todavía puede cambiar. + Partial { + turn: TurnId, + committed: String, + volatile: String, + }, + + /// Transcripción definitiva de la intervención completa. + Transcript { + turn: TurnId, + text: String, + audio_secs: f32, + decode: Duration, + }, + + /// La intervención no contenía nada transcribible. + Discarded { turn: TurnId }, + + /// Primer fragmento de texto que devuelve el modelo. + ReplyStarted { turn: TurnId, ttft: Duration }, + + /// Un trozo de respuesta según llega del modelo. + ReplyDelta { turn: TurnId, text: String }, + + /// Una frase completa lista para sintetizar. + Sentence { + turn: TurnId, + index: usize, + text: String, + }, + + /// Respuesta completa del modelo para este turno. + ReplyDone { turn: TurnId, text: String }, + + /// El modelo ha pedido ejecutar una herramienta. + ToolRequested { + turn: TurnId, + name: String, + arguments: String, + }, + + /// Resultado de esa ejecución. + ToolFinished { + turn: TurnId, + name: String, + ok: bool, + output: String, + took: Duration, + }, + + /// Primer audio audible del turno: la métrica que de verdad percibe quien habla. + AudioStarted { turn: TurnId, latency: Duration }, + + /// Se ha terminado de reproducir todo lo del turno. + AudioFinished { turn: TurnId }, + + /// Se ha cortado la reproducción. + Interrupted { + turn: TurnId, + reason: InterruptReason, + }, + + /// Aviso no fatal; el turno continúa. + Warning { turn: TurnId, message: String }, + + /// Fallo que aborta el turno. + Failed { turn: TurnId, message: String }, + + /// Cierre ordenado. + Shutdown, +} + +impl Event { + pub fn turn(&self) -> Option<TurnId> { + match self { + Event::SpeechStarted { turn, .. } + | Event::Partial { turn, .. } + | Event::Transcript { turn, .. } + | Event::Discarded { turn } + | Event::ReplyStarted { turn, .. } + | Event::ReplyDelta { turn, .. } + | Event::Sentence { turn, .. } + | Event::ReplyDone { turn, .. } + | Event::ToolRequested { turn, .. } + | Event::ToolFinished { turn, .. } + | Event::AudioStarted { turn, .. } + | Event::AudioFinished { turn } + | Event::Interrupted { turn, .. } + | Event::Warning { turn, .. } + | Event::Failed { turn, .. } => Some(*turn), + Event::Shutdown => None, + } + } +} diff --git a/crates/asist-core/src/http.rs b/crates/asist-core/src/http.rs new file mode 100644 index 0000000..32b6680 --- /dev/null +++ b/crates/asist-core/src/http.rs @@ -0,0 +1,351 @@ +//! Cliente HTTP/1.1 mínimo para hablar con los servidores locales. +//! +//! Los tres motores (llama-server, tts-server) escuchan en localhost sobre HTTP +//! plano, así que no hace falta TLS ni un runtime asíncrono. A cambio de las +//! ~300 líneas de aquí se gana lo único que ninguna librería genérica da +//! cómodo y que este pipeline necesita de verdad: **abortar una respuesta a +//! media descarga**. Cuando el usuario interrumpe al asistente hay que dejar de +//! leer audio ya generado y cerrar el socket en ese mismo instante; con un +//! cliente que sólo expone `read_to_end` eso no se puede hacer. + +use std::io::{BufRead, BufReader, Read, Write}; +use std::net::TcpStream; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use crate::error::{Error, Result}; + +/// Bandera compartida que aborta una lectura en curso. +/// +/// Se comprueba entre bloques, así que el corte ocurre como muy tarde un +/// `read` después de activarla —del orden de milisegundos con los tamaños de +/// bloque que usan los servidores. +#[derive(Clone, Default)] +pub struct Cancel(Arc<AtomicBool>); + +impl Cancel { + pub fn new() -> Self { + Self::default() + } + + pub fn cancel(&self) { + self.0.store(true, Ordering::SeqCst); + } + + pub fn is_cancelled(&self) -> bool { + self.0.load(Ordering::SeqCst) + } + + pub fn reset(&self) { + self.0.store(false, Ordering::SeqCst); + } +} + +#[derive(Clone, Debug)] +pub struct HttpClient { + /// `host:puerto` del servidor local. + pub authority: String, + pub connect_timeout: Duration, + /// Tiempo máximo sin recibir un solo byte. No es el tiempo total: una + /// síntesis larga puede tardar minutos siempre que siga fluyendo. + pub read_timeout: Duration, +} + +impl HttpClient { + pub fn new(authority: impl Into<String>) -> Self { + Self { + authority: authority.into(), + connect_timeout: Duration::from_secs(5), + read_timeout: Duration::from_secs(120), + } + } + + pub fn with_read_timeout(mut self, timeout: Duration) -> Self { + self.read_timeout = timeout; + self + } + + fn connect(&self, path: &str) -> Result<TcpStream> { + let addr = resolve(&self.authority, path)?; + let stream = TcpStream::connect_timeout(&addr, self.connect_timeout).map_err(|e| { + Error::Transport { + url: format!("http://{}{}", self.authority, path), + source: e, + } + })?; + stream.set_read_timeout(Some(self.read_timeout))?; + // Los cuerpos son pequeños y la latencia manda: nada de Nagle. + let _ = stream.set_nodelay(true); + Ok(stream) + } + + fn send(&self, method: &str, path: &str, body: Option<&[u8]>) -> Result<Response> { + let mut stream = self.connect(path)?; + let mut head = format!( + "{method} {path} HTTP/1.1\r\nHost: {}\r\nConnection: close\r\nAccept: */*\r\n", + self.authority + ); + if let Some(body) = body { + head.push_str("Content-Type: application/json\r\n"); + head.push_str(&format!("Content-Length: {}\r\n", body.len())); + } + head.push_str("\r\n"); + + let url = format!("http://{}{}", self.authority, path); + let io = |e: std::io::Error| Error::Transport { + url: url.clone(), + source: e, + }; + stream.write_all(head.as_bytes()).map_err(io)?; + if let Some(body) = body { + stream.write_all(body).map_err(io)?; + } + stream.flush().map_err(io)?; + + Response::read_head(BufReader::new(stream), url) + } + + /// GET que devuelve el cuerpo completo ya decodificado. + pub fn get(&self, path: &str) -> Result<Vec<u8>> { + self.send("GET", path, None)?.into_body() + } + + /// POST de JSON que devuelve el cuerpo completo. + pub fn post_json(&self, path: &str, body: &serde_json::Value) -> Result<Vec<u8>> { + let payload = serde_json::to_vec(body)?; + self.send("POST", path, Some(&payload))?.into_body() + } + + /// POST de JSON cuya respuesta se consume a trozos según llega. + /// + /// `on_chunk` recibe cada bloque en cuanto está disponible y devuelve + /// `false` para cortar; junto a `cancel` son las dos vías por las que el + /// barge-in detiene una síntesis a medias. + pub fn post_json_streaming( + &self, + path: &str, + body: &serde_json::Value, + cancel: &Cancel, + mut on_chunk: impl FnMut(&[u8]) -> bool, + ) -> Result<()> { + let payload = serde_json::to_vec(body)?; + let mut response = self.send("POST", path, Some(&payload))?; + response.error_for_status()?; + response.stream(cancel, &mut on_chunk) + } + + /// POST de JSON que entrega la respuesta línea a línea (SSE de llama-server). + pub fn post_json_lines( + &self, + path: &str, + body: &serde_json::Value, + cancel: &Cancel, + mut on_line: impl FnMut(&str) -> bool, + ) -> Result<()> { + let mut pending = Vec::new(); + self.post_json_streaming(path, body, cancel, |chunk| { + pending.extend_from_slice(chunk); + while let Some(nl) = pending.iter().position(|b| *b == b'\n') { + let line: Vec<u8> = pending.drain(..=nl).collect(); + let line = String::from_utf8_lossy(&line); + if !on_line(line.trim_end_matches(['\r', '\n'])) { + return false; + } + } + true + }) + } + + /// Sondea `/health` y devuelve `true` si el servidor está listo. + pub fn healthy(&self, path: &str) -> bool { + self.get(path).is_ok() + } +} + +fn resolve(authority: &str, path: &str) -> Result<std::net::SocketAddr> { + use std::net::ToSocketAddrs; + authority + .to_socket_addrs() + .map_err(|e| Error::Transport { + url: format!("http://{authority}{path}"), + source: e, + })? + .next() + .ok_or_else(|| Error::Config(format!("no se pudo resolver «{authority}»"))) +} + +/// Codificación del cuerpo anunciada por el servidor. +enum Body { + Chunked, + Length(usize), + /// Sin longitud ni chunked: el cuerpo acaba cuando cierra el socket. + UntilClose, +} + +pub struct Response { + pub status: u16, + url: String, + reader: BufReader<TcpStream>, + body: Body, +} + +impl Response { + fn read_head(mut reader: BufReader< |