//! 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, 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 { 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, 0x0E], "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(); }