Compare commits

..
2 changed files with 22 additions and 78 deletions

View file

@ -197,22 +197,6 @@ async fn main() -> Result<()> {
} }
} }
/// Banner message for a failing matrix `ensure_all` sweep, shared by both
/// the initial and periodic `record_err` call sites in `cmd_serve` so the
/// wording can't drift between them.
fn matrix_sweep_banner(ctx: sweep_health::SweepFailure) -> String {
let age = ctx.since_last_ok.map_or_else(
|| "no success this session".to_owned(),
|d| format!("last ok {} ago", sweep_health::fmt_age(d)),
);
format!(
"matrix user/space sweep failing ({} consecutive, {age}) \
some agents may be missing matrix accounts, space membership, \
or chat-room invites",
ctx.consecutive
)
}
/// Start the coordinator daemon: open the broker, run migrations, spawn /// Start the coordinator daemon: open the broker, run migrations, spawn
/// background tasks (auto-update, vacuums, crash-watcher, reminder-scheduler, /// background tasks (auto-update, vacuums, crash-watcher, reminder-scheduler,
/// dashboard), then serve the admin socket until a signal arrives. /// dashboard), then serve the admin socket until a signal arrives.
@ -430,24 +414,11 @@ async fn cmd_serve(
let mut matrix_shutdown = coord.shutdown_rx(); let mut matrix_shutdown = coord.shutdown_rx();
tokio::spawn(async move { tokio::spawn(async move {
let interval = std::time::Duration::from_mins(30); let interval = std::time::Duration::from_mins(30);
// Debounced banner: a lone bad sweep (homeserver mid-restart, a matrix::ensure_all().await;
// transient HTTP blip) shouldn't flap the dashboard, but a sweep
// that's been failing for hours (missing agent invites, a broken
// admin token) should surface. Cleared the moment a sweep is clean.
let mut health = sweep_health::SweepHealth::new("matrix_ensure_all", "warn", 2);
if matrix::ensure_all().await {
health.record_ok();
} else {
health.record_err(matrix_sweep_banner);
}
loop { loop {
tokio::select! { tokio::select! {
() = tokio::time::sleep(interval) => { () = tokio::time::sleep(interval) => {
if matrix::ensure_all().await { matrix::ensure_all().await;
health.record_ok();
} else {
health.record_err(matrix_sweep_banner);
}
} }
_ = matrix_shutdown.changed() => { _ = matrix_shutdown.changed() => {
tracing::info!("matrix ensure_all: shutdown signal received"); tracing::info!("matrix ensure_all: shutdown signal received");

View file

@ -1503,28 +1503,20 @@ async fn resolve_room_alias(
} }
/// Sweep every existing container (manager + sub-agents) and ensure /// Sweep every existing container (manager + sub-agents) and ensure
/// each has a matrix user + token on the local homeserver. Called at /// each has a matrix user + token on the local homeserver. Called once
/// hive-c0re startup, alongside `forge::ensure_all`, and then /// at hive-c0re startup, alongside `forge::ensure_all`. No-op when the
/// periodically (see the caller in `main.rs`). No-op when the
/// hive-matrix container isn't running. Per-step failures are logged /// hive-matrix container isn't running. Per-step failures are logged
/// but don't abort the sweep. /// but don't abort the sweep.
/// pub async fn ensure_all() {
/// Returns `true` when every step of the sweep succeeded, `false` when at
/// least one step failed — the caller feeds this into a
/// [`crate::stats::sweep_health::SweepHealth`] to raise a debounced
/// dashboard banner on persistent failure (this sweep re-runs every 30
/// minutes, so a one-off blip self-heals without ever bannering).
pub async fn ensure_all() -> bool {
if !is_present().await { if !is_present().await {
tracing::debug!("matrix: hive-matrix container absent, skipping user sweep"); tracing::debug!("matrix: hive-matrix container absent, skipping user sweep");
return true; return;
} }
let mut ok = true;
let register_token = match ensure_register_token() { let register_token = match ensure_register_token() {
Ok(t) => t, Ok(t) => t,
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: ensure_register_token failed"); tracing::warn!(error = ?e, "matrix: ensure_register_token failed");
return false; return;
} }
}; };
// One HTTP client for the whole sweep — connection pool is // One HTTP client for the whole sweep — connection pool is
@ -1536,7 +1528,7 @@ pub async fn ensure_all() -> bool {
Ok(c) => c, Ok(c) => c,
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: build HTTP client failed; skipping sweep"); tracing::warn!(error = ?e, "matrix: build HTTP client failed; skipping sweep");
return false; return;
} }
}; };
// Provision hive admin user FIRST so it's the first registered // Provision hive admin user FIRST so it's the first registered
@ -1544,11 +1536,10 @@ pub async fn ensure_all() -> bool {
// registered user admin automatically). // registered user admin automatically).
if let Err(e) = ensure_admin_user(&client, &register_token).await { if let Err(e) = ensure_admin_user(&client, &register_token).await {
tracing::warn!(error = ?e, "matrix: ensure_admin_user failed"); tracing::warn!(error = ?e, "matrix: ensure_admin_user failed");
ok = false;
} }
let Ok(containers) = crate::lifecycle::list().await else { let Ok(containers) = crate::lifecycle::list().await else {
tracing::warn!("matrix: nixos-container list failed; skipping user sweep"); tracing::warn!("matrix: nixos-container list failed; skipping user sweep");
return false; return;
}; };
let mut agent_names: Vec<String> = Vec::new(); let mut agent_names: Vec<String> = Vec::new();
for c in &containers { for c in &containers {
@ -1559,45 +1550,33 @@ pub async fn ensure_all() -> bool {
agent_names.push(name.to_owned()); agent_names.push(name.to_owned());
} }
if !provision_space(&client, &agent_names).await { // Provision the hive Space and invite all agents (+ the admin account).
ok = false;
}
ok
}
/// Provision the hive Space + default chat room and invite every agent
/// (+ the admin account) to both. Split out of [`ensure_all`] purely to
/// keep that function under the `too_many_lines` threshold — this is the
/// tail half of the same sequential sweep and shares its aggregate-bool,
/// log-and-continue failure handling.
async fn provision_space(client: &reqwest::Client, agent_names: &[String]) -> bool {
let mut ok = true;
// server_name is needed to form full Matrix user IDs for invites. // server_name is needed to form full Matrix user IDs for invites.
let admin_token = match read_admin_token() { let admin_token = match read_admin_token() {
Ok(t) => t, Ok(t) => t,
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: skipping hive space provisioning (no admin token)"); tracing::warn!(error = ?e, "matrix: skipping hive space provisioning (no admin token)");
return false; return;
} }
}; };
// server_name first — the agent invites need it (fully-qualified user ids). // server_name first — the agent invites need it (fully-qualified user ids).
let server_name = match discover_server_name(client).await { let server_name = match discover_server_name(&client).await {
Ok(s) => s, Ok(s) => s,
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: discover_server_name failed; skipping space provisioning"); tracing::warn!(error = ?e, "matrix: discover_server_name failed; skipping space provisioning");
return false; return;
} }
}; };
let room_id = match ensure_hive_space(client, &admin_token).await { let room_id = match ensure_hive_space(&client, &admin_token).await {
Ok(id) => id, Ok(id) => id,
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: ensure_hive_space failed"); tracing::warn!(error = ?e, "matrix: ensure_hive_space failed");
return false; return;
} }
}; };
// Invite @hive admin first, then all agents. // Invite @hive admin first, then all agents.
if let Err(e) = invite_to_room( if let Err(e) = invite_to_room(
client, &client,
&admin_token, &admin_token,
&room_id, &room_id,
HIVE_ADMIN_LOCALPART, HIVE_ADMIN_LOCALPART,
@ -1606,12 +1585,10 @@ async fn provision_space(client: &reqwest::Client, agent_names: &[String]) -> bo
.await .await
{ {
tracing::warn!(error = ?e, "matrix: invite @hive to space failed"); tracing::warn!(error = ?e, "matrix: invite @hive to space failed");
ok = false;
} }
for name in agent_names { for name in &agent_names {
if let Err(e) = invite_to_room(client, &admin_token, &room_id, name, &server_name).await { if let Err(e) = invite_to_room(&client, &admin_token, &room_id, name, &server_name).await {
tracing::warn!(%name, error = ?e, "matrix: invite agent to space failed"); tracing::warn!(%name, error = ?e, "matrix: invite agent to space failed");
ok = false;
} }
} }
@ -1620,10 +1597,10 @@ async fn provision_space(client: &reqwest::Client, agent_names: &[String]) -> bo
// rooms to chat in (Matrix semantics — children aren't auto-joined), so // rooms to chat in (Matrix semantics — children aren't auto-joined), so
// without this the Space is empty. The restricted join rule additionally // without this the Space is empty. The restricted join rule additionally
// lets the operator (a Space member) join from the Space hierarchy. // lets the operator (a Space member) join from the Space hierarchy.
match ensure_hive_chat_room(client, &admin_token, &room_id, &server_name).await { match ensure_hive_chat_room(&client, &admin_token, &room_id, &server_name).await {
Ok(chat_room_id) => { Ok(chat_room_id) => {
if let Err(e) = invite_to_room( if let Err(e) = invite_to_room(
client, &client,
&admin_token, &admin_token,
&chat_room_id, &chat_room_id,
HIVE_ADMIN_LOCALPART, HIVE_ADMIN_LOCALPART,
@ -1632,23 +1609,19 @@ async fn provision_space(client: &reqwest::Client, agent_names: &[String]) -> bo
.await .await
{ {
tracing::warn!(error = ?e, "matrix: invite @hive to chat room failed"); tracing::warn!(error = ?e, "matrix: invite @hive to chat room failed");
ok = false;
} }
for name in agent_names { for name in &agent_names {
if let Err(e) = if let Err(e) =
invite_to_room(client, &admin_token, &chat_room_id, name, &server_name).await invite_to_room(&client, &admin_token, &chat_room_id, name, &server_name).await
{ {
tracing::warn!(%name, error = ?e, "matrix: invite agent to chat room failed"); tracing::warn!(%name, error = ?e, "matrix: invite agent to chat room failed");
ok = false;
} }
} }
} }
Err(e) => { Err(e) => {
tracing::warn!(error = ?e, "matrix: ensure_hive_chat_room failed"); tracing::warn!(error = ?e, "matrix: ensure_hive_chat_room failed");
ok = false;
} }
} }
ok
} }
#[cfg(test)] #[cfg(test)]