Nuevo módulo hammerd::crashes: traduce la señal de ciclo de vida normalizada (fuente = arje-bus BusEvent) a proto::Event::Crashed y la bombea al EventBus → agent.sock (capa de IA). Traducción + bucle de bombeo desacoplados del transporte y testeados (3 tests). Resta sólo el adaptador de transporte arje-bus ↔ hammerd. Roadmap B.2 actualizado. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
132 lines
5.6 KiB
Rust
132 lines
5.6 KiB
Rust
//! 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<Lifecycle>`] 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<Event> {
|
|
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<Lifecycle>, 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");
|
|
}
|
|
}
|