aboutsummaryrefslogtreecommitdiffstats
path: root/crates/asist-app/src/supervisor.rs
diff options
context:
space:
mode:
Diffstat (limited to 'crates/asist-app/src/supervisor.rs')
-rw-r--r--crates/asist-app/src/supervisor.rs228
1 files changed, 228 insertions, 0 deletions
diff --git a/crates/asist-app/src/supervisor.rs b/crates/asist-app/src/supervisor.rs
new file mode 100644
index 0000000..17360da
--- /dev/null
+++ b/crates/asist-app/src/supervisor.rs
@@ -0,0 +1,228 @@
+//! Arranque y parada de los servidores locales.
+//!
+//! El asistente puede levantar `llama-server` y `tts-server` él mismo, para
+//! que ponerlo en marcha sea un solo comando. Los procesos se lanzan en su
+//! propio grupo y se paran con SIGTERM antes de recurrir a SIGKILL: matar en
+//! seco a llama-server deja la GPU ocupada hasta que el driver la recupera.
+
+use std::path::Path;
+use std::process::{Child, Command, Stdio};
+use std::time::{Duration, Instant};
+
+use asist_core::config::{Config, LlamaProcess, TtsProcess};
+use asist_core::error::{Error, Result};
+
+pub struct Supervisor {
+ children: Vec<Managed>,
+ log_dir: std::path::PathBuf,
+}
+
+struct Managed {
+ name: &'static str,
+ child: Child,
+}
+
+impl Supervisor {
+ pub fn new(log_dir: impl AsRef<Path>) -> Result<Self> {
+ let log_dir = log_dir.as_ref().to_path_buf();
+ std::fs::create_dir_all(&log_dir)?;
+ Ok(Self {
+ children: Vec::new(),
+ log_dir,
+ })
+ }
+
+ /// Lanza los servidores que la configuración pida y que no estén ya arriba.
+ ///
+ /// Reutilizar uno que ya escucha es deliberado: durante el desarrollo se
+ /// reinicia el asistente muchas veces y volver a cargar los modelos
+ /// cuesta más de un minuto.
+ pub fn start(&mut self, config: &Config, llm_up: bool, tts_up: bool) -> Result<()> {
+ if !config.supervisor.manage {
+ return Ok(());
+ }
+ if llm_up {
+ tracing::info!(target: "supervisor", "llama-server ya está escuchando, se reutiliza");
+ } else {
+ let command = llama_command(&config.supervisor.llama, config)?;
+ self.spawn("llama-server", command)?;
+ }
+ if tts_up {
+ tracing::info!(target: "supervisor", "tts-server ya está escuchando, se reutiliza");
+ } else {
+ let command = tts_command(&config.supervisor.tts, config)?;
+ self.spawn("tts-server", command)?;
+ }
+ Ok(())
+ }
+
+ fn spawn(&mut self, name: &'static str, mut command: Command) -> Result<()> {
+ let log_path = self.log_dir.join(format!("{name}.log"));
+ let log = std::fs::File::create(&log_path)?;
+ let errors = log.try_clone()?;
+
+ // La salida va al fichero: los modelos escupen cientos de líneas y
+ // taparían la transcripción en pantalla. Cuando algo falla, el
+ // mensaje de error apunta aquí.
+ command
+ .stdin(Stdio::null())
+ .stdout(Stdio::from(log))
+ .stderr(Stdio::from(errors));
+
+ let child = command.spawn().map_err(|e| {
+ Error::Config(format!(
+ "no se pudo lanzar {name} ({:?}): {e}",
+ command.get_program()
+ ))
+ })?;
+ tracing::info!(
+ target: "supervisor",
+ proceso = name,
+ pid = child.id(),
+ registro = %log_path.display(),
+ "lanzado"
+ );
+ self.children.push(Managed { name, child });
+ Ok(())
+ }
+
+ /// Comprueba si alguno se ha muerto solo, y devuelve su nombre.
+ pub fn crashed(&mut self) -> Option<(&'static str, Option<i32>)> {
+ for managed in &mut self.children {
+ if let Ok(Some(status)) = managed.child.try_wait() {
+ return Some((managed.name, status.code()));
+ }
+ }
+ None
+ }
+
+ pub fn log_path(&self, name: &str) -> std::path::PathBuf {
+ self.log_dir.join(format!("{name}.log"))
+ }
+
+ /// Para todo lo lanzado. SIGTERM primero para que liberen la GPU.
+ pub fn shutdown(&mut self) {
+ for managed in &mut self.children {
+ terminate(&mut managed.child, managed.name);
+ }
+ self.children.clear();
+ }
+}
+
+impl Drop for Supervisor {
+ fn drop(&mut self) {
+ self.shutdown();
+ }
+}
+
+fn terminate(child: &mut Child, name: &str) {
+ if matches!(child.try_wait(), Ok(Some(_))) {
+ return;
+ }
+ // SIGTERM directo: `Child::kill` manda SIGKILL, que no da al servidor
+ // ocasión de soltar la memoria de la GPU.
+ #[cfg(unix)]
+ unsafe {
+ libc_kill(child.id() as i32, 15);
+ }
+ #[cfg(not(unix))]
+ let _ = child.kill();
+
+ let deadline = Instant::now() + Duration::from_secs(10);
+ while Instant::now() < deadline {
+ match child.try_wait() {
+ Ok(Some(_)) => {
+ tracing::info!(target: "supervisor", proceso = name, "detenido");
+ return;
+ }
+ Ok(None) => std::thread::sleep(Duration::from_millis(100)),
+ Err(_) => break,
+ }
+ }
+ tracing::warn!(target: "supervisor", proceso = name, "no respondió a SIGTERM, se fuerza");
+ let _ = child.kill();
+ let _ = child.wait();
+}
+
+#[cfg(unix)]
+unsafe fn libc_kill(pid: i32, signal: i32) {
+ // Se declara aquí el único símbolo de libc que hace falta, en vez de
+ // arrastrar el crate entero por una llamada.
+ unsafe extern "C" {
+ fn kill(pid: i32, sig: i32) -> i32;
+ }
+ unsafe {
+ kill(pid, signal);
+ }
+}
+
+fn require(path: &Path, what: &str) -> Result<()> {
+ if path.exists() {
+ return Ok(());
+ }
+ Err(Error::Config(format!(
+ "falta {what}: {}. Ejecuta scripts/bootstrap.sh para compilar los \
+ motores y enlazar los modelos",
+ path.display()
+ )))
+}
+
+fn llama_command(process: &LlamaProcess, config: &Config) -> Result<Command> {
+ require(&process.binary, "el binario de llama-server")?;
+ require(&process.model, "el modelo del LLM")?;
+
+ let mut command = Command::new(&process.binary);
+ command
+ .arg("--model")
+ .arg(&process.model)
+ .arg("--host")
+ .arg(&config.llm.host)
+ .arg("--port")
+ .arg(config.llm.port.to_string());
+
+ if process.mmproj.exists() {
+ command.arg("--mmproj").arg(&process.mmproj);
+ }
+ // Sin esto, la plantilla del modelo abre <think> y no lo cierra nunca: el
+ // asistente se pasa entre 7 y 9 s razonando antes de la primera palabra.
+ if process.chat_template.exists() {
+ command
+ .arg("--jinja")
+ .arg("--chat-template-file")
+ .arg(&process.chat_template);
+ } else {
+ tracing::warn!(
+ target: "supervisor",
+ ruta = %process.chat_template.display(),
+ "no está la plantilla sin razonamiento: el modelo tardará varios \
+ segundos en empezar a hablar"
+ );
+ }
+ command.args(&process.extra_args);
+ Ok(command)
+}
+
+fn tts_command(process: &TtsProcess, config: &Config) -> Result<Command> {
+ require(&process.binary, "el binario de tts-server")?;
+ require(&process.model, "el modelo hablante del TTS")?;
+ require(&process.codec, "el codec del TTS")?;
+
+ let mut command = Command::new(&process.binary);
+ command
+ .arg("--model")
+ .arg(&process.model)
+ .arg("--codec")
+ .arg(&process.codec)
+ .arg("--host")
+ .arg(&config.tts.host)
+ .arg("--port")
+ .arg(config.tts.port.to_string())
+ .arg("--lang")
+ .arg(&config.tts.language)
+ // El ajuste con más efecto de todo el sistema: de fábrica son 24 s, que
+ // en la práctica significa no devolver nada hasta terminar la frase.
+ .arg("--codec-chunk-dur")
+ .arg(process.codec_chunk_dur.to_string());
+ command.args(&process.extra_args);
+ Ok(command)
+}