test(hammerd): e2e del adaptador arje_link contra un bus falso de arje

Cierra el último tramo de validación de B.2 (adoptar arje como init): el
adaptador arje_link ya estaba escrito y cableado en main, pero su único test
era un round-trip del espejo contra sí mismo — no probaba la ruta de transporte
real ni el contrato con arje-bus.

arje_link_e2e levanta un socket falso de arje-bus (mismo wire: frames
u32-BE-len + postcard), valida que el adaptador emite el frame Subscribe exacto
([00 01 00 0D]), le inyecta un EnteCrashed(Killed 11), y verifica que sale como
Event::Crashed{code:139} por el EventBus — ejercitando run() entero
(connect→subscribe→leer→traducir→publicar) sin necesitar un init vivo. run()
pasa a pub(crate) para que el test lo maneje. El contrato de bytes con arje-bus
queda fijado del lado arje por arje-bus::contrato_wire.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Sergio
2026-06-25 16:56:25 +00:00
co-authored by Claude Opus 4.8
parent 813d965d5c
commit 2207f18587
2 changed files with 128 additions and 1 deletions
+2 -1
View File
@@ -148,7 +148,8 @@ fn read_frame(r: &mut impl Read) -> std::io::Result<Vec<u8>> {
/// 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<()> {
/// `pub(crate)` para que el e2e (`tests/arje_link_e2e.rs`) lo maneje contra un bus falso.
pub(crate) 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");
+126
View File
@@ -0,0 +1,126 @@
//! E2E del adaptador de transporte arje-bus → bus de agente (último tramo de B.2).
//!
//! No necesita un init vivo: levantamos un **socket falso de arje-bus** que habla el mismo
//! wire (frames `u32-BE-len + postcard`), validamos que el adaptador manda el `Subscribe`
//! correcto, le inyectamos un `EnteCrashed`, y verificamos que sale como `Event::Crashed`
//! por el `EventBus`. Así se ejercita la ruta real de `arje_link::run`: connect → subscribe →
//! leer frame → traducir → publicar.
//!
//! El contrato de bytes con arje-bus (que acá replicamos como `Tx*`) está fijado del lado de
//! arje por `arje-bus::contrato_wire` — si arje reordena un enum, ese test falla primero.
// Importamos los módulos privados del binario (mismo patrón que `bus_e2e.rs`). El orden
// importa: `arje_link` y `crashes` referencian `crate::events`.
#[path = "../src/events.rs"]
mod events;
#[path = "../src/crashes.rs"]
mod crashes;
#[path = "../src/arje_link.rs"]
mod arje_link;
use std::io::{Read, Write};
use std::os::unix::net::UnixListener;
use std::path::PathBuf;
use std::thread;
use std::time::Duration;
use hammer_core::proto::Event;
use serde::Serialize;
use ulid::Ulid;
// ── Lado emisor del wire (lo que arje-bus serializaría). Discriminantes fijados por
// `arje-bus::contrato_wire`: BusPayload Event=2, BusEvent EnteCrashed=0, LifecycleStatus
// Killed=1. ─────────────────────────────────────────────────────────────────────────────
#[derive(Serialize)]
struct TxMsg {
from: Option<Ulid>,
seq: u64,
payload: TxPayload,
}
#[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),
}
fn frame(body: &[u8]) -> Vec<u8> {
let mut v = (body.len() as u32).to_be_bytes().to_vec();
v.extend_from_slice(body);
v
}
#[test]
fn crashed_de_un_bus_falso_llega_como_event_crashed() {
let tmp = tempfile::tempdir().unwrap();
let sock = tmp.path().join("ente-bus.sock");
let listener = UnixListener::bind(&sock).unwrap();
// Servidor falso de arje-bus.
let server = thread::spawn(move || {
let (mut conn, _) = listener.accept().unwrap();
// 1) El adaptador debe mandar el frame Subscribe exacto: [None, seq=1, Request, Subscribe=13].
let mut len = [0u8; 4];
conn.read_exact(&mut len).unwrap();
let n = u32::from_be_bytes(len) as usize;
let mut body = vec![0u8; n];
conn.read_exact(&mut body).unwrap();
assert_eq!(body, vec![0x00, 0x01, 0x00, 0x0D], "frame Subscribe inesperado");
// 2) Le inyectamos un EnteCrashed por señal (SIGSEGV=11 → 128+11=139).
let id = Ulid::from_string("01ARZ3NDEKTSV4RRFFQ69G5FAV").unwrap();
let ev = postcard::to_stdvec(&TxMsg {
from: Some(id),
seq: 0,
payload: TxPayload::Event(TxEvent::EnteCrashed {
id,
label: "demonio".into(),
status: TxStatus::Killed(11),
}),
})
.unwrap();
conn.write_all(&frame(&ev)).unwrap();
conn.flush().unwrap();
// Al salir del closure se cierra `conn` → el adaptador sale por EOF.
});
let bus = events::EventBus::new();
let rx = bus.subscribe();
let bus_for_run = bus.clone();
let sock_for_run = PathBuf::from(&sock);
let client = thread::spawn(move || {
// Devuelve Err por el EOF final del servidor; es el cierre esperado.
let _ = arje_link::run(&sock_for_run, &bus_for_run);
});
// El crash debe materializarse como Event::Crashed en el bus de agente.
let ev = rx
.recv_timeout(Duration::from_secs(5))
.expect("el adaptador no publicó el crash");
match ev {
Event::Crashed { service, code } => {
assert_eq!(service, "demonio");
assert_eq!(code, 139, "Killed(11) → 128+11");
}
other => panic!("esperaba Event::Crashed, fue {other:?}"),
}
server.join().unwrap();
client.join().unwrap();
}