feat(bus): B.2 sink del CRASHED — hammerd::crashes traduce eventos de arje
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>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
e9f535a1bb
commit
7c099afe4f
@@ -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.
|
||||
|
||||
@@ -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<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");
|
||||
}
|
||||
}
|
||||
@@ -14,6 +14,7 @@ use clap::Parser;
|
||||
|
||||
mod bus;
|
||||
mod control;
|
||||
mod crashes;
|
||||
mod events;
|
||||
mod watcher;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user