diff --git a/hive-c0re/src/lifecycle/mod.rs b/hive-c0re/src/lifecycle/mod.rs index 11ab3993..7203e9a5 100644 --- a/hive-c0re/src/lifecycle/mod.rs +++ b/hive-c0re/src/lifecycle/mod.rs @@ -935,48 +935,9 @@ async fn priv_run_inner(kind: &str, name: &str, node_id: Option) -> Result< match result { Ok(()) => Ok(()), - Err(e) => { - let journal = if kind == "update" { - container_journal_tail(&container).await - } else { - String::new() - }; - match log_id { - Some(id) => bail!("{e:#}; see build log #{id}{journal}"), - None => bail!("{e:#}{journal}"), - } - } - } -} - -/// On a failed `nixos-container update`, the stderr nixos-container -/// itself prints is often terse ("failed to reload container") — the -/// real reason (which unit failed `switch-to-configuration` during -/// the reload phase) lands in the *container's* own journal, not on -/// the host. Fetch the tail of it so a failed rebuild self-documents -/// the failing unit in the error string, no second round-trip. -/// -/// Scoped to `update`: that's the reload-phase case, and the -/// container is still up (running the old generation) so -/// `journalctl -M` works. Best-effort — returns "" for other verbs -/// or when the journal can't be read (machine gone, journalctl -/// missing); it never produces an error of its own. -async fn container_journal_tail(container: &str) -> String { - // `-M` enters the container namespace and needs root, so the read - // is delegated to hive-priv (hive-c0re itself runs unprivileged). - let res = crate::priv_client::read_container_journal( - container, - hive_priv_sock::JournalQuery { - lines: 40, - ..Default::default() + Err(e) => match log_id { + Some(id) => bail!("{e:#}; see build log #{id}"), + None => bail!("{e:#}"), }, - ) - .await; - match res { - Ok((stdout, _)) if !stdout.is_empty() => format!( - "\n--- last 40 journal lines from container '{container}' ---\n{}", - stdout.trim_end() - ), - _ => String::new(), } } diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index 160fcf58..c9ecaa1b 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -344,17 +344,10 @@ async fn cmd_serve( } } }); - // Knowledge periodic pull: hourly fallback in case the webhook is - // missed (e.g. hive-c0re was down during a push). Deliberately does - // NOT also fire an immediate pull at startup the way this task used - // to: `auto_update::run`'s `NodeKind::KnowledgePull` DAG node (spawned - // separately, a few lines up) already does that unconditionally on - // every boot. The two used to run concurrently with no lock between - // them, both `git pull --ff-only`-ing the same working tree — a real - // race, and the likely root cause of the "local changes would be - // overwritten" wedge this file's `pull()` now defends against - // (`reset --hard` before every pull). Removing the redundant caller - // fixes the race at its source instead of just self-healing after it. + // Knowledge periodic pull: hourly fallback in case a knowledge event from + // the swarm is missed (e.g. hive-c0re was down when it was sent). Sleeps + // before its first pull: the boot `NodeKind::KnowledgePull` DAG node + // already pulls once at startup. let mut knowledge_shutdown = coord.shutdown_rx(); let knowledge_coord = coord.clone(); tokio::spawn(async move { @@ -394,15 +387,11 @@ async fn cmd_serve( } } }); - // Matrix user sweep: same shape — ensure every container has - // an account on the local matrix-tuwunel homeserver with an - // access_token persisted to `/matrix-token`. No-op when - // the hive-matrix container isn't running. - // - // Re-submitted every 30 minutes so that token files deleted by - // `hive-matrix-daemon` (stale-token recovery — `M_UNKNOWN_TOKEN`) - // get re-provisioned without requiring a hive-c0re restart. The - // startup pass is the boot `MatrixSweep` node. + // Matrix sweep (`matrix::ensure_all`): the hive's own account, the hive + // Space and chat room, and every agent container's invite to both. It + // creates no agent accounts or tokens — swarm-controller mints those. + // No-op when the hive-matrix container isn't running. Re-submitted every + // 30 minutes; the startup pass is the boot `MatrixSweep` node. let mut matrix_shutdown = coord.shutdown_rx(); let matrix_coord = coord.clone(); tokio::spawn(async move { diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 570dfb83..c36cfbb4 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -287,8 +287,8 @@ async fn dispatch(req: &HostRequest, coord: Arc) -> HostResponse { } } -/// Create + start the container for `name`, rolling back socket -/// registration and notifying the manager on failure. +/// Create + start the container for `name` and bind its MCP listener. On a +/// failed spawn nothing was registered: post a swarm notice and return the error. async fn handle_spawn(coord: &Arc, name: &str) -> Result { tracing::info!(%name, "spawn"); let agent_dir = crate::paths::agent_runtime_dir(name); @@ -762,7 +762,6 @@ async fn handle_matrix_reset_password(name: &str) -> Result { Ok(HostResponse::messages(vec![ format!("matrix: password for @{name}:{server_name} reset"), format!("password persisted at: {}", pw_path.display()), - format!("next: hivectl matrix create-user {name} # mints a fresh access token"), ])) } diff --git a/hive-c0re/src/workers/knowledge.rs b/hive-c0re/src/workers/knowledge.rs index a2910053..bd6cb705 100644 --- a/hive-c0re/src/workers/knowledge.rs +++ b/hive-c0re/src/workers/knowledge.rs @@ -111,9 +111,8 @@ pub async fn pull(coord: &Coordinator) -> Result<()> { // `LOCAL_DIR` after the initial clone — this working tree exists to mirror // `origin/main`, not to be edited in place. A tracked file left dirty // by any other means (a stray manual edit on the host, an interrupted - // prior operation, or — the actual root cause here — two unsynchronized - // boot-time pull callers racing on this same working tree, since fixed - // in `main.rs`) would otherwise abort the `--ff-only` merge below with + // prior operation, or two pulls running on this working tree at once) + // would otherwise abort the `--ff-only` merge below with // "local changes would be overwritten"; an untracked file left behind // the same ways aborts it with "untracked working tree files would be // overwritten" instead — same wedge, just `clean`'s failure message diff --git a/hive-c0re/src/workers/mcp_sockets.rs b/hive-c0re/src/workers/mcp_sockets.rs index 49d9b860..5e1c040f 100644 --- a/hive-c0re/src/workers/mcp_sockets.rs +++ b/hive-c0re/src/workers/mcp_sockets.rs @@ -6,8 +6,8 @@ //! re-register all running agents. //! //! After startup, listeners are managed event-driven: -//! - `run_create` calls `register_agent` eagerly on first-spawn. -//! - `run_reconcile` calls `register_agent` immediately after `start_with_fallback`. +//! - `run_start` (the job queue's `Start` node) and `server::handle_spawn` +//! call `register_agent` once the container is started. //! - `kill`/`destroy` paths call `unregister_agent`. //! //! No recurring poll is needed because c0re owns the listener lifecycle. An diff --git a/hive-screen-mcp/src/main.rs b/hive-screen-mcp/src/main.rs index 47b83501..89a737ed 100644 --- a/hive-screen-mcp/src/main.rs +++ b/hive-screen-mcp/src/main.rs @@ -29,16 +29,37 @@ use rmcp::{ transport::stdio, }; use serde::Deserialize; +use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; use tokio::process::Command; +/// Upper bound on one `grim` / `wtype` run. Against a live local compositor +/// these finish in well under a second, a long `type_text` included. The +/// agent's turn waits on the tool call, so a wedged compositor would +/// otherwise hold the turn for as long as it stays wedged. +const CMD_TIMEOUT: Duration = Duration::from_secs(30); + +/// Upper bound on one RFB exchange (connect, handshake, pointer events) with +/// the neatvnc server on `127.0.0.1`: a few round trips of a few bytes each. +/// A server that is down refuses the connect at once; this catches one that +/// accepts and then never speaks. +const VNC_TIMEOUT: Duration = Duration::from_secs(10); + /// Run a command and return its result. /// `Ok(stdout_trimmed)` on success (may be empty string when the tool /// produces no output). `Err(human_readable_message)` on failure -/// (non-zero exit or failed spawn). -async fn run_cmd(program: &str, args: &[&str]) -> Result { - match Command::new(program).args(args).output().await { +/// (non-zero exit, failed spawn, or no exit within `limit`, in which case +/// the child is killed). +async fn run_cmd(program: &str, args: &[&str], limit: Duration) -> Result { + let output = Command::new(program).args(args).kill_on_drop(true).output(); + let Ok(result) = tokio::time::timeout(limit, output).await else { + return Err(format!( + "`{program}` timed out after {}s and was killed", + limit.as_secs_f32() + )); + }; + match result { Ok(out) if out.status.success() => { Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned()) } @@ -176,8 +197,24 @@ fn rfb_pointer_event(button_mask: u8, x: u16, y: u16) -> [u8; 6] { /// /// `events` is a slice of `(button_mask, x, y)` tuples. Each is encoded /// as an RFB `PointerEvent` and written in order; the connection is then -/// closed. Errors at any step short-circuit and return a `String`. -async fn rfb_send_pointer_events(port: u16, events: &[(u8, u16, u16)]) -> Result<(), String> { +/// closed. Errors at any step short-circuit and return a `String`, as does +/// the whole exchange taking longer than `limit`. +async fn rfb_send_pointer_events( + port: u16, + events: &[(u8, u16, u16)], + limit: Duration, +) -> Result<(), String> { + tokio::time::timeout(limit, rfb_exchange(port, events)) + .await + .unwrap_or_else(|_| { + Err(format!( + "VNC: no reply from localhost:{port} within {}s", + limit.as_secs_f32() + )) + }) +} + +async fn rfb_exchange(port: u16, events: &[(u8, u16, u16)]) -> Result<(), String> { let mut stream = TcpStream::connect(("127.0.0.1", port)) .await .map_err(|e| format!("VNC: connect to localhost:{port}: {e}"))?; @@ -210,7 +247,7 @@ impl ScreenMcp { .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| d.as_millis()); let path = format!("/tmp/hive-screenshot-{ms}.png"); - match run_cmd("grim", &["-t", "png", &path]).await { + match run_cmd("grim", &["-t", "png", &path], CMD_TIMEOUT).await { Ok(_) => format!("screenshot saved to `{path}` — use the Read tool to view it"), Err(e) => e, } @@ -225,7 +262,7 @@ impl ScreenMcp { `services.hyperhive.agent.gui.enable = true`)." )] async fn type_text(&self, Parameters(args): Parameters) -> String { - cmd_result(run_cmd("wtype", &[&args.text]).await) + cmd_result(run_cmd("wtype", &[&args.text], CMD_TIMEOUT).await) } #[tool( @@ -255,7 +292,7 @@ impl ScreenMcp { wtype_args.push(m.to_owned()); } let arg_refs: Vec<&str> = wtype_args.iter().map(String::as_str).collect(); - cmd_result(run_cmd("wtype", &arg_refs).await) + cmd_result(run_cmd("wtype", &arg_refs, CMD_TIMEOUT).await) } #[tool( @@ -269,7 +306,7 @@ impl ScreenMcp { let (Ok(x), Ok(y)) = (u16::try_from(args.x), u16::try_from(args.y)) else { return "mouse_move: x and y must be in range 0–65535".to_owned(); }; - match rfb_send_pointer_events(vnc_port(), &[(0, x, y)]).await { + match rfb_send_pointer_events(vnc_port(), &[(0, x, y)], VNC_TIMEOUT).await { Ok(()) => "ok".to_owned(), Err(e) => e, } @@ -299,7 +336,7 @@ impl ScreenMcp { }; // Move to position, press, release — all in one VNC connection. let events = [(0, x, y), (btn_mask, x, y), (0, x, y)]; - match rfb_send_pointer_events(vnc_port(), &events).await { + match rfb_send_pointer_events(vnc_port(), &events, VNC_TIMEOUT).await { Ok(()) => "ok".to_owned(), Err(e) => e, } @@ -358,3 +395,88 @@ async fn main() -> Result<()> { service.waiting().await?; Ok(()) } + +#[cfg(test)] +mod tests { + use super::{Duration, rfb_send_pointer_events, run_cmd}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + /// A child that would create `marker` half a second after it starts. + fn marker_cmd(marker: &std::path::Path) -> String { + format!("sleep 0.5; touch '{}'", marker.display()) + } + + fn marker_path(test: &str) -> std::path::PathBuf { + let path = + std::env::temp_dir().join(format!("hive-screen-mcp-{test}-{}", std::process::id())); + let _ = std::fs::remove_file(&path); + path + } + + #[tokio::test] + async fn run_cmd_times_out_and_kills_the_child() { + let marker = marker_path("killed"); + let err = run_cmd( + "sh", + &["-c", &marker_cmd(&marker)], + Duration::from_millis(50), + ) + .await + .expect_err("a child running past the limit is an error"); + assert!(err.contains("timed out after"), "{err}"); + tokio::time::sleep(Duration::from_secs(1)).await; + assert!(!marker.exists(), "the timed-out child kept running"); + } + + #[tokio::test] + async fn run_cmd_waits_for_a_child_within_the_limit() { + let marker = marker_path("finished"); + run_cmd("sh", &["-c", &marker_cmd(&marker)], Duration::from_secs(10)) + .await + .expect("a child inside the limit succeeds"); + assert!(marker.exists()); + let _ = std::fs::remove_file(&marker); + } + + #[tokio::test] + async fn rfb_times_out_on_a_server_that_never_speaks() { + let listener = TcpListener::bind(("127.0.0.1", 0)).await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let server = tokio::spawn(async move { + let (conn, _) = listener.accept().await.unwrap(); + tokio::time::sleep(Duration::from_secs(5)).await; + drop(conn); + }); + let err = rfb_send_pointer_events(port, &[(0, 1, 1)], Duration::from_millis(200)) + .await + .expect_err("a silent server is an error"); + assert!(err.contains("no reply"), "{err}"); + server.abort(); + } + + #[tokio::test] + async fn rfb_completes_against_a_server_that_answers() { + let listener = TcpListener::bind(("127.0.0.1", 0)).await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let server = tokio::spawn(async move { + let (mut conn, _) = listener.accept().await.unwrap(); + let mut buf = [0u8; 12]; + conn.write_all(b"RFB 003.008\n").await.unwrap(); + conn.read_exact(&mut buf).await.unwrap(); + conn.write_all(&[1, 1]).await.unwrap(); + conn.read_exact(&mut buf[..1]).await.unwrap(); + conn.write_all(&0u32.to_be_bytes()).await.unwrap(); + conn.read_exact(&mut buf[..1]).await.unwrap(); + // ServerInit: width, height, pixel format, empty name. + conn.write_all(&[0u8; 2 + 2 + 16 + 4]).await.unwrap(); + let mut event = [0u8; 6]; + conn.read_exact(&mut event).await.unwrap(); + event + }); + rfb_send_pointer_events(port, &[(0, 1, 2)], Duration::from_secs(10)) + .await + .expect("a well-behaved server completes"); + assert_eq!(server.await.unwrap(), [5, 0, 0, 1, 0, 2]); + } +}