//! E2E del bus de agente (Fase 5): //! //! Arrancamos `serve_agent_bus` en un thread con un socket bajo tempdir, conectamos como //! cliente, completamos el handshake y ejercitamos los caminos no-pesados (`query`, errores). //! //! `compile` real no se ejercita aquí — depende del lab + red, ya cubierto por el e2e de //! grep gated en `TAKANA_NETWORK_TESTS`. Lo que sí verificamos es que un `compile` SIN la //! capacidad devuelve `Error{code:"no_cap"}`, lo que prueba el dispatcher + el gate. // El binario no expone una API pública; importamos sus módulos privados desde // `path = ...`. `proto` se movió a takana-core en Fase 6. #[path = "../src/events.rs"] mod events; #[path = "../src/control.rs"] mod control; #[path = "../src/bus.rs"] mod bus; use std::io::{BufRead, BufReader, Write}; use std::os::unix::net::UnixStream; use std::path::PathBuf; use std::sync::Arc; use std::thread; use std::time::Duration; use takana_core::proto::{Cap, Command, Event, Peer, RecipeInline}; /// Helper: lanza un bus en background y devuelve la ruta del socket. Limpia al final del /// test gracias al `tempfile::TempDir` que se le pasa al caller. fn start_bus( sock: &std::path::Path, policy: bus::CapsPolicy, store_root: PathBuf, init_control: PathBuf, ) -> events::EventBus { let events = events::EventBus::new(); let ctx = bus::BusContext { store_root, init_control, events: events.clone(), }; let sock = sock.to_path_buf(); thread::Builder::new() .name("bus-test".into()) .spawn(move || { let _ = bus::serve_agent_bus(&sock, policy, ctx); }) .unwrap(); // El caller usa `wait_for_sock` para esperar a que `bind` exista en disco. events } fn wait_for_sock(sock: &std::path::Path) -> UnixStream { for _ in 0..200 { if let Ok(s) = UnixStream::connect(sock) { return s; } std::thread::sleep(Duration::from_millis(10)); } panic!("no pude conectar al bus en {}", sock.display()); } fn send(stream: &mut UnixStream, cmd: &Command) { let line = serde_json::to_string(cmd).unwrap(); stream.write_all(line.as_bytes()).unwrap(); stream.write_all(b"\n").unwrap(); stream.flush().unwrap(); } fn recv_event(reader: &mut BufReader) -> Event { let mut line = String::new(); reader.read_line(&mut line).expect("read event"); serde_json::from_str(line.trim()).expect("event JSON válido") } fn permissive_policy() -> bus::CapsPolicy { Arc::new(|_peer: &Peer| { vec![Cap::Query, Cap::Compile, Cap::Inject, Cap::Init] }) } fn read_only_policy() -> bus::CapsPolicy { Arc::new(|_peer: &Peer| vec![Cap::Query]) } #[test] fn handshake_returns_welcome_with_peer_creds() { let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); start_bus( &sock, permissive_policy(), d.path().join("store"), d.path().join("init.ctl"), ); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send( &mut writer, &Command::Hello { ver: 1, client: "test".into() }, ); let ev = recv_event(&mut reader); match ev { Event::Welcome { ver, caps, peer } => { assert_eq!(ver, 1); assert!(caps.contains(&Cap::Query)); assert_eq!(peer.uid, unsafe { libc::getuid() }); assert_eq!(peer.pid, std::process::id() as i32); } other => panic!("esperaba Welcome, llegó {other:?}"), } } #[test] fn missing_cap_returns_error_no_cap() { let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); start_bus( &sock, read_only_policy(), d.path().join("store"), d.path().join("init.ctl"), ); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send( &mut writer, &Command::Hello { ver: 1, client: "test".into() }, ); let _ = recv_event(&mut reader); // welcome send( &mut writer, &Command::Compile { recipe: RecipeInline { name: "x".into(), repo: Some("git://x".into()), commit: Some("abc".into()), tarball: None, sha256: None, patch: None, compiler: "zig-cc".into(), target: "x86_64-linux-musl".into(), link: "static".into(), flags: vec![], evidence: Default::default(), }, }, ); match recv_event(&mut reader) { Event::Error { code, msg } => { assert_eq!(code, "no_cap"); assert!(msg.contains("Compile"), "{msg}"); } other => panic!("esperaba Error{{no_cap}}, llegó {other:?}"), } } #[test] fn query_file_returns_value_for_existing_path() { let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); start_bus( &sock, permissive_policy(), d.path().join("store"), d.path().join("init.ctl"), ); let target = d.path().join("data.txt"); std::fs::write(&target, b"contenido").unwrap(); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send(&mut writer, &Command::Hello { ver: 1, client: "t".into() }); let _ = recv_event(&mut reader); // welcome send( &mut writer, &Command::Query { what: "file".into(), path: Some(target.display().to_string()), name: None, expr: None, }, ); match recv_event(&mut reader) { Event::QueryResult { what, value } => { assert_eq!(what, "file"); assert_eq!(value["exists"], serde_json::json!(true)); assert_eq!(value["size"], serde_json::json!(9)); } other => panic!("esperaba QueryResult, llegó {other:?}"), } } #[test] fn query_expr_evaluates_against_host_path() { let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); start_bus( &sock, permissive_policy(), d.path().join("store"), d.path().join("init.ctl"), ); // file: sobre un archivo real; el daemon evalúa con fs_root=None ⇒ ruta absoluta. let target = d.path().join("payload.bin"); std::fs::write(&target, b"X").unwrap(); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send(&mut writer, &Command::Hello { ver: 1, client: "t".into() }); let _ = recv_event(&mut reader); // welcome send( &mut writer, &Command::Query { what: "expr".into(), path: None, name: None, expr: Some(format!("file:{}", target.display())), }, ); match recv_event(&mut reader) { Event::QueryResult { what, value } => { assert_eq!(what, "expr"); assert_eq!(value["found"], serde_json::json!(true)); assert_eq!(value["size"], serde_json::json!(1)); } other => panic!("esperaba QueryResult, llegó {other:?}"), } // Y una expresión inválida vuelve como JSON con `error`, no como Event::Error // (la query es bien formada; el problema es el contenido). send( &mut writer, &Command::Query { what: "expr".into(), path: None, name: None, expr: Some("foo:bar".into()), }, ); match recv_event(&mut reader) { Event::QueryResult { what, value } => { assert_eq!(what, "expr"); assert!(value["error"].as_str().unwrap().contains("kind 'foo'")); } other => panic!("esperaba QueryResult, llegó {other:?}"), } } #[test] fn modified_events_fan_out_to_subscribers() { // Validamos sólo el camino del fan-out a través del bus, sin watcher real: // disparamos `publish(Modified)` directamente y vemos que llega por el socket. let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); let bus_handle = start_bus( &sock, permissive_policy(), d.path().join("store"), d.path().join("init.ctl"), ); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send(&mut writer, &Command::Hello { ver: 1, client: "t".into() }); let _ = recv_event(&mut reader); // welcome // El forwarder de eventos arranca dentro de handle_connection; damos un margen para // que la suscripción esté registrada antes de publicar. std::thread::sleep(Duration::from_millis(50)); bus_handle.publish(&Event::Modified { path: "/bin/grep".into(), op: "replace".into(), ts: "2026-06-09T00:00:00Z".into(), }); match recv_event(&mut reader) { Event::Modified { path, op, ts } => { assert_eq!(path, "/bin/grep"); assert_eq!(op, "replace"); assert_eq!(ts, "2026-06-09T00:00:00Z"); } other => panic!("esperaba Modified, llegó {other:?}"), } } #[test] fn init_command_writes_to_fifo() { let d = tempfile::tempdir().unwrap(); let sock = d.path().join("agent.sock"); let fifo = d.path().join("init.ctl"); control::ensure_fifo(&fifo).unwrap(); // Reader del FIFO en background: lee una línea y termina. let fifo_for_reader = fifo.clone(); let reader_thread = std::thread::spawn(move || { let f = std::fs::OpenOptions::new() .read(true) .open(&fifo_for_reader) .unwrap(); let mut r = std::io::BufReader::new(f); let mut s = String::new(); std::io::BufRead::read_line(&mut r, &mut s).unwrap(); s }); std::thread::sleep(Duration::from_millis(50)); start_bus( &sock, permissive_policy(), d.path().join("store"), fifo.clone(), ); let stream = wait_for_sock(&sock); let mut writer = stream.try_clone().unwrap(); let mut reader = BufReader::new(stream); send(&mut writer, &Command::Hello { ver: 1, client: "t".into() }); let _ = recv_event(&mut reader); send(&mut writer, &Command::Init { cmd: "start web".into() }); match recv_event(&mut reader) { Event::InitAck { cmd } => assert_eq!(cmd, "start web"), other => panic!("esperaba InitAck, llegó {other:?}"), } let from_fifo = reader_thread.join().unwrap(); assert_eq!(from_fifo, "start web\n"); }