diff --git a/hive-c0re/src/forge/mod.rs b/hive-c0re/src/forge/mod.rs index 498d076e..30c64186 100644 --- a/hive-c0re/src/forge/mod.rs +++ b/hive-c0re/src/forge/mod.rs @@ -418,10 +418,10 @@ pub async fn ensure_all() { } } -/// Ensure a Forgejo `pull_request` org-webhook for `agent-configs` exists and -/// points at hive-c0re's `/webhook/config-pr` endpoint. Idempotent — lists -/// existing hooks first and skips creation when one is already targeting the -/// correct URL. +/// Ensure a Forgejo `pull_request` org-webhook for `agent-configs` targets +/// hive-c0re's `/webhook/config-pr` and signs with `webhook_secret`. A hook +/// already at that URL is kept if [`crate::webhook_secret::is_registered`], +/// else deleted and recreated (Forgejo's edit-hook API ignores `secret`). /// /// `hive_domain` is the public domain name of the hive; the webhook URL is /// `https:///webhook/config-pr` (routed through the gateway, @@ -442,8 +442,8 @@ pub async fn ensure_all() { /// /// Returns an error if: /// - `hive_domain` produces a URL that `url::Url::parse` rejects. -/// - The Forgejo `org_create_hook` API call fails (transport error, auth -/// failure, or the `agent-configs` org does not exist). +/// - Deleting the hook being replaced, or `org_create_hook`, fails (transport +/// error, auth failure, or the `agent-configs` org does not exist). /// - The HTTP call times out (10 s limit). /// /// Listing failures are treated as best-effort: they fall through to the @@ -467,31 +467,39 @@ pub async fn ensure_config_pr_webhook( .await .map_err(anyhow::Error::from) .and_then(|r| r.map_err(anyhow::Error::from)); + let listed_ok = listed.is_ok(); match listed { Ok(hooks) => { - let already_exists = hooks.iter().any(|h| { - h.config - .as_ref() - .and_then(|c| c.get("url")) - .map(String::as_str) - == Some(target_url.as_str()) - }); - if already_exists { + let secret_registered = crate::webhook_secret::is_registered(webhook_secret); + if config_pr_hook_is_current(&hooks, &target_url, secret_registered) { tracing::debug!(%target_url, "forge: config-pr webhook already configured"); return Ok(()); } - // Delete stale hooks that point at our path but a different base - // (e.g. old loopback hooks from before the SSRF-bypass migration). for h in &hooks { - let hook_url = h - .config - .as_ref() - .and_then(|c| c.get("url")) - .map_or("", String::as_str); - if hook_url.ends_with("/webhook/config-pr") - && hook_url != target_url - && let Some(id) = h.id - { + let hook_url = hook_url(h); + let Some(id) = h.id else { continue }; + if hook_url == target_url { + // A failed delete must not fall through to create: the + // old hook would stay beside the new one, and a later + // pass sees the URL present and never removes it. + tracing::info!( + hook_url, + org = CONFIG_ORG, + "forge: replacing config-pr webhook (secret not the registered one)" + ); + tokio::time::timeout( + HTTP_TIMEOUT, + client.org_delete_hook(CONFIG_ORG, id).send(), + ) + .await + .map_err(anyhow::Error::from) + .and_then(|r| r.map_err(anyhow::Error::from)) + .with_context(|| { + format!("delete config-pr webhook {id} to replace its secret") + })?; + } else if hook_url.ends_with("/webhook/config-pr") { + // Our path on a different base (e.g. old loopback hooks + // from before the SSRF-bypass migration). tracing::info!( hook_url, org = CONFIG_ORG, @@ -531,12 +539,63 @@ pub async fn ensure_config_pr_webhook( .and_then(|r| r.map_err(anyhow::Error::from)) .with_context(|| format!("create config-pr webhook on org {CONFIG_ORG}"))?; tracing::info!(%target_url, "forge: config-pr webhook created on org {CONFIG_ORG}"); + // Only a listed pass knows no hook with an older secret survived beside + // this one; an unlisted pass leaves the record stale so the next + // registration lists and replaces again. + if listed_ok && let Err(e) = crate::webhook_secret::record_registered(webhook_secret) { + tracing::warn!(error = ?e, "forge: recording the config-pr webhook secret failed"); + } Ok(()) } +/// The `url` a Forgejo hook delivers to, or `""` when it carries none. +fn hook_url(hook: &forgejo_api::structs::Hook) -> &str { + hook.config + .as_ref() + .and_then(|c| c.get("url")) + .map_or("", String::as_str) +} + +/// Whether the config-PR hook needs no work: one targets `target_url`, and +/// the current secret is the one it was registered with. +fn config_pr_hook_is_current( + hooks: &[forgejo_api::structs::Hook], + target_url: &str, + secret_registered: bool, +) -> bool { + secret_registered && hooks.iter().any(|h| hook_url(h) == target_url) +} + #[cfg(test)] mod tests { - use super::describe_forge_admin; + use super::{config_pr_hook_is_current, describe_forge_admin}; + + const TARGET: &str = "https://hive.example/webhook/config-pr"; + + fn hook(url: &str) -> forgejo_api::structs::Hook { + serde_json::from_value(serde_json::json!({ "id": 1, "config": { "url": url } })) + .expect("hook json") + } + + #[test] + fn a_hook_at_the_target_with_the_registered_secret_is_kept() { + assert!(config_pr_hook_is_current(&[hook(TARGET)], TARGET, true)); + } + + /// Forgejo can't be asked which secret a hook signs with, so a hook at + /// the right URL is not enough: without the matching record it is + /// replaced, or every delivery fails HMAC. + #[test] + fn a_hook_at_the_target_with_a_changed_secret_is_replaced() { + assert!(!config_pr_hook_is_current(&[hook(TARGET)], TARGET, false)); + } + + #[test] + fn a_hook_on_another_base_is_not_ours() { + let stale = hook("http://127.0.0.1:7000/webhook/config-pr"); + assert!(!config_pr_hook_is_current(&[stale], TARGET, true)); + assert!(!config_pr_hook_is_current(&[], TARGET, true)); + } #[test] fn describe_keeps_the_verb_path_and_drops_every_value() { diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 3d953998..0b709e6a 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -316,14 +316,16 @@ async fn run_destroy_bookkeeping(coord: &Arc, agent: &str, purge: b if let Err(e) = sync_meta_after_lifecycle(coord).await { tracing::warn!(error = ?e, %agent, "meta sync after destroy failed"); } - let _ = coord.approvals.fail_pending_for_agent( + if let Err(e) = coord.approvals.fail_pending_for_agent( agent, if purge { "agent purged" } else { "agent destroyed" }, - ); + ) { + tracing::warn!(%agent, error = ?e, "approvals: failing pending on destroy failed"); + } // Drop the durable power intent — a future agent of the same name seeds // fresh from its observed state. if let Err(e) = coord.power.remove(agent) { diff --git a/hive-c0re/src/matrix.rs b/hive-c0re/src/matrix.rs index 4440c4d1..8ce10717 100644 --- a/hive-c0re/src/matrix.rs +++ b/hive-c0re/src/matrix.rs @@ -518,6 +518,28 @@ mod extract_new_password_tests { } } +/// The `event_id` of an admin-room command, from its `PUT .../send` +/// response. It is the anchor separating the bot's reply to this command +/// from older replies in the room, so a response without one is an error: +/// unanchored, an earlier reply (an older reset password) would be returned +/// as this command's result. +async fn sent_event_id(resp: reqwest::Response) -> Result { + let status = resp.status(); + if !status.is_success() { + let body = resp.json::().await.unwrap_or_default(); + anyhow::bail!("matrix: admin room send failed: HTTP {status}, body: {body}"); + } + let body = resp + .json::() + .await + .context("matrix: parse admin room send response")?; + body["event_id"] + .as_str() + .filter(|id| !id.is_empty()) + .map(str::to_owned) + .with_context(|| format!("matrix: admin room send response has no event_id: {body}")) +} + /// Send a command to the Matrix admin room and poll for a bot response. /// /// Strategy: send the command, capture its `event_id`, then poll backwards @@ -552,18 +574,7 @@ async fn admin_room_send_and_poll( .send() .await .context("matrix: PUT admin room message")?; - if !send_resp.status().is_success() { - let body = send_resp - .json::() - .await - .unwrap_or_default(); - anyhow::bail!("matrix: admin room send failed: {body}"); - } - let send_json = send_resp - .json::() - .await - .unwrap_or_default(); - let our_event_id = send_json["event_id"].as_str().unwrap_or("").to_owned(); + let our_event_id = sent_event_id(send_resp).await?; // Poll for bot response: fetch the 20 most recent events (newest-first) // on each tick. Walk the list until we hit our own command event_id; @@ -586,12 +597,7 @@ async fn admin_room_send_and_poll( for event in events { // Stop as soon as we reach our own command — everything // older (further into the list) predates our request. - // If our_event_id is empty (malformed PUT response), we skip - // this guard and inspect all 20 events — slight risk of a - // false match from an older response, but an acceptable fallback. - if !our_event_id.is_empty() - && event["event_id"].as_str() == Some(our_event_id.as_str()) - { + if event["event_id"].as_str() == Some(our_event_id.as_str()) { break; } if event["type"].as_str() != Some("m.room.message") { @@ -1505,8 +1511,8 @@ async fn room_membership( /// Invite a fully-qualified Matrix user id (`@user:server`) to `room_id` /// using the sender token. Idempotent: a user who is already a member or /// already has a pending invite is left untouched (no fresh invite is sent, -/// so they are not re-notified), and a 403 `M_FORBIDDEN` / `M_BAD_STATE` -/// from a racing invite is still treated as success. +/// so they are not re-notified). A refused invite is an error unless the +/// user turns out to be invited or joined anyway — see [`refused_invite`]. async fn invite_user_id( client: &reqwest::Client, sender_token: &str, @@ -1541,16 +1547,31 @@ async fn invite_user_id( tracing::debug!(%user_id, %room_id, "matrix: invited to room"); return Ok(()); } - // 403 with M_FORBIDDEN or M_BAD_STATE typically means the user is - // already a member or has a pending invite — both are fine. - if status == StatusCode::FORBIDDEN { - let body = resp.json::().await.unwrap_or_default(); - let errcode = body["errcode"].as_str().unwrap_or(""); - if errcode == "M_FORBIDDEN" || errcode == "M_BAD_STATE" { - tracing::debug!(%user_id, %room_id, %errcode, "matrix: invite skipped (already member/invited)"); - return Ok(()); - } - anyhow::bail!("matrix: invite {user_id} to {room_id}: HTTP {status}, body: {body}"); + // The pre-check misses a member when its read failed or a concurrent + // invite landed after it, so a refusal re-reads the membership. + let membership = if status == StatusCode::FORBIDDEN { + room_membership(client, sender_token, &encoded_room_id, user_id).await + } else { + None + }; + refused_invite(resp, membership.as_deref(), user_id, room_id).await +} + +/// Outcome of an invite POST the homeserver did not accept. A 403 is success +/// only when `membership`, re-read after the refusal, shows the user invited +/// or joined: `M_FORBIDDEN` covers both "already in the room" and real +/// refusals (banned target, sender lacks power), so the errcode can't tell +/// them apart. +async fn refused_invite( + resp: reqwest::Response, + membership: Option<&str>, + user_id: &str, + room_id: &str, +) -> Result<()> { + let status = resp.status(); + if status == StatusCode::FORBIDDEN && matches!(membership, Some("invite" | "join")) { + tracing::debug!(%user_id, %room_id, "matrix: invite refused, user already member/invited"); + return Ok(()); } let body = resp.json::().await.unwrap_or_default(); anyhow::bail!("matrix: invite {user_id} to {room_id}: HTTP {status}, body: {body}") @@ -1846,6 +1867,75 @@ mod tests { assert_eq!(extract_access_token(&body).unwrap(), "syt_abc123"); } + /// A homeserver response built in memory, so no client (and no TLS + /// roots) is needed to drive the response-handling halves. + fn response(status: u16, body: &'static str) -> reqwest::Response { + axum::http::Response::builder() + .status(status) + .body(body) + .expect("valid mock response") + .into() + } + + const FORBIDDEN_BODY: &str = r#"{"errcode":"M_FORBIDDEN","error":"refused"}"#; + + /// A banned target or a sender without power gets exactly this 403; + /// with no membership behind it, the invite did not happen. + #[tokio::test] + async fn a_refused_invite_without_membership_is_an_error() { + for membership in [None, Some("ban"), Some("leave")] { + let outcome = + refused_invite(response(403, FORBIDDEN_BODY), membership, "@a:x", "!r:x").await; + assert!(outcome.is_err(), "membership {membership:?} must not pass"); + } + } + + #[tokio::test] + async fn a_refused_invite_for_a_member_is_success() { + for membership in ["invite", "join"] { + refused_invite( + response(403, FORBIDDEN_BODY), + Some(membership), + "@a:x", + "!r:x", + ) + .await + .expect("already invited/joined is the goal state"); + } + } + + /// Membership only excuses a 403; any other failure stays a failure. + #[tokio::test] + async fn a_non_403_invite_failure_or_garbage_body_is_an_error() { + for (status, body) in [(500, "not json"), (403, "not json"), (429, "{}")] { + let outcome = refused_invite(response(status, body), None, "@a:x", "!r:x").await; + assert!(outcome.is_err(), "HTTP {status} {body:?} must fail"); + } + let outcome = refused_invite(response(500, "{}"), Some("join"), "@a:x", "!r:x").await; + assert!(outcome.is_err(), "a 500 is not excused by membership"); + } + + #[tokio::test] + async fn an_admin_send_without_an_event_id_is_an_error() { + for (status, body) in [ + (500, r#"{"event_id":"$e"}"#), + (200, "not json"), + (200, "{}"), + (200, r#"{"event_id":""}"#), + ] { + let outcome = sent_event_id(response(status, body)).await; + assert!(outcome.is_err(), "HTTP {status} {body:?} must fail"); + } + } + + #[tokio::test] + async fn an_admin_send_returns_its_event_id() { + let id = sent_event_id(response(200, r#"{"event_id":"$e"}"#)) + .await + .expect("well-formed send response"); + assert_eq!(id, "$e"); + } + #[test] fn extract_access_token_errors_on_missing_field() { let body = serde_json::json!({"user_id": "@alice:matrix.example.org"}); diff --git a/hive-c0re/src/paths.rs b/hive-c0re/src/paths.rs index 5586640a..3f1cdf3e 100644 --- a/hive-c0re/src/paths.rs +++ b/hive-c0re/src/paths.rs @@ -61,8 +61,9 @@ pub fn db_dir() -> PathBuf { /// `webhook-secret` — hex-encoded 32-byte HMAC secret shared between /// hive-c0re's webhook handlers and the Forgejo webhook registrations. -/// Generated on first startup and persisted; Forgejo is re-registered -/// whenever the secret changes. +/// Generated on first startup and persisted. The config-PR hook is replaced +/// at the next boot-time registration when this no longer matches +/// [`forge_config_pr_webhook_fingerprint`]. #[must_use] pub fn webhook_secret_file() -> PathBuf { state_root().join("webhook-secret") @@ -74,6 +75,15 @@ pub fn forge_dir() -> PathBuf { state_root().join("forge") } +/// `forge/config-pr-webhook-secret-sha256` — SHA-256 of the webhook secret +/// the config-PR hook was last registered with. Forgejo never returns a +/// hook's secret, so this is the only record of which one it signs with; +/// delete to force the hook to be replaced. +#[must_use] +pub fn forge_config_pr_webhook_fingerprint() -> PathBuf { + forge_dir().join("config-pr-webhook-secret-sha256") +} + /// `forge/core-avatar-set` — marker: core account avatar uploaded. #[must_use] pub fn forge_core_avatar_marker() -> PathBuf { diff --git a/hive-c0re/src/webhook_secret.rs b/hive-c0re/src/webhook_secret.rs index e3207d2e..e4f44fca 100644 --- a/hive-c0re/src/webhook_secret.rs +++ b/hive-c0re/src/webhook_secret.rs @@ -44,6 +44,36 @@ fn load_or_generate_at(path: &std::path::Path) -> Result { Ok(secret) } +/// Whether `secret` is the one the config-PR hook was last registered with, +/// per [`crate::paths::forge_config_pr_webhook_fingerprint()`]. A missing or +/// unreadable record counts as "not registered". +pub fn is_registered(secret: &str) -> bool { + is_registered_at(&crate::paths::forge_config_pr_webhook_fingerprint(), secret) +} + +/// Record `secret` as the one the config-PR hook is now registered with. +pub fn record_registered(secret: &str) -> Result<()> { + record_registered_at(&crate::paths::forge_config_pr_webhook_fingerprint(), secret) +} + +fn is_registered_at(path: &std::path::Path, secret: &str) -> bool { + std::fs::read_to_string(path).is_ok_and(|raw| raw.trim() == fingerprint(secret)) +} + +fn record_registered_at(path: &std::path::Path, secret: &str) -> Result<()> { + std::fs::create_dir_all(path.parent().unwrap_or(path)) + .with_context(|| format!("create dir for {}", path.display()))?; + std::fs::write(path, format!("{}\n", fingerprint(secret))) + .with_context(|| format!("write webhook secret fingerprint to {}", path.display())) +} + +/// Hex SHA-256 of `secret` — stored in its place so the record on disk is +/// not a second copy of the key. +fn fingerprint(secret: &str) -> String { + use sha2::{Digest as _, Sha256}; + hex_encode(&Sha256::digest(secret.as_bytes())) +} + /// Read 32 random bytes from `/dev/urandom` and hex-encode them. fn generate_hex_secret() -> Result { use std::io::Read as _; @@ -106,7 +136,46 @@ fn hex_decode(s: &str) -> Option> { #[cfg(test)] mod tests { - use super::{hex_encode, load_or_generate_at, verify_signature}; + use super::{ + hex_encode, is_registered_at, load_or_generate_at, record_registered_at, verify_signature, + }; + + /// No record is the state of every hive before its first replacement, + /// and must read as "not registered" so that hook gets replaced once. + #[test] + fn a_secret_with_no_fingerprint_on_record_is_not_registered() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("forge").join("fingerprint"); + assert!(!is_registered_at(&path, &"a".repeat(64))); + } + + #[test] + fn a_recorded_secret_is_registered_and_a_changed_one_is_not() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("forge").join("fingerprint"); + let old = "a".repeat(64); + let new = "b".repeat(64); + + record_registered_at(&path, &old).expect("record"); + assert!( + is_registered_at(&path, &old), + "unchanged secret: keep the hook" + ); + assert!( + !is_registered_at(&path, &new), + "changed secret: replace the hook" + ); + assert!( + !std::fs::read_to_string(&path) + .expect("read back") + .contains(&old), + "the record must not be a copy of the secret" + ); + + record_registered_at(&path, &new).expect("re-record"); + assert!(is_registered_at(&path, &new)); + assert!(!is_registered_at(&path, &old)); + } /// Compute the `sha256=` header Forgejo would send for `secret` + /// `body`, so the "matches" test below isn't asserting against a