feat(arje-link): transporte arje-bus → bus de agente, cierra B.2 end-to-end
hammerd::arje_link se suscribe al bus del init (ENTE_BUS_SOCK), relee el frame postcard de arje-bus con un mirror mínimo de suscriptor (sin arrastrar el crate-graph de arje ⇒ hammer sigue standalone) y reenvía cada BusEvent → crashes → Event::Crashed → agent.sock. Wire verificado byte-a-byte contra arje-bus real (ulid string, frame Subscribe=[00,01,00,0d]); 2 tests de round-trip local + frame. Se lanza en thread si ENTE_BUS_SOCK está definido (no-op si no). Roadmap B.2 marcado ✅ (resta sólo el smoke contra init vivo). 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
f6b337f6cb
commit
a23630c50d
@@ -25,6 +25,8 @@ nix = { version = "0.30", default-features = false, features = ["fanotify", "fs"
|
||||
libc = "0.2"
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
postcard.workspace = true
|
||||
ulid.workspace = true
|
||||
|
||||
[dev-dependencies]
|
||||
tempfile.workspace = true
|
||||
|
||||
@@ -0,0 +1,260 @@
|
||||
//! Adaptador de transporte arje-bus → bus de agente de hammer (último tramo de B.2).
|
||||
//!
|
||||
//! arje (el init/PID 1) supervisa cada Ente y, al morir, difunde `BusEvent` por su bus
|
||||
//! (`arje-bus`: `BusRequest::Subscribe` + `BusPayload::Event`). Este módulo se **suscribe** a
|
||||
//! ese stream y traduce cada evento al vocabulario de hammer (`crashes::Lifecycle` →
|
||||
//! `Event::Crashed`), publicándolo en el [`EventBus`] → `/run/agent.sock` → la capa de IA.
|
||||
//!
|
||||
//! ## Por qué un mirror y no una dependencia de `arje-bus`
|
||||
//!
|
||||
//! `arje-bus` arrastra el grafo de crates de arje (arje-card → card-core → …). Acoplarlo aquí
|
||||
//! rompería el build hermético/standalone de hammer. En su lugar **releemos el frame postcard**:
|
||||
//! el bus usa frames `u32-BE-len + postcard(BusMessage)`, y reproducimos sólo el subconjunto que
|
||||
//! un suscriptor necesita. El layout está verificado byte-a-byte contra `arje-bus` real (mismas
|
||||
//! versiones de `ulid` 1.2 y `postcard` 1.1; `ulid` serializa como string). Si arje reordena las
|
||||
//! variantes del enum, el contrato se rompe — por eso los discriminantes están documentados y
|
||||
//! cubiertos por un test de round-trip local.
|
||||
//!
|
||||
//! Varios campos del mirror (`from`, `seq`, `id`) no se leen: existen para **consumir** los
|
||||
//! bytes del wire al deserializar. Por eso el `allow(dead_code)` a nivel de módulo.
|
||||
#![allow(dead_code)]
|
||||
|
||||
use std::io::{Read, Write};
|
||||
use std::os::unix::net::UnixStream;
|
||||
use std::path::PathBuf;
|
||||
use std::thread;
|
||||
|
||||
use serde::Deserialize;
|
||||
use ulid::Ulid;
|
||||
|
||||
use crate::crashes::{self, LifeStatus, Lifecycle};
|
||||
use crate::events::EventBus;
|
||||
|
||||
/// Igual que `arje_bus::ENV_BUS_SOCK`: el env donde arje-zero publica la ruta del socket.
|
||||
pub const ENV_BUS_SOCK: &str = "ENTE_BUS_SOCK";
|
||||
|
||||
/// Tope de frame, igual que `arje_bus::MAX_FRAME` (1 MiB) — protección contra OOM.
|
||||
const MAX_FRAME: usize = 1 << 20;
|
||||
|
||||
// ── Mirror del wire de arje-bus (subconjunto de suscriptor) ──────────────────────────────
|
||||
//
|
||||
// Discriminantes (postcard = varint del índice de variante), verificados contra arje-bus:
|
||||
// BusPayload: Request=0, Response=1, Event=2
|
||||
// BusRequest::Subscribe = 13 (último; ver arje-bus/src/lib.rs)
|
||||
// BusEvent: EnteCrashed=0, EnteRestarting=1, EnteExited=2
|
||||
// LifecycleStatus: Exited=0, Killed=1
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
struct BusMessage {
|
||||
from: Option<Ulid>,
|
||||
seq: u64,
|
||||
payload: BusPayload,
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
enum BusPayload {
|
||||
/// Índice 0. Nunca lo recibe un suscriptor (los Invoke van a proveedores, no a subs);
|
||||
/// existe sólo para alinear el discriminante de `Response`/`Event`.
|
||||
Request,
|
||||
/// Índice 1. El único que recibimos por aquí es el `Ok` del ack de `Subscribe`.
|
||||
Response(MiniResponse),
|
||||
/// Índice 2. El stream de eventos de ciclo de vida.
|
||||
Event(BusEvent),
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
enum MiniResponse {
|
||||
Ok,
|
||||
Error(String),
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
enum BusEvent {
|
||||
EnteCrashed { id: Ulid, label: String, status: LifecycleStatus },
|
||||
EnteRestarting { id: Ulid, label: String, delay_ms: u64 },
|
||||
EnteExited { id: Ulid, label: String },
|
||||
}
|
||||
|
||||
#[derive(Deserialize, Debug)]
|
||||
enum LifecycleStatus {
|
||||
Exited(i32),
|
||||
Killed(i32),
|
||||
}
|
||||
|
||||
impl BusEvent {
|
||||
/// Traduce el evento de arje a la señal normalizada de hammer. El `label` del Ente es el
|
||||
/// `service`; una muerte por señal usa la convención de shell `128 + signum` para que el
|
||||
/// código sea siempre ≠ 0.
|
||||
fn to_lifecycle(self) -> Lifecycle {
|
||||
match self {
|
||||
BusEvent::EnteCrashed { label, status, .. } => {
|
||||
let code = match status {
|
||||
LifecycleStatus::Exited(c) => c,
|
||||
LifecycleStatus::Killed(sig) => 128 + sig,
|
||||
};
|
||||
Lifecycle { service: label, status: LifeStatus::Crashed { code } }
|
||||
}
|
||||
BusEvent::EnteRestarting { label, delay_ms, .. } => {
|
||||
Lifecycle { service: label, status: LifeStatus::Restarting { delay_ms } }
|
||||
}
|
||||
BusEvent::EnteExited { label, .. } => {
|
||||
Lifecycle { service: label, status: LifeStatus::Exited }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Cuerpo del frame `Subscribe` (sin el prefijo de longitud), hand-encodeado: `from=None` +
|
||||
/// `seq` (varint) + `BusPayload::Request`(0) + `BusRequest::Subscribe`(13). Verificado
|
||||
/// byte-a-byte contra `postcard::to_stdvec` de arje-bus.
|
||||
fn subscribe_body(seq: u64) -> Vec<u8> {
|
||||
let mut v = vec![0x00u8]; // Option::None
|
||||
let mut s = seq; // u64 como varint LEB128
|
||||
loop {
|
||||
let b = (s & 0x7f) as u8;
|
||||
s >>= 7;
|
||||
if s == 0 {
|
||||
v.push(b);
|
||||
break;
|
||||
} else {
|
||||
v.push(b | 0x80);
|
||||
}
|
||||
}
|
||||
v.push(0x00); // BusPayload::Request
|
||||
v.push(0x0D); // BusRequest::Subscribe (índice 13)
|
||||
v
|
||||
}
|
||||
|
||||
fn write_frame(w: &mut impl Write, body: &[u8]) -> std::io::Result<()> {
|
||||
w.write_all(&(body.len() as u32).to_be_bytes())?;
|
||||
w.write_all(body)?;
|
||||
w.flush()
|
||||
}
|
||||
|
||||
fn read_frame(r: &mut impl Read) -> std::io::Result<Vec<u8>> {
|
||||
let mut len = [0u8; 4];
|
||||
r.read_exact(&mut len)?;
|
||||
let n = u32::from_be_bytes(len) as usize;
|
||||
if n > MAX_FRAME {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidData,
|
||||
format!("frame oversize: {n} > {MAX_FRAME}"),
|
||||
));
|
||||
}
|
||||
let mut buf = vec![0u8; n];
|
||||
r.read_exact(&mut buf)?;
|
||||
Ok(buf)
|
||||
}
|
||||
|
||||
/// Conecta al bus de arje, se suscribe y bombea eventos de ciclo de vida al `EventBus` hasta
|
||||
/// que la conexión se cierra. Bloqueante — pensado para correr en su propio thread.
|
||||
fn run(sock: &PathBuf, bus: &EventBus) -> std::io::Result<()> {
|
||||
let mut stream = UnixStream::connect(sock)?;
|
||||
write_frame(&mut stream, &subscribe_body(1))?;
|
||||
tracing::info!(sock = %sock.display(), "arje-link: suscrito al bus de init");
|
||||
loop {
|
||||
let body = read_frame(&mut stream)?;
|
||||
let msg: BusMessage = match postcard::from_bytes(&body) {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
// Un frame que no entendemos no debe matar el puente: lo saltamos.
|
||||
tracing::warn!(error = %e, "arje-link: frame indescifrable, ignoro");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
match msg.payload {
|
||||
BusPayload::Event(ev) => {
|
||||
let life = ev.to_lifecycle();
|
||||
if let Some(event) = crashes::to_event(&life) {
|
||||
bus.publish(&event);
|
||||
}
|
||||
}
|
||||
// El ack del Subscribe y cualquier otra respuesta: nada que hacer.
|
||||
BusPayload::Response(_) | BusPayload::Request => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Lanza el puente en un thread si `ENTE_BUS_SOCK` está definido (arje corriendo). Si no, es un
|
||||
/// no-op: hammerd sigue sin la fuente de crashes (los tests del bus no la necesitan). El thread
|
||||
/// es resiliente: si el bus muere, loguea y termina sin tumbar al daemon.
|
||||
pub fn spawn_if_configured(bus: EventBus) {
|
||||
let Ok(path) = std::env::var(ENV_BUS_SOCK) else {
|
||||
tracing::debug!("arje-link: {ENV_BUS_SOCK} no definido; sin fuente de crashes");
|
||||
return;
|
||||
};
|
||||
let sock = PathBuf::from(path);
|
||||
let _ = thread::Builder::new().name("arje-link".into()).spawn(move || {
|
||||
if let Err(e) = run(&sock, &bus) {
|
||||
tracing::warn!(error = %e, "arje-link: puente terminó (¿bus de init caído?)");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn subscribe_body_es_el_frame_esperado() {
|
||||
// Verificado contra arje-bus real: [from=None, seq=1, Request, Subscribe=13].
|
||||
assert_eq!(subscribe_body(1), vec![0x00, 0x01, 0x00, 0x0D]);
|
||||
// seq grande → varint multibyte (300 = 0xAC 0x02).
|
||||
assert_eq!(subscribe_body(300), vec![0x00, 0xAC, 0x02, 0x00, 0x0D]);
|
||||
}
|
||||
|
||||
/// Round-trip local con el MISMO `postcard`/`ulid` que arje-bus: serializamos un
|
||||
/// `BusMessage::Event` con tipos espejo `Serialize` y lo deserializamos con los `Deserialize`
|
||||
/// reales del módulo. Si los discriminantes/orden de campos se desalinean, esto rompe.
|
||||
#[test]
|
||||
fn evento_crashed_roundtrips_por_el_wire() {
|
||||
use serde::Serialize;
|
||||
#[derive(Serialize)]
|
||||
enum TxPayload {
|
||||
#[allow(dead_code)]
|
||||
Request,
|
||||
#[allow(dead_code)]
|
||||
Response,
|
||||
Event(TxEvent),
|
||||
}
|
||||
#[derive(Serialize)]
|
||||
enum TxEvent {
|
||||
EnteCrashed { id: Ulid, label: String, status: TxStatus },
|
||||
}
|
||||
#[derive(Serialize)]
|
||||
enum TxStatus {
|
||||
#[allow(dead_code)]
|
||||
Exited(i32),
|
||||
Killed(i32),
|
||||
}
|
||||
#[derive(Serialize)]
|
||||
struct TxMsg {
|
||||
from: Option<Ulid>,
|
||||
seq: u64,
|
||||
payload: TxPayload,
|
||||
}
|
||||
|
||||
let id = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").unwrap();
|
||||
let bytes = postcard::to_stdvec(&TxMsg {
|
||||
from: Some(id),
|
||||
seq: 7,
|
||||
payload: TxPayload::Event(TxEvent::EnteCrashed {
|
||||
id,
|
||||
label: "demonio".into(),
|
||||
status: TxStatus::Killed(11),
|
||||
}),
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
let msg: BusMessage = postcard::from_bytes(&bytes).unwrap();
|
||||
match msg.payload {
|
||||
BusPayload::Event(ev) => {
|
||||
let life = ev.to_lifecycle();
|
||||
assert_eq!(life.service, "demonio");
|
||||
// Killed(11) → 128 + 11 = 139.
|
||||
assert_eq!(life.status, LifeStatus::Crashed { code: 139 });
|
||||
}
|
||||
other => panic!("esperaba Event, fue {other:?}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -12,6 +12,7 @@ use std::thread;
|
||||
|
||||
use clap::Parser;
|
||||
|
||||
mod arje_link;
|
||||
mod bus;
|
||||
mod control;
|
||||
mod crashes;
|
||||
@@ -70,6 +71,10 @@ fn main() -> anyhow::Result<()> {
|
||||
|
||||
let event_bus = events::EventBus::new();
|
||||
|
||||
// Puente con el bus del init (arje): si ENTE_BUS_SOCK está definido, se suscribe a los
|
||||
// eventos de ciclo de vida y publica los crashes en el bus de agente. No-op si no hay init.
|
||||
arje_link::spawn_if_configured(event_bus.clone());
|
||||
|
||||
// FIFO de control humano: lo creamos siempre que se pueda (no es fatal).
|
||||
let init_fifo_ok = match control::ensure_fifo(&args.init_control) {
|
||||
Ok(path) => {
|
||||
|
||||
Reference in New Issue
Block a user