diff --git a/crates/hammer-core/src/proto.rs b/crates/hammer-core/src/proto.rs index ea4f2834..7afa8afa 100644 --- a/crates/hammer-core/src/proto.rs +++ b/crates/hammer-core/src/proto.rs @@ -117,8 +117,10 @@ pub enum Event { }, /// La línea fue enviada al FIFO de control. InitAck { cmd: String }, - /// Un servicio supervisado murió. (Fase 5 sólo lo declara; la supervisión real llega - /// con el init propio del track posterior.) + /// Un servicio supervisado murió. La supervisión real la provee el init propio + /// (arje, PID 1): difunde `arje_bus::BusEvent::EnteCrashed` y el sink + /// `hammerd::crashes` lo traduce a este evento. Falta sólo el adaptador de + /// transporte que conecta ambos buses (ver `hammerd::crashes`). Crashed { service: String, code: i32 }, /// Mutación detectada por el watcher. Hammerd lo re-emite por el bus después de /// registrarlo en el diario. diff --git a/crates/hammerd/src/crashes.rs b/crates/hammerd/src/crashes.rs new file mode 100644 index 00000000..a6673823 --- /dev/null +++ b/crates/hammerd/src/crashes.rs @@ -0,0 +1,131 @@ +//! Puente de crashes: la **fuente real** del `Event::Crashed` del bus de agente. +//! +//! Hasta ahora `Event::Crashed` sólo lo declaraba el protocolo y lo disparaban los +//! tests: la Fase 5 difirió la "supervisión real" al init propio del track posterior. +//! Ese init es **arje** (ver `docs/adr/0007-arje-como-init-propio.md`): como PID 1 +//! supervisa cada Ente, detecta su muerte en `on_death` y la difunde por su bus como +//! `BusEvent::EnteCrashed` (crate `arje-bus`, `BusRequest::Subscribe` + `BusPayload::Event`). +//! +//! Este módulo es el **sink** de ese flujo (B.2 "exponer el CRASHED a la capa de IA"): +//! traduce la señal de ciclo de vida normalizada al vocabulario del bus de agente de +//! hammer (`Event::Crashed`) y la bombea al [`EventBus`] → toda conexión del +//! `/run/agent.sock` (la capa de IA) la recibe. +//! +//! ## Frontera de acoplamiento +//! +//! La **traducción** y el **bucle de bombeo** viven acá y son agnósticos del transporte +//! —se testean sin un init vivo—. Lo único que falta para cerrar B.2 end-to-end es el +//! *adaptador de transporte*: un thread que conecte al socket de arje +//! (`$ARJE_BUS_SOCK`), mande un `Subscribe` y reenvíe cada `BusEvent` a un +//! [`std::sync::mpsc::Sender`] hacia [`pump`]. Ese adaptador es lo que +//! decide cómo los dos repos comparten el wire (dep directa a `arje-bus` vs. crate de +//! proto compartido vs. relectura del frame postcard) — decisión registrada en +//! `docs/09-trust-model.md` / el PLAN de arje, y validable sólo contra un init real. +//! +//! Las APIs quedan listas (`pub`) aunque el transporte aún no las invoque: ese es el +//! patrón "vocabulario declarado, flujo por cablear" que arje-zero usa en su `events.rs`. +#![allow(dead_code)] + +use std::sync::mpsc::Receiver; + +use hammer_core::proto::Event; + +use crate::events::EventBus; + +/// Señal de ciclo de vida normalizada desde el supervisor (arje), independiente del +/// wire concreto para no acoplar `hammerd` al crate `arje-bus`. El adaptador de +/// transporte construye estos valores a partir de `arje_bus::BusEvent`. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Lifecycle { + /// Identidad humana del Ente que cambió de estado (su `label` en arje). Es lo que + /// la capa de IA verá como `service`. + pub service: String, + pub status: LifeStatus, +} + +/// Estado terminal/transicional de un Ente, ya normalizado. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum LifeStatus { + /// Terminó de forma anómala. `code` es el exit-code; para muertes por señal el + /// adaptador usa la convención de shell `128 + signum` para que sea siempre ≠ 0. + Crashed { code: i32 }, + /// arje programó un reinicio tras backoff. Observabilidad — el `proto` actual no + /// tiene evento dedicado, así que no se traduce (todavía). + Restarting { delay_ms: u64 }, + /// Terminó limpio (exit 0) y no se reinicia. Sin evento en el `proto` actual. + Exited, +} + +/// Traduce una señal de ciclo de vida al evento del bus de agente. Devuelve `None` +/// para los estados que el `proto` actual no expone (reinicios / salidas limpias): +/// hoy sólo el crash es accionable por la IA. +pub fn to_event(life: &Lifecycle) -> Option { + match &life.status { + LifeStatus::Crashed { code } => Some(Event::Crashed { + service: life.service.clone(), + code: *code, + }), + LifeStatus::Restarting { .. } | LifeStatus::Exited => None, + } +} + +/// Bucle de bombeo: drena señales de la fuente (el adaptador de transporte de arje) y +/// publica las accionables en el `EventBus`. Bloquea hasta que la fuente se cierra; +/// pensado para correr en su propio thread, como los demás subsistemas de `hammerd`. +pub fn pump(source: Receiver, bus: EventBus) { + for life in source { + if let Some(ev) = to_event(&life) { + bus.publish(&ev); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn crash_traduce_a_event_crashed() { + let life = Lifecycle { + service: "web".into(), + status: LifeStatus::Crashed { code: 139 }, // 128 + SIGSEGV + }; + assert_eq!( + to_event(&life), + Some(Event::Crashed { service: "web".into(), code: 139 }) + ); + } + + #[test] + fn restarting_y_exited_no_producen_evento() { + let r = Lifecycle { service: "x".into(), status: LifeStatus::Restarting { delay_ms: 100 } }; + let e = Lifecycle { service: "x".into(), status: LifeStatus::Exited }; + assert_eq!(to_event(&r), None); + assert_eq!(to_event(&e), None); + } + + #[test] + fn pump_publica_solo_los_crashes_en_el_bus() { + let bus = EventBus::new(); + let sub = bus.subscribe(); + + let (tx, rx) = std::sync::mpsc::channel(); + // Mezcla: un exited (descartado), un crash (publicado), un restarting (descartado). + tx.send(Lifecycle { service: "a".into(), status: LifeStatus::Exited }).unwrap(); + tx.send(Lifecycle { service: "b".into(), status: LifeStatus::Crashed { code: 1 } }).unwrap(); + tx.send(Lifecycle { service: "c".into(), status: LifeStatus::Restarting { delay_ms: 5 } }).unwrap(); + drop(tx); // cierra la fuente → pump termina. + + pump(rx, bus); + + // Sólo el crash de "b" llegó al suscriptor. + match sub.try_recv().expect("debía haber un Crashed") { + Event::Crashed { service, code } => { + assert_eq!(service, "b"); + assert_eq!(code, 1); + } + other => panic!("esperaba Crashed, fue {other:?}"), + } + assert!(sub.try_recv().is_err(), "no debía publicarse nada más"); + } +} diff --git a/crates/hammerd/src/main.rs b/crates/hammerd/src/main.rs index f9a9a55f..1e989c33 100644 --- a/crates/hammerd/src/main.rs +++ b/crates/hammerd/src/main.rs @@ -14,6 +14,7 @@ use clap::Parser; mod bus; mod control; +mod crashes; mod events; mod watcher; diff --git a/docs/10-roadmap.md b/docs/10-roadmap.md index d1821303..4cd1f8f5 100644 --- a/docs/10-roadmap.md +++ b/docs/10-roadmap.md @@ -285,8 +285,19 @@ Diseño completo en [SDD 11 — Bootstrap from-scratch](11-bootstrap.md). Resume `SWAP_BWRAP=1`. **Pendiente:** correr in-VM acumulando los swaps (`KVM=1 MEM=24576 SWAP_MAKE=1 SWAP_BUSYBOX=1 SWAP_LINUX_HEADERS=1 SWAP_BWRAP=1 ./scripts/selfhost-verify.sh`) para el `✓ REPRODUCIBLE`, y seguir con la última pieza: rust/llvm (la grande — ya hay infra de deps). -- ⏭️ **También pendiente (Stage 1):** bus único (B.2: exponer el `CRASHED` a la capa de IA) y - atestación arje (A1/A2). +- 🔨 **B.2 — `CRASHED` real a la capa de IA (en progreso):** el `Event::Crashed` del bus de + agente ya no es sólo declarativo. **Fuente (lado arje, ✅):** `arje-bus` ganó + `BusRequest::Subscribe` + `BusPayload::Event(BusEvent)`; arje-zero difunde en `on_death` + `EnteCrashed{id,label,status}` / `EnteRestarting{delay_ms}` / `EnteExited` a las conexiones + suscritas, purgando las muertas (4 tests). **Sink (lado hammer, ✅):** `hammerd::crashes` + traduce la señal normalizada → `Event::Crashed` y la bombea al `EventBus` → + `/run/agent.sock` (3 tests). **Falta el último tramo:** el *adaptador de transporte* —un + thread que conecte a `$ARJE_BUS_SOCK`, mande `Subscribe` y reenvíe cada `BusEvent` a + `crashes::pump` (vía `BusClient::{subscribe,next_event}`)— y la decisión de cómo los dos + repos comparten el wire (dep directa a `arje-bus` vs. proto compartido vs. relectura del + frame postcard). Es el único paso que necesita un init vivo para validarse end-to-end. +- ⏭️ **También pendiente (Stage 1):** atestación arje (A1/A2 — ya cableada en el repo arje, + resta sólo enchufar su veredicto de boot a este roadmap). ## Notas de entorno - Desarrollo principal: laptop del autor.