//! 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), } /// Cómo terminó una respuesta en streaming. #[derive(Debug, Clone, Default)] pub struct StreamOutcome { pub text: String, pub tool_calls: Vec, /// Tiempo hasta el primer fragmento con contenido. pub ttft: Option, 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::(&body) else { return; }; let Some(template) = props.get("chat_template").and_then(Value::as_str) else { return; }; let opens = template.matches("").count(); let closes = template.matches("").count(); if opens > closes { tracing::warn!( target: "llm", "la plantilla de chat abre sin cerrarlo: el modelo razonará \ varios segundos antes de contestar. Arranca llama-server con \ --chat-template-file config/qwen35-no-think.jinja" ); } } /// Pregunta al modelo por una imagen, en una petición aparte de la /// conversación. /// /// El servidor tiene cargado el proyector multimodal (`--mmproj`), así que /// acepta partes `image_url` con la imagen en base64. Se hace fuera del /// historial a propósito: una imagen ocupa cientos de tokens de contexto y /// arrastrarla turno tras turno saldría carísimo para lo poco que aporta /// una vez descrita. Lo que vuelve a la conversación es el texto. /// /// Medido en esta máquina, el coste depende mucho del tamaño: 1,3 s a /// 320x240, 2,9 s a 640x480 y 7,8 s a 1280x720. pub fn look(&self, image_jpeg: &[u8], question: &str, cancel: &Cancel) -> Result { use base64::Engine; let encoded = base64::engine::general_purpose::STANDARD.encode(image_jpeg); let mut body = json!({ "messages": [{ "role": "user", "content": [ { "type": "text", "text": question }, { "type": "image_url", "image_url": { "url": format!("data:image/jpeg;base64,{encoded}") } } ] }], "stream": true, "temperature": self.config.temperature, "max_tokens": self.config.max_tokens, }); if !self.config.model.is_empty() { body["model"] = json!(self.config.model); } // En streaming aunque no se use el texto parcial: una petición de // visión tarda segundos y sin flujo el socket puede quedarse callado // hasta pasado el plazo de lectura. let started = Instant::now(); let mut text = String::new(); 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::(payload) else { return true; }; let Some(choice) = event.get("choices").and_then(|c| c.get(0)) else { return true; }; if let Some(chunk) = choice .get("delta") .and_then(|d| d.get("content")) .and_then(Value::as_str) { text.push_str(chunk); } if choice.get("finish_reason").is_some_and(|r| !r.is_null()) { finished = true; return false; } true }); match result { Ok(()) => {} Err(Error::Cancelled) if finished || cancel.is_cancelled() => {} Err(err) => return Err(err), } tracing::debug!( target: "llm", bytes = image_jpeg.len(), ms = started.elapsed().as_millis(), "visión" ); let text = text.trim().to_string(); if text.is_empty() { return Err(Error::Llm("el modelo no describió la imagen".into())); } Ok(text) } 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 { 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::(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, } #[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 { 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"); } }