aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-core
diff options
context:
space:
mode:
Diffstat (limited to 'crates/asist-core')
-rw-r--r--crates/asist-core/Cargo.toml15
-rw-r--r--crates/asist-core/src/config.rs508
-rw-r--r--crates/asist-core/src/error.rs82
-rw-r--r--crates/asist-core/src/event.rs127
-rw-r--r--crates/asist-core/src/http.rs351
-rw-r--r--crates/asist-core/src/lib.rs17
-rw-r--r--crates/asist-core/src/telemetry.rs255
-rw-r--r--crates/asist-core/src/text.rs330
-rw-r--r--crates/asist-core/src/tools.rs563
9 files changed, 2248 insertions, 0 deletions
diff --git a/crates/asist-core/Cargo.toml b/crates/asist-core/Cargo.toml
new file mode 100644
index 0000000..bd8c3ae
--- /dev/null
+++ b/crates/asist-core/Cargo.toml
@@ -0,0 +1,15 @@
+[package]
+name = "asist-core"
+version.workspace = true
+edition.workspace = true
+rust-version.workspace = true
+license.workspace = true
+
+[dependencies]
+anyhow.workspace = true
+crossbeam-channel.workspace = true
+serde.workspace = true
+serde_json.workspace = true
+thiserror.workspace = true
+toml.workspace = true
+tracing.workspace = true
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<TcpStream>, url: String) -> Result<Self> {
+ let io = |e: std::io::Error| Error::Transport {
+ url: url.clone(),
+ source: e,
+ };
+
+ let mut line = String::new();
+ reader.read_line(&mut line).map_err(io)?;
+ let status = line
+ .split_whitespace()
+ .nth(1)
+ .and_then(|s| s.parse::<u16>().ok())
+ .ok_or_else(|| Error::Config(format!("respuesta HTTP ilegible de {url}: {line:?}")))?;
+
+ let mut body = Body::UntilClose;
+ loop {
+ let mut header = String::new();
+ if reader.read_line(&mut header).map_err(io)? == 0 {
+ break;
+ }
+ let header = header.trim_end();
+ if header.is_empty() {
+ break;
+ }
+ let Some((name, value)) = header.split_once(':') else {
+ continue;
+ };
+ let (name, value) = (name.trim().to_ascii_lowercase(), value.trim());
+ match name.as_str() {
+ "transfer-encoding" if value.eq_ignore_ascii_case("chunked") => {
+ body = Body::Chunked;
+ }
+ "content-length" => {
+ if let Ok(n) = value.parse() {
+ body = Body::Length(n);
+ }
+ }