La evidencia corre del lado del BUILD (hammer-agent no depende de hammer-build:
separación PROPONE/CONSTRUYE). Piezas:
- proto: RecipeInline lleva `evidence`; Event::BuildReady lleva
`verdict: Option<EvidenceVerdict>` (+ EvidenceCheckVerdict). Ambos con
skip_serializing_if ⇒ wire compatible con clientes pre-H1c.
- hammerd/bus: run_compile ejecuta la evidencia declarada tras sellar el
artefacto (run_evidence vía swm_bridge::recipe_from_source_patch) y adjunta
el veredicto en BuildReady. Si ni se pudo ejecutar ⇒ veredicto fallido
sintético (no pasa en silencio). El lab reporta; el gate vive en el agente.
- client: compile() devuelve CompileOutcome { artifact, verdict }.
- orchestrator: VERIFY lee el veredicto; si all_passed=false empuja
VerifyCheck::fail y ABORTA antes de hidratar (artefacto sellado, no toca el
sistema). Nuevo constructor VerifyCheck::fail.
Tests: roundtrip del veredicto en proto; stub-bus refleja evidencia→verdict;
nueva prueba de integración orchestrator_evidence_gate (veredicto fallido ⇒
run() aborta con "no se propone" y no hidrata). Suite completa verde (33 suites).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
314 lines
11 KiB
Rust
314 lines
11 KiB
Rust
//! `AgentClient`: cliente síncrono del bus de agente. Ver `docs/07-agent-bus.md` §3.
|
|
//!
|
|
//! Modelo: una conexión = un thread del que llama. El reader thread interno demultiplexa
|
|
//! eventos por tipo: respuestas correlativas a un comando se devuelven al caller via
|
|
//! `recv_*_event_blocking`; eventos asíncronos (`Modified`, `Crashed`) se acumulan en una
|
|
//! cola que el caller drena con `drain_async()`.
|
|
//!
|
|
//! No hay correlation-id en el protocolo todavía (el SDD 07 no lo exige), así que la
|
|
//! correlación se hace por **tipo**: tras un `Compile` esperamos `BuildReady|BuildFailed`,
|
|
//! cualquier otra cosa que no sea `Modified`/`Crashed` se reporta como protocol error.
|
|
|
|
use std::collections::VecDeque;
|
|
use std::io::{BufRead, BufReader, Write};
|
|
use std::os::unix::net::UnixStream;
|
|
use std::path::Path;
|
|
use std::sync::{Arc, Condvar, Mutex};
|
|
use std::thread;
|
|
use std::time::{Duration, Instant};
|
|
|
|
use hammer_core::proto::{Cap, Command, Event, Peer, RecipeInline};
|
|
|
|
#[derive(Debug, thiserror::Error)]
|
|
pub enum ClientError {
|
|
#[error("io: {0}")]
|
|
Io(#[from] std::io::Error),
|
|
#[error("json: {0}")]
|
|
Json(#[from] serde_json::Error),
|
|
#[error("protocolo: {0}")]
|
|
Protocol(String),
|
|
#[error("timeout esperando {0}")]
|
|
Timeout(&'static str),
|
|
#[error("el bus devolvió error code={code} msg={msg}")]
|
|
Bus { code: String, msg: String },
|
|
#[error("la conexión se cerró")]
|
|
Closed,
|
|
}
|
|
|
|
pub type Result<T> = std::result::Result<T, ClientError>;
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct Welcome {
|
|
pub ver: u32,
|
|
pub caps: Vec<Cap>,
|
|
pub peer: Peer,
|
|
}
|
|
|
|
/// Resultado de `compile`: el hash del artefacto sellado más el veredicto de la evidencia (H1c).
|
|
/// `verdict` es `None` cuando la receta no declaró evidencia.
|
|
#[derive(Debug, Clone)]
|
|
pub struct CompileOutcome {
|
|
pub artifact: String,
|
|
pub verdict: Option<hammer_core::proto::EvidenceVerdict>,
|
|
}
|
|
|
|
/// Estado compartido entre el reader thread y el caller.
|
|
struct Inbox {
|
|
/// Eventos asíncronos (Modified/Crashed) en orden de llegada.
|
|
async_events: VecDeque<Event>,
|
|
/// Próximo evento "respuesta" al último comando. Sólo cabe uno a la vez en este
|
|
/// cliente síncrono (el caller envía + espera antes del siguiente).
|
|
pending: Option<Event>,
|
|
closed: bool,
|
|
}
|
|
|
|
#[derive(Clone)]
|
|
struct Shared {
|
|
inbox: Arc<Mutex<Inbox>>,
|
|
cv: Arc<Condvar>,
|
|
}
|
|
|
|
pub struct AgentClient {
|
|
writer: UnixStream,
|
|
shared: Shared,
|
|
pub welcome: Welcome,
|
|
}
|
|
|
|
impl AgentClient {
|
|
/// Conecta al socket, hace handshake `Hello` y arranca el reader thread. Devuelve un
|
|
/// cliente listo para enviar comandos.
|
|
pub fn connect(sock_path: &Path) -> Result<Self> {
|
|
Self::connect_with_client_name(sock_path, "hammer-agent")
|
|
}
|
|
|
|
pub fn connect_with_client_name(sock_path: &Path, client: &str) -> Result<Self> {
|
|
let stream = UnixStream::connect(sock_path)?;
|
|
let reader_stream = stream.try_clone()?;
|
|
let writer = stream;
|
|
|
|
let shared = Shared {
|
|
inbox: Arc::new(Mutex::new(Inbox {
|
|
async_events: VecDeque::new(),
|
|
pending: None,
|
|
closed: false,
|
|
})),
|
|
cv: Arc::new(Condvar::new()),
|
|
};
|
|
// Reader thread: drena el socket línea a línea hasta EOF.
|
|
let s_for_reader = shared.clone();
|
|
thread::Builder::new()
|
|
.name("agent-client-reader".into())
|
|
.spawn(move || reader_loop(reader_stream, s_for_reader))
|
|
.ok();
|
|
|
|
let mut me = Self {
|
|
writer,
|
|
shared,
|
|
welcome: Welcome {
|
|
ver: 0,
|
|
caps: vec![],
|
|
peer: Peer { uid: 0, gid: 0, pid: 0 },
|
|
},
|
|
};
|
|
me.send(&Command::Hello {
|
|
ver: hammer_core::proto::PROTOCOL_VERSION,
|
|
client: client.to_string(),
|
|
})?;
|
|
let ev = me.recv_pending(Duration::from_secs(5))?;
|
|
match ev {
|
|
Event::Welcome { ver, caps, peer } => {
|
|
me.welcome = Welcome { ver, caps, peer };
|
|
}
|
|
Event::Error { code, msg } => return Err(ClientError::Bus { code, msg }),
|
|
other => {
|
|
return Err(ClientError::Protocol(format!(
|
|
"esperaba Welcome, llegó {other:?}"
|
|
)));
|
|
}
|
|
}
|
|
Ok(me)
|
|
}
|
|
|
|
/// Bloquea hasta `BuildReady` o `BuildFailed`. `timeout` es el wall-clock máximo —
|
|
/// el lab puede tardar varios minutos, así que el default debería ser generoso.
|
|
pub fn compile(&mut self, recipe: RecipeInline, timeout: Duration) -> Result<CompileOutcome> {
|
|
self.send(&Command::Compile { recipe })?;
|
|
match self.recv_pending(timeout)? {
|
|
Event::BuildReady { artifact, verdict, .. } => Ok(CompileOutcome { artifact, verdict }),
|
|
Event::BuildFailed { reason, .. } => Err(ClientError::Bus {
|
|
code: "build_failed".into(),
|
|
msg: reason,
|
|
}),
|
|
Event::Error { code, msg } => Err(ClientError::Bus { code, msg }),
|
|
other => Err(ClientError::Protocol(format!(
|
|
"esperaba BuildReady|BuildFailed, llegó {other:?}"
|
|
))),
|
|
}
|
|
}
|
|
|
|
pub fn inject(
|
|
&mut self,
|
|
artifact: &str,
|
|
target: &str,
|
|
overlay: Option<&str>,
|
|
timeout: Duration,
|
|
) -> Result<usize> {
|
|
self.send(&Command::Inject {
|
|
artifact: artifact.to_string(),
|
|
target: target.to_string(),
|
|
overlay: overlay.map(String::from),
|
|
})?;
|
|
match self.recv_pending(timeout)? {
|
|
Event::Injected { files, .. } => Ok(files),
|
|
Event::Error { code, msg } => Err(ClientError::Bus { code, msg }),
|
|
other => Err(ClientError::Protocol(format!("esperaba Injected, llegó {other:?}"))),
|
|
}
|
|
}
|
|
|
|
pub fn query_file(&mut self, path: &str, timeout: Duration) -> Result<serde_json::Value> {
|
|
self.send(&Command::Query {
|
|
what: "file".into(),
|
|
path: Some(path.to_string()),
|
|
name: None,
|
|
expr: None,
|
|
})?;
|
|
match self.recv_pending(timeout)? {
|
|
Event::QueryResult { value, .. } => Ok(value),
|
|
Event::Error { code, msg } => Err(ClientError::Bus { code, msg }),
|
|
other => Err(ClientError::Protocol(format!(
|
|
"esperaba QueryResult, llegó {other:?}"
|
|
))),
|
|
}
|
|
}
|
|
|
|
/// Evalúa una expresión del mini-lenguaje (`bin:grep`, `pin:musl`, …) en el daemon.
|
|
/// El JSON devuelto sigue el schema de `hammer_core::query::eval`.
|
|
pub fn query_expr(&mut self, expr: &str, timeout: Duration) -> Result<serde_json::Value> {
|
|
self.send(&Command::Query {
|
|
what: "expr".into(),
|
|
path: None,
|
|
name: None,
|
|
expr: Some(expr.to_string()),
|
|
})?;
|
|
match self.recv_pending(timeout)? {
|
|
Event::QueryResult { value, .. } => Ok(value),
|
|
Event::Error { code, msg } => Err(ClientError::Bus { code, msg }),
|
|
other => Err(ClientError::Protocol(format!(
|
|
"esperaba QueryResult, llegó {other:?}"
|
|
))),
|
|
}
|
|
}
|
|
|
|
pub fn init(&mut self, cmd: &str, timeout: Duration) -> Result<()> {
|
|
self.send(&Command::Init { cmd: cmd.to_string() })?;
|
|
match self.recv_pending(timeout)? {
|
|
Event::InitAck { .. } => Ok(()),
|
|
Event::Error { code, msg } => Err(ClientError::Bus { code, msg }),
|
|
other => Err(ClientError::Protocol(format!("esperaba InitAck, llegó {other:?}"))),
|
|
}
|
|
}
|
|
|
|
/// Vacía la cola de eventos asíncronos. No bloquea.
|
|
pub fn drain_async(&self) -> Vec<Event> {
|
|
let Ok(mut g) = self.shared.inbox.lock() else { return vec![] };
|
|
std::mem::take(&mut g.async_events).into_iter().collect()
|
|
}
|
|
|
|
/// Espera el próximo evento asíncrono (típicamente para verificar que un `Modified`
|
|
/// específico haya llegado tras una operación).
|
|
pub fn next_async(&self, timeout: Duration) -> Result<Event> {
|
|
let g = self.shared.inbox.lock().unwrap();
|
|
let (mut g, _) = self
|
|
.shared
|
|
.cv
|
|
.wait_timeout_while(g, timeout, |i| {
|
|
!i.closed && i.async_events.is_empty()
|
|
})
|
|
.unwrap();
|
|
if let Some(ev) = g.async_events.pop_front() {
|
|
Ok(ev)
|
|
} else if g.closed {
|
|
Err(ClientError::Closed)
|
|
} else {
|
|
Err(ClientError::Timeout("async event"))
|
|
}
|
|
}
|
|
|
|
fn send(&mut self, cmd: &Command) -> Result<()> {
|
|
let line = serde_json::to_string(cmd)?;
|
|
self.writer.write_all(line.as_bytes())?;
|
|
self.writer.write_all(b"\n")?;
|
|
self.writer.flush()?;
|
|
Ok(())
|
|
}
|
|
|
|
fn recv_pending(&self, timeout: Duration) -> Result<Event> {
|
|
let deadline = Instant::now() + timeout;
|
|
let g = self.shared.inbox.lock().unwrap();
|
|
let (mut g, _) = self
|
|
.shared
|
|
.cv
|
|
.wait_timeout_while(g, timeout, |i| !i.closed && i.pending.is_none())
|
|
.unwrap();
|
|
if let Some(ev) = g.pending.take() {
|
|
return Ok(ev);
|
|
}
|
|
if g.closed {
|
|
return Err(ClientError::Closed);
|
|
}
|
|
// Si el wait expiró exactamente en el momento de un evento, vuelvo a comprobar.
|
|
if Instant::now() >= deadline {
|
|
return Err(ClientError::Timeout("response"));
|
|
}
|
|
Err(ClientError::Timeout("response"))
|
|
}
|
|
}
|
|
|
|
fn reader_loop(stream: UnixStream, shared: Shared) {
|
|
let reader = BufReader::new(stream);
|
|
for line in reader.lines() {
|
|
let line = match line {
|
|
Ok(l) => l,
|
|
Err(e) => {
|
|
tracing::warn!(error = %e, "agent-client: read error, cerrando");
|
|
break;
|
|
}
|
|
};
|
|
if line.trim().is_empty() {
|
|
continue;
|
|
}
|
|
let ev: Event = match serde_json::from_str(&line) {
|
|
Ok(e) => e,
|
|
Err(e) => {
|
|
tracing::warn!(error = %e, line = %line, "agent-client: línea no parseable");
|
|
continue;
|
|
}
|
|
};
|
|
let is_async = matches!(
|
|
&ev,
|
|
Event::Modified { .. } | Event::Crashed { .. }
|
|
);
|
|
let mut g = match shared.inbox.lock() {
|
|
Ok(g) => g,
|
|
Err(_) => break,
|
|
};
|
|
if is_async {
|
|
g.async_events.push_back(ev);
|
|
} else {
|
|
// Si ya hay un pending sin recoger, eso es un bug del caller (envió dos
|
|
// comandos antes de leer la primera respuesta). Lo logueamos y reemplazamos.
|
|
if g.pending.is_some() {
|
|
tracing::warn!(
|
|
"agent-client: pending ya tenía respuesta; el caller no la recogió"
|
|
);
|
|
}
|
|
g.pending = Some(ev);
|
|
}
|
|
shared.cv.notify_all();
|
|
}
|
|
if let Ok(mut g) = shared.inbox.lock() {
|
|
g.closed = true;
|
|
shared.cv.notify_all();
|
|
}
|
|
}
|