hive-screen-mcp: bound grim/wtype and VNC calls; hive-c0re: make messages match the code

hive-screen-mcp ran `grim` / `wtype` through an unbounded
`Command::output()` and spoke RFB to neatvnc with no deadline, so a
wedged compositor or a VNC server that accepts and never speaks held the
agent's turn forever. Each subprocess now has a 30s limit and is killed
when it hits it; each RFB exchange has a 10s limit. Both come back to the
agent as the tool's text result, like every other failure in this crate.

hive-c0re operator-facing text that described behaviour the code lacks:
- `hivectl matrix reset-password` printed a "next: hivectl matrix
  create-user" hint that fails for every target (agents are refused, a
  non-agent hits M_USER_IN_USE). The line is gone.
- a failed `nixos-container update` appended the container's journal
  tail, read with `journalctl -M`. Since the job DAG, `update` only runs
  from the `Swap` node on a stopped container, so the read always came
  back empty. The helper is removed; the error still points at the
  build log.
- the matrix sweep comment in main.rs said it re-provisions agent token
  files; `ensure_all` creates no agent accounts or tokens.
- the knowledge-pull comments named a webhook caller that no longer
  exists and claimed a race was "fixed at its source".
- `handle_spawn`'s doc and the `mcp_sockets` module doc named callers of
  `register_agent` / a rollback that do not exist.

Refs #4723
This commit is contained in:
atlas 2026-09-26 23:28:09 +02:00
commit 19cc1b12e2
6 changed files with 150 additions and 80 deletions

View file

@ -935,48 +935,9 @@ async fn priv_run_inner(kind: &str, name: &str, node_id: Option<u64>) -> Result<
match result { match result {
Ok(()) => Ok(()), Ok(()) => Ok(()),
Err(e) => { Err(e) => match log_id {
let journal = if kind == "update" { Some(id) => bail!("{e:#}; see build log #{id}"),
container_journal_tail(&container).await None => bail!("{e:#}"),
} 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()
}, },
)
.await;
match res {
Ok((stdout, _)) if !stdout.is_empty() => format!(
"\n--- last 40 journal lines from container '{container}' ---\n{}",
stdout.trim_end()
),
_ => String::new(),
} }
} }

View file

@ -344,17 +344,10 @@ async fn cmd_serve(
} }
} }
}); });
// Knowledge periodic pull: hourly fallback in case the webhook is // Knowledge periodic pull: hourly fallback in case a knowledge event from
// missed (e.g. hive-c0re was down during a push). Deliberately does // the swarm is missed (e.g. hive-c0re was down when it was sent). Sleeps
// NOT also fire an immediate pull at startup the way this task used // before its first pull: the boot `NodeKind::KnowledgePull` DAG node
// to: `auto_update::run`'s `NodeKind::KnowledgePull` DAG node (spawned // already pulls once at startup.
// 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.
let mut knowledge_shutdown = coord.shutdown_rx(); let mut knowledge_shutdown = coord.shutdown_rx();
let knowledge_coord = coord.clone(); let knowledge_coord = coord.clone();
tokio::spawn(async move { tokio::spawn(async move {
@ -394,15 +387,11 @@ async fn cmd_serve(
} }
} }
}); });
// Matrix user sweep: same shape — ensure every container has // Matrix sweep (`matrix::ensure_all`): the hive's own account, the hive
// an account on the local matrix-tuwunel homeserver with an // Space and chat room, and every agent container's invite to both. It
// access_token persisted to `<state>/matrix-token`. No-op when // creates no agent accounts or tokens — swarm-controller mints those.
// the hive-matrix container isn't running. // No-op when the hive-matrix container isn't running. Re-submitted every
// // 30 minutes; the startup pass is the boot `MatrixSweep` node.
// 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.
let mut matrix_shutdown = coord.shutdown_rx(); let mut matrix_shutdown = coord.shutdown_rx();
let matrix_coord = coord.clone(); let matrix_coord = coord.clone();
tokio::spawn(async move { tokio::spawn(async move {

View file

@ -287,8 +287,8 @@ async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
} }
} }
/// Create + start the container for `name`, rolling back socket /// Create + start the container for `name` and bind its MCP listener. On a
/// registration and notifying the manager on failure. /// failed spawn nothing was registered: post a swarm notice and return the error.
async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostResponse> { async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostResponse> {
tracing::info!(%name, "spawn"); tracing::info!(%name, "spawn");
let agent_dir = crate::paths::agent_runtime_dir(name); let agent_dir = crate::paths::agent_runtime_dir(name);
@ -762,7 +762,6 @@ async fn handle_matrix_reset_password(name: &str) -> Result<HostResponse> {
Ok(HostResponse::messages(vec![ Ok(HostResponse::messages(vec![
format!("matrix: password for @{name}:{server_name} reset"), format!("matrix: password for @{name}:{server_name} reset"),
format!("password persisted at: {}", pw_path.display()), format!("password persisted at: {}", pw_path.display()),
format!("next: hivectl matrix create-user {name} # mints a fresh access token"),
])) ]))
} }

View file

@ -111,9 +111,8 @@ pub async fn pull(coord: &Coordinator) -> Result<()> {
// `LOCAL_DIR` after the initial clone — this working tree exists to mirror // `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 // `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 // by any other means (a stray manual edit on the host, an interrupted
// prior operation, or — the actual root cause here — two unsynchronized // prior operation, or two pulls running on this working tree at once)
// boot-time pull callers racing on this same working tree, since fixed // would otherwise abort the `--ff-only` merge below with
// in `main.rs`) would otherwise abort the `--ff-only` merge below with
// "local changes would be overwritten"; an untracked file left behind // "local changes would be overwritten"; an untracked file left behind
// the same ways aborts it with "untracked working tree files would be // the same ways aborts it with "untracked working tree files would be
// overwritten" instead — same wedge, just `clean`'s failure message // overwritten" instead — same wedge, just `clean`'s failure message

View file

@ -6,8 +6,8 @@
//! re-register all running agents. //! re-register all running agents.
//! //!
//! After startup, listeners are managed event-driven: //! After startup, listeners are managed event-driven:
//! - `run_create` calls `register_agent` eagerly on first-spawn. //! - `run_start` (the job queue's `Start` node) and `server::handle_spawn`
//! - `run_reconcile` calls `register_agent` immediately after `start_with_fallback`. //! call `register_agent` once the container is started.
//! - `kill`/`destroy` paths call `unregister_agent`. //! - `kill`/`destroy` paths call `unregister_agent`.
//! //!
//! No recurring poll is needed because c0re owns the listener lifecycle. An //! No recurring poll is needed because c0re owns the listener lifecycle. An

View file

@ -29,16 +29,37 @@ use rmcp::{
transport::stdio, transport::stdio,
}; };
use serde::Deserialize; use serde::Deserialize;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpStream; use tokio::net::TcpStream;
use tokio::process::Command; 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. /// Run a command and return its result.
/// `Ok(stdout_trimmed)` on success (may be empty string when the tool /// `Ok(stdout_trimmed)` on success (may be empty string when the tool
/// produces no output). `Err(human_readable_message)` on failure /// produces no output). `Err(human_readable_message)` on failure
/// (non-zero exit or failed spawn). /// (non-zero exit, failed spawn, or no exit within `limit`, in which case
async fn run_cmd(program: &str, args: &[&str]) -> Result<String, String> { /// the child is killed).
match Command::new(program).args(args).output().await { async fn run_cmd(program: &str, args: &[&str], limit: Duration) -> Result<String, String> {
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(out) if out.status.success() => {
Ok(String::from_utf8_lossy(&out.stdout).trim().to_owned()) 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 /// `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 /// as an RFB `PointerEvent` and written in order; the connection is then
/// closed. Errors at any step short-circuit and return a `String`. /// closed. Errors at any step short-circuit and return a `String`, as does
async fn rfb_send_pointer_events(port: u16, events: &[(u8, u16, u16)]) -> Result<(), String> { /// 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)) let mut stream = TcpStream::connect(("127.0.0.1", port))
.await .await
.map_err(|e| format!("VNC: connect to localhost:{port}: {e}"))?; .map_err(|e| format!("VNC: connect to localhost:{port}: {e}"))?;
@ -210,7 +247,7 @@ impl ScreenMcp {
.duration_since(std::time::UNIX_EPOCH) .duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_millis()); .map_or(0, |d| d.as_millis());
let path = format!("/tmp/hive-screenshot-{ms}.png"); 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"), Ok(_) => format!("screenshot saved to `{path}` — use the Read tool to view it"),
Err(e) => e, Err(e) => e,
} }
@ -225,7 +262,7 @@ impl ScreenMcp {
`services.hyperhive.agent.gui.enable = true`)." `services.hyperhive.agent.gui.enable = true`)."
)] )]
async fn type_text(&self, Parameters(args): Parameters<TypeTextArgs>) -> String { async fn type_text(&self, Parameters(args): Parameters<TypeTextArgs>) -> String {
cmd_result(run_cmd("wtype", &[&args.text]).await) cmd_result(run_cmd("wtype", &[&args.text], CMD_TIMEOUT).await)
} }
#[tool( #[tool(
@ -255,7 +292,7 @@ impl ScreenMcp {
wtype_args.push(m.to_owned()); wtype_args.push(m.to_owned());
} }
let arg_refs: Vec<&str> = wtype_args.iter().map(String::as_str).collect(); 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( #[tool(
@ -269,7 +306,7 @@ impl ScreenMcp {
let (Ok(x), Ok(y)) = (u16::try_from(args.x), u16::try_from(args.y)) else { 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(); 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(), Ok(()) => "ok".to_owned(),
Err(e) => e, Err(e) => e,
} }
@ -299,7 +336,7 @@ impl ScreenMcp {
}; };
// Move to position, press, release — all in one VNC connection. // Move to position, press, release — all in one VNC connection.
let events = [(0, x, y), (btn_mask, x, y), (0, x, y)]; 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(), Ok(()) => "ok".to_owned(),
Err(e) => e, Err(e) => e,
} }
@ -358,3 +395,88 @@ async fn main() -> Result<()> {
service.waiting().await?; service.waiting().await?;
Ok(()) 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]);
}
}