//! Optional matrix-tuwunel wiring: the hive's appservice identity (host) + //! per-agent account creation → `/matrix-token`. No-op //! when the `hive-matrix` container isn't running, so operators who //! haven't flipped `services.hyperhive.deploy.matrix.enable = true` pay //! nothing. //! //! Accounts are created **as the hive's appservice**, not by presenting a //! shared registration token in a UIAA flow. The difference that matters //! here is not the round-trip count: an appservice token is an *identity* //! the homeserver knows, so the secret never has to be the same on both //! sides of the wire, and account creation does not depend on registration //! being open to anyone who learns a token. //! //! See `docs/integrations/matrix.md::Provisioning flow (appservice)` for the //! registration file's shape, how its token reaches both halves, and the //! host/container bind-mount layout. use std::path::PathBuf; use anyhow::{Context, Result}; use reqwest::StatusCode; /// Client-server API base this daemon provisions against, from /// `HIVE_MATRIX_API_URL` (set by the hyperhive NixOS module from /// `services.hyperhive.swarm.matrix.apiUrl`). /// /// `None` means **this hive has no homeserver to provision against** and /// every matrix path no-ops — see [`is_present`]. There is deliberately no /// fallback: `localhost:8008` is right only when the homeserver happens to /// share this daemon's netns, and an address baked into the binary is one /// that builds fine and then talks to the wrong machine. The nix module /// supplies the loopback address when it is itself the thing running /// tuwunel, where it is not a guess but a fact about what it just started. /// /// Also not an agent-facing address either way. An agent has its own netns; /// agents are handed the gateway vhost via `HIVE_MATRIX_URL`, and get /// nothing at all when the hive has no vhost to offer. fn matrix_http() -> Option<&'static str> { static BASE: std::sync::OnceLock> = std::sync::OnceLock::new(); BASE.get_or_init(|| std::env::var("HIVE_MATRIX_API_URL").ok()) .as_deref() } /// [`matrix_http`] for the call sites that propagate with `?`. /// /// # Errors /// When no homeserver is configured. Reaching one of these paths at all /// means an [`is_present`] gate was skipped, so the message names that /// rather than the missing variable. fn matrix_base() -> Result<&'static str> { matrix_http().context( "matrix: no homeserver configured \ (services.hyperhive.swarm.matrix.apiUrl / HIVE_MATRIX_API_URL) — \ this path should have been gated on matrix::is_present()", ) } /// HTTP timeout for registration round-trips. Account creation is one /// POST; even the slow path should finish well inside this budget. const HTTP_TIMEOUT_SECS: u64 = 10; /// Length (bytes) of the throwaway per-agent matrix password. Random /// 32-byte hex — agents never log in with the password (they /// authenticate by `access_token`), so it's protocol overhead. We /// store it nowhere. const PASSWORD_BYTES: usize = 32; /// Matrix localpart this hive acts as. Not an agent; has no state dir. /// /// **Per hive, derived from the hive's name** — `hive-`, built by /// `swarm_secret_client::matrix::hive_localpart`, which is the one place the /// derivation lives. It used to be the bare constant `hive`, and a swarm runs /// one homeserver: every hive on it logged in as the same `@hive:`, so the /// homeserver could not attribute an action to a hive and revoking one hive's /// standing revoked all of them. /// /// Also the `sender_localpart` of the hive's appservice registration, /// which is what creates this account on a homeserver that has never had /// one: the homeserver creates a registration's sender user itself, at /// startup, before it accepts a request. See /// `nix/host-modules/hive-matrix.nix`, where the same string is rendered as /// `hiveLocalpart`. **The two must match**; nothing wires an override /// across. /// /// An ordinary account, with no homeserver-admin standing: what it /// provisions — the Space, the chat room, the invites — it provisions as /// the creator of those rooms. /// /// # Errors /// When `HYPERHIVE_HIVE_NAME` is unset or empty. That is the same variable /// [`stored_sender_token`] logs in with and `hive-c0re`'s meta flake asserts /// for every hive, so its absence is a broken deployment rather than a /// supported one — and there is no safe default to fall back to: the only one /// available is the shared account this replaced. pub fn hive_localpart() -> Result { let hive = std::env::var("HYPERHIVE_HIVE_NAME") .ok() .filter(|h| !h.is_empty()) .context( "matrix: HYPERHIVE_HIVE_NAME is unset, so this hive has no matrix account of its own \ — the localpart is derived from the hive name and has no swarm-wide fallback", )?; Ok(swarm_secret_client::matrix::hive_localpart(&hive)) } /// Display name of the hive Space. Plain text, no special characters, so /// the Space stays rediscoverable by name (no room alias needed) even when /// the persisted room-id file is lost — preventing duplicate spaces from /// being created on the next sweep. pub const HIVE_SPACE_NAME: &str = "hive"; /// Display name of the default "hive chat" room — the `m.space.child` of /// the hive Space that every agent + the operator can join. Plain text so /// it stays rediscoverable by name (mirrors [`HIVE_SPACE_NAME`]) when the /// persisted room-id file is lost, preventing duplicate chat rooms. pub const HIVE_CHAT_ROOM_NAME: &str = "hive-chat"; /// Topic for the default hive chat room. const HIVE_CHAT_ROOM_TOPIC: &str = "Hive-wide chat for all agents and the operator. Auto-provisioned by hive-c0re."; /// Host path for the appservice sender account's matrix access token. Outside every /// purgeable path — not deleted by `destroy --purge` on any agent. #[must_use] pub fn sender_token_path() -> PathBuf { crate::paths::matrix_sender_token() } /// Password file for the hive's own `@hive-:` account. Stored OUTSIDE /// the purgeable `agent_state_root` tree so it survives `destroy --purge`. /// /// Path: `/var/lib/hyperhive/matrix/creds/-password` fn password_path(name: &str) -> PathBuf { crate::paths::matrix_creds_dir().join(format!("{name}-password")) } /// Host path where the hive Matrix Space room ID is persisted. /// Outside every purgeable path — not deleted by `destroy --purge`. #[must_use] pub fn hive_space_room_id_path() -> PathBuf { crate::paths::matrix_space_room_id() } /// Host path where the default hive chat room ID is persisted. Outside /// every purgeable path — not deleted by `destroy --purge`. #[must_use] pub fn hive_chat_room_id_path() -> PathBuf { crate::paths::matrix_chat_room_id() } /// Whether this hive has a homeserver to provision against. /// /// **A configured API URL, not a local container.** It used to scan /// `nixos-container list` for `hive-matrix`, which answers a different /// question — "is the homeserver a container on this host" — and so made a /// remote homeserver silently no-op no matter how it was addressed. The /// nix module still supplies the loopback URL whenever it runs tuwunel /// itself, so a co-located hive behaves exactly as before. #[must_use] pub fn is_present() -> bool { matrix_http().is_some() } /// Read `n` cryptographic-quality bytes from `/dev/urandom` and return /// them hex-encoded. Avoids pulling a workspace `rand` dep just for /// 32 bytes of randomness; the kernel's CSPRNG is more than enough for /// a long-lived shared secret on the same host. fn random_hex(n: usize) -> Result { use std::io::Read; let mut buf = vec![0_u8; n]; let mut f = std::fs::File::open("/dev/urandom").context("open /dev/urandom")?; f.read_exact(&mut buf).context("read /dev/urandom")?; let mut hex = String::with_capacity(n * 2); for b in &buf { use std::fmt::Write as _; write!(hex, "{b:02x}").ok(); } Ok(hex) } /// Read the hive's appservice token — the `as_token` of the registration /// the homeserver loaded at boot. Every account this module creates is /// authorised by it. /// /// **Reads, never mints**, unlike the registration token it replaced. /// That token was the whole agreement, so whichever side wrote it first /// was right; this one has a second half — the registration file naming /// it, which only the nix side writes. A token minted here would be a /// token the homeserver has never heard of, and the failure would surface /// as every request being refused rather than as a missing file. /// /// # Errors /// When the file is absent or empty. That means the host activation /// script has not run on this generation yet; callers log it and leave /// existing accounts alone rather than trying to proceed. pub fn read_appservice_token() -> Result { let path = crate::paths::matrix_appservice_token(); std::fs::read_to_string(&path) .ok() .map(|s| s.trim().to_owned()) .filter(|s| !s.is_empty()) .with_context(|| { format!( "matrix appservice token not found at {} — it is minted by the \ hive-matrix activation script, which also renders the registration \ file naming it; deploy the hive-matrix module (or re-run \ `nixos-rebuild switch`) before provisioning matrix users", path.display() ) }) } /// Build the localpart of a matrix user id for `agent`. Matrix /// usernames are 1-255 chars from the `[a-z0-9._=-/]` alphabet; agent /// names already conform (hyperhive enforces a strict subset), so no /// escaping is needed at the boundary. fn user_localpart(agent: &str) -> &str { agent } /// Send the registration POST as the appservice and parse the response. /// Returns `Ok((status, body))` on any completed HTTP round-trip /// (including the `M_USER_IN_USE` 400 the caller treats as "already /// exists"); errors only on transport failure. async fn register_post( client: &reqwest::Client, as_token: &str, body: &serde_json::Value, ) -> Result<(StatusCode, serde_json::Value)> { let base = matrix_base()?; let url = format!("{base}/_matrix/client/v3/register"); let resp = client .post(&url) .bearer_auth(as_token) .json(body) .send() .await .context("matrix: POST /register")?; let status = resp.status(); let json = resp .json::() .await .context("matrix: parse /register response")?; Ok((status, json)) } /// Generate a throwaway random password for matrix UIAA registration. /// `PASSWORD_BYTES` raw bytes ⇒ 64-char hex string. The hive sender /// account authenticates by `access_token`, so this password is /// protocol overhead never used for login. pub fn random_password() -> Result { random_hex(PASSWORD_BYTES) } /// Create the matrix account for `agent` as the hive's appservice and /// return an access token for it. One round-trip: an appservice-typed /// registration needs no UIAA stage at all, so there is no session to /// carry and no shared secret to present. /// /// The account is created **by** the appservice but is an ordinary user /// afterwards — it gets its own device and its own access token, and the /// agent authenticates with that rather than with anything the hive /// holds. The `as_token` never leaves the host. /// /// Its one caller, the hive sender account's provisioning, always passes /// a [`random_password`] throwaway — the account authenticates by /// `access_token`, never `m.login.password`. /// /// # Errors /// Propagates the homeserver's own body, which is what the /// `M_USER_IN_USE` callers match on. A `M_EXCLUSIVE` body means the /// localpart falls outside the appservice's namespace — the registration /// file's `namespaces.users` regex is the place to look, not this call. async fn register_user( client: &reqwest::Client, agent: &str, as_token: &str, password: &str, ) -> Result { let localpart = user_localpart(agent); let body = serde_json::json!({ // What makes this an appservice registration rather than an // ordinary one. Without it the homeserver treats the request as a // normal client's and asks for a UIAA flow — even holding the // as_token. "type": "m.login.application_service", "username": localpart, "password": password, // device_id stays stable across re-runs so a re-mint doesn't // strand orphan devices in tuwunel. "device_id": format!("hyperhive-{agent}"), "initial_device_display_name": format!("hyperhive ({agent})"), "inhibit_login": false, }); let (status, body) = register_post(client, as_token, &body).await?; if !status.is_success() { anyhow::bail!("matrix: /register as appservice HTTP {status}, body: {body}"); } extract_access_token(&body) } /// Log in as an **existing** account using the hive's appservice token, /// and return a fresh access token for it. No password involved: the /// appservice is authorised for every localpart in its namespace, so it /// can mint a session for one without knowing anything about the account. /// /// This is the recovery path that used to need a stored password or an /// admin-room password reset — an account whose token file was lost is /// re-tokened from the hive's own identity instead. The device id matches /// [`register_user`]'s, so a re-login replaces that device's token rather /// than accumulating devices. async fn appservice_login(client: &reqwest::Client, as_token: &str, agent: &str) -> Result { let base = matrix_base()?; let url = format!("{base}/_matrix/client/v3/login"); let body = serde_json::json!({ "type": "m.login.application_service", "identifier": { "type": "m.id.user", "user": user_localpart(agent), }, "device_id": format!("hyperhive-{agent}"), "initial_device_display_name": format!("hyperhive ({agent})"), }); let resp = client .post(&url) .bearer_auth(as_token) .json(&body) .send() .await .context("matrix: POST /login as appservice")?; let status = resp.status(); let json = resp .json::() .await .context("matrix: parse appservice /login response")?; if !status.is_success() { anyhow::bail!("matrix: appservice /login HTTP {status} for {agent}, body: {json}"); } extract_access_token(&json) } /// Pull `access_token` out of a successful /register response. fn extract_access_token(body: &serde_json::Value) -> Result { body["access_token"] .as_str() .map(str::to_owned) .with_context(|| format!("matrix: missing access_token in response: {body}")) } /// Login with `m.login.password` and return the access token. Fallback /// for when registration fails with `M_USER_IN_USE` — the account /// already exists in the homeserver but the token file was lost. Fails /// if the stored password no longer matches (e.g. homeserver wiped); /// the only caller is the hive sender account's own recovery path /// (`hivectl matrix sync-admin`). async fn login_user(client: &reqwest::Client, agent: &str, password: &str) -> Result { let base = matrix_base()?; let url = format!("{base}/_matrix/client/v3/login"); let body = serde_json::json!({ "type": "m.login.password", "identifier": { "type": "m.id.user", "user": user_localpart(agent), }, "password": password, "device_id": format!("hyperhive-{agent}"), "initial_device_display_name": format!("hyperhive ({agent})"), }); let resp = client .post(&url) .json(&body) .send() .await .context("matrix: POST /login")?; let status = resp.status(); let json = resp .json::() .await .context("matrix: parse /login response")?; if !status.is_success() { anyhow::bail!("matrix: /login HTTP {status} for agent {agent}, body: {json}"); } extract_access_token(&json) } /// Percent-encode a matrix room ID for use in a URL path segment. /// Only `:` needs encoding; `!` and alphanumerics are path-safe. fn encode_room_id_for_url(room_id: &str) -> String { room_id.replace(':', "%3A") } /// Ensure the hive's `@hive-:` matrix user exists and that its access token is /// persisted at [`sender_token_path()`]. /// /// **Nothing here depends on registration order, and nothing here is /// privileged.** The account used to have to be the first ever /// registered, to win tuwunel's automatic first-user grant — a rule that /// cannot fire for an appservice-created account at all. It is now an /// ordinary account: the homeserver creates it because it is the /// appservice registration's `sender_localpart`, and everything the hive /// provisions with it, it provisions as the creator of those rooms. /// /// Idempotent — skips the account work when the token file already exists /// and is non-empty. /// /// The token is taken from the **swarm secret store** when it is there: /// `swarm-matrix-ctl`, the oneshot inside the matrix container, publishes /// it under an identity of its own, and taking it from there is what lets a /// hive that holds no `as_token` have an admin at all. The mint ladder below /// stays as the fallback for a store that is empty, unconfigured or /// unreachable — which is every swarm whose matrix container predates that /// binary. pub async fn ensure_hive_user(client: &reqwest::Client, as_token: &str) -> Result<()> { use std::os::unix::fs::PermissionsExt; let path = sender_token_path(); if path.exists() && let Ok(existing) = std::fs::read_to_string(&path) && !existing.trim().is_empty() { tracing::debug!("matrix: the sender token is already present"); return Ok(()); } if let Some(token) = stored_sender_token().await { return persist_sender_token(&path, &token); } // Per hive, and fatal when it cannot be derived: the fallback ladder below // must not mint under some other hive's name. let localpart = hive_localpart()?; let password = random_password()?; let access_token = match register_user(client, &localpart, as_token, &password).await { Ok(token) => { let pw_path = password_path(&localpart); if let Some(parent) = pw_path.parent() { std::fs::create_dir_all(parent).ok(); } if let Err(e) = std::fs::write(&pw_path, format!("{password}\n")) { tracing::warn!(error = ?e, %localpart, "matrix: failed to persist the sender account password"); } else { let _ = std::fs::set_permissions(&pw_path, std::fs::Permissions::from_mode(0o600)); } token } Err(reg_err) if reg_err.to_string().contains("M_USER_IN_USE") => { // The expected path, not an edge case: this account is the // appservice's own `sender_localpart`, so the homeserver // creates it when it loads the registration — before // hive-c0re gets a chance to ask. An appservice login needs // no password, which is just as well since an account the // homeserver created has none. tracing::info!(%localpart, "matrix: the sender account already exists, logging in as the appservice"); match appservice_login(client, as_token, &localpart).await { Ok(token) => token, Err(e) => { tracing::warn!(error = ?e, %localpart, "matrix: appservice login for the sender account failed; falling back to the stored password"); let pw_path = password_path(&localpart); let stored = std::fs::read_to_string(&pw_path) .ok() .map(|s| s.trim().to_owned()) .filter(|s| !s.is_empty()) .with_context(|| { format!( "matrix: @{localpart}: exists, appservice login failed, and no \ password is stored at {} — check that the registration file's \ namespace covers @{localpart} and that the homeserver \ loaded it", pw_path.display() ) })?; login_user(client, &localpart, &stored).await? } } } Err(other) => return Err(other), }; persist_sender_token(&path, &access_token) } /// Write the appservice sender account's access token to `path`, 0600, creating the /// directory if it is not there. /// /// Shared by both arms of [`ensure_hive_user`] rather than duplicated into /// the store one: the file's mode is the only thing keeping an unprivileged /// reader off the hive's matrix credential, and a second copy of that decision /// is one that can be edited alone. fn persist_sender_token(path: &std::path::Path, access_token: &str) -> Result<()> { use std::os::unix::fs::PermissionsExt; if let Some(parent) = path.parent() { std::fs::create_dir_all(parent).ok(); } std::fs::write(path, format!("{access_token}\n")) .with_context(|| format!("matrix: write the sender token to {}", path.display()))?; let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)); tracing::info!(path = %path.display(), "matrix: provisioned the sender token"); Ok(()) } /// Fetch the sender token `swarm-matrix-ctl` published, under /// this hive's own store identity. /// /// The cert role is the hive's name, straight out of `HYPERHIVE_HIVE_NAME` — /// the same role string `workers::credential` logs in with, and already in /// this process's environment, so the store read costs no plumbing through /// [`ensure_all`]. That name is also the path's own hive segment, which is /// what makes this read reach **this** hive's token and 403 on any other's: /// the grant behind it is the hive's own `swarm/hives//*` stanza, not /// the swarm-wide `swarm/services/*` one the shared token used to sit under. /// /// `None`, never an error, for every way this can come up empty — no hive /// name, no `BAO_*` identity, an unreachable store, nothing at the path. All /// four mean the same thing to the caller ("mint it the old way"), and three /// of them are the ordinary state of a swarm that has not deployed `swarm-matrix-ctl` /// yet, so raising would turn a supported deployment into a warning every /// sweep. /// /// 🩸 Logs the store **path** and never the value. async fn stored_sender_token() -> Option { let hive = std::env::var("HYPERHIVE_HIVE_NAME") .ok() .filter(|h| !h.is_empty())?; let path = match swarm_secret_client::matrix::sender_token_path(&hive) { Ok(path) => path, Err(e) => { tracing::warn!(error = %e, "matrix: this hive's name does not form a store path"); return None; } }; let store = match swarm_secret_client::SecretStore::from_env(&hive).await { Ok(store) => store, Err(e) => { tracing::debug!(error = %e, "matrix: no swarm secret store to read the sender token from"); return None; } }; match store .read::(&path) .await { Ok(credential) if !credential.value.trim().is_empty() => { tracing::info!(%path, "matrix: taking the sender token from the swarm store"); Some(credential.value) } Ok(_) => { tracing::warn!(%path, "matrix: the stored sender token is empty; minting instead"); None } Err(e) => { tracing::debug!(%path, error = %e, "matrix: no sender token in the store; minting instead"); None } } } /// Discover the matrix `server_name` from the running homeserver via /// `GET /_matrix/key/v2/server` (unauthenticated federation key endpoint). /// The response JSON always includes `"server_name"` per the matrix spec. pub async fn discover_server_name(client: &reqwest::Client) -> Result { let base = matrix_base()?; let url = format!("{base}/_matrix/key/v2/server"); let resp = client .get(&url) .send() .await .context("matrix: GET /_matrix/key/v2/server")?; let status = resp.status(); let body = resp .json::() .await .context("matrix: parse /_matrix/key/v2/server response")?; if !status.is_success() { anyhow::bail!("matrix: /_matrix/key/v2/server returned HTTP {status}, body: {body}"); } body["server_name"] .as_str() .map(str::to_owned) .with_context(|| { format!("matrix: /_matrix/key/v2/server response missing server_name field: {body}") }) } /// Read the sender token from disk. Returns an error if it /// is absent — callers should gate their homeserver calls on this. pub fn read_sender_token() -> Result { let path = sender_token_path(); std::fs::read_to_string(&path) .ok() .map(|s| s.trim().to_owned()) .filter(|s| !s.is_empty()) .with_context(|| { format!( "the matrix sender token was not found at {} — \ ensure hive-c0re has started at least once with matrix enabled \ (it provisions the hive's @hive-: account on boot)", path.display() ) }) } /// Persist the hive Space room id to [`hive_space_room_id_path()`] (0600). fn persist_space_room_id(room_id: &str) -> Result<()> { use std::os::unix::fs::PermissionsExt; let path = hive_space_room_id_path(); if let Some(parent) = path.parent() { std::fs::create_dir_all(parent).ok(); } std::fs::write(&path, format!("{room_id}\n")) .with_context(|| format!("matrix: write space room_id to {}", path.display()))?; let _ = std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)); Ok(()) } /// Scan the admin account's joined rooms for the canonical hive Space: the /// `m.space` whose name is [`HIVE_SPACE_NAME`]. Returns the first match /// (deterministic per the homeserver's joined-rooms order) so a lost /// room-id file recovers the existing space instead of spawning a /// duplicate. `None` if the homeserver is unreachable or no match exists. /// /// Name-based (not alias-based) rediscovery keeps the Space free of any /// special-char room alias — the hardcoded plain name is the anchor. async fn find_space_by_name(client: &reqwest::Client, sender_token: &str) -> Option { let base = matrix_http()?; let joined_url = format!("{base}/_matrix/client/v3/joined_rooms"); let joined: serde_json::Value = client .get(&joined_url) .bearer_auth(sender_token) .send() .await .ok()? .json() .await .ok()?; let rooms = joined["joined_rooms"].as_array()?; for room in rooms { let Some(room_id) = room.as_str() else { continue; }; let encoded = encode_room_id_for_url(room_id); // Must be an m.space (m.room.create `type`). let create_url = format!("{base}/_matrix/client/v3/rooms/{encoded}/state/m.room.create/"); let is_space = match client .get(&create_url) .bearer_auth(sender_token) .send() .await { Ok(r) if r.status().is_success() => r .json::() .await .ok() .is_some_and(|c| c["type"].as_str() == Some("m.space")), _ => false, }; if !is_space { continue; } // …and named HIVE_SPACE_NAME (m.room.name `name`). let name_url = format!("{base}/_matrix/client/v3/rooms/{encoded}/state/m.room.name/"); let name_matches = match client.get(&name_url).bearer_auth(sender_token).send().await { Ok(r) if r.status().is_success() => r .json::() .await .ok() .is_some_and(|n| n["name"].as_str() == Some(HIVE_SPACE_NAME)), _ => false, }; if name_matches { return Some(room_id.to_owned()); } } None } /// Create (or recover) the hive Matrix Space and persist its room ID to /// [`hive_space_room_id_path()`]. The Space is a private `m.space` owned by /// `@hive-:`, identified by its hardcoded name [`HIVE_SPACE_NAME`] (no alias). /// /// Dedup strategy (single canonical space): /// 1. If the room-id file exists, reuse it. /// 2. Otherwise, rediscover by scanning the admin's joined rooms for the /// `m.space` named [`HIVE_SPACE_NAME`] and adopt it (re-persisting the /// file). This recovers the existing space after a state wipe instead /// of creating a duplicate. /// 3. Only if neither yields a room do we `createRoom`. /// /// # Errors /// /// Returns an error if the homeserver is unreachable, `createRoom` fails, /// or the room-ID file cannot be written. pub async fn ensure_hive_space(client: &reqwest::Client, sender_token: &str) -> Result { let base = matrix_base()?; // 1. Stored room id wins (fast path). if let Ok(existing) = std::fs::read_to_string(hive_space_room_id_path()) { let trimmed = existing.trim().to_owned(); if !trimmed.is_empty() { tracing::debug!(room_id = %trimmed, "matrix: hive space already provisioned"); return Ok(trimmed); } } // 2. No stored id — rediscover the existing space by its hardcoded name // before creating a new one (prevents duplicate spaces after a wipe). if let Some(room_id) = find_space_by_name(client, sender_token).await { persist_space_room_id(&room_id)?; tracing::info!(%room_id, "matrix: recovered hive space by name"); return Ok(room_id); } // 3. Create the space (plain hardcoded name, no alias). let url = format!("{base}/_matrix/client/v3/createRoom"); let body = serde_json::json!({ "name": HIVE_SPACE_NAME, "creation_content": { "type": "m.space" }, "preset": "private_chat", "visibility": "private", }); let resp = client .post(&url) .bearer_auth(sender_token) .json(&body) .send() .await .context("matrix: POST /createRoom (hive space)")?; let status = resp.status(); let json = resp .json::() .await .context("matrix: parse /createRoom response")?; if !status.is_success() { anyhow::bail!("matrix: createRoom HTTP {status}, body: {json}"); } let room_id = json["room_id"] .as_str() .with_context(|| format!("matrix: createRoom missing room_id: {json}"))? .to_owned(); persist_space_room_id(&room_id)?; tracing::info!(%room_id, "matrix: created hive space"); Ok(room_id) } /// Persist the hive chat room id to [`hive_chat_room_id_path()`] (0600). fn persist_chat_room_id(room_id: &str) -> Result<()> { use std::os::unix::fs::PermissionsExt; let path = hive_chat_room_id_path(); if let Some(parent) = path.parent() { std::fs::create_dir_all(parent).ok(); } std::fs::write(&path, format!("{room_id}\n")) .with_context(|| format!("matrix: write chat room_id to {}", path.display()))?; let _ = std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)); Ok(()) } /// Scan the admin account's joined rooms for the canonical hive chat room: /// a non-space room named [`HIVE_CHAT_ROOM_NAME`]. Mirrors /// [`find_space_by_name`] so a lost room-id file recovers the existing chat /// room instead of spawning a duplicate. `None` if the homeserver is /// unreachable or no match exists. async fn find_chat_room_by_name(client: &reqwest::Client, sender_token: &str) -> Option { let base = matrix_http()?; let joined_url = format!("{base}/_matrix/client/v3/joined_rooms"); let joined: serde_json::Value = client .get(&joined_url) .bearer_auth(sender_token) .send() .await .ok()? .json() .await .ok()?; let rooms = joined["joined_rooms"].as_array()?; for room in rooms { let Some(room_id) = room.as_str() else { continue; }; let encoded = encode_room_id_for_url(room_id); // Skip the Space itself (and any other m.space). let create_url = format!("{base}/_matrix/client/v3/rooms/{encoded}/state/m.room.create/"); let is_space = match client .get(&create_url) .bearer_auth(sender_token) .send() .await { Ok(r) if r.status().is_success() => r .json::() .await .ok() .is_some_and(|c| c["type"].as_str() == Some("m.space")), _ => false, }; if is_space { continue; } // …and named HIVE_CHAT_ROOM_NAME (m.room.name `name`). let name_url = format!("{base}/_matrix/client/v3/rooms/{encoded}/state/m.room.name/"); let name_matches = match client.get(&name_url).bearer_auth(sender_token).send().await { Ok(r) if r.status().is_success() => r .json::() .await .ok() .is_some_and(|n| n["name"].as_str() == Some(HIVE_CHAT_ROOM_NAME)), _ => false, }; if name_matches { return Some(room_id.to_owned()); } } None } /// Whether the state event still has to be written. /// /// `current` is the room's existing content for this `(type, state_key)`, /// or `None` when it could not be read. **Unreadable means write**: the /// re-apply exists so a missing link gets repaired, and one redundant /// event is cheaper than a hierarchy that never reconverges. /// /// Comparison is `serde_json::Value` equality, which compares objects as /// maps — key order in the response does not matter. fn state_needs_write(current: Option<&serde_json::Value>, desired: &serde_json::Value) -> bool { current != Some(desired) } /// Read the room's current content for one state event. `None` on any /// failure — absent state, transport error, or an unparseable body are /// all "we don't know", and [`state_needs_write`] turns that into a write. /// /// One failure is louder than the rest, and it is the only one that hides: /// **a non-404 HTTP failure means the read broke while the write still /// works**, so the guard degrades to writing every time and the sweep /// resumes waking the hive with nothing else to show for it. A transport /// error needs no warning of its own — the PUT immediately after it fails /// too, loudly — and a 404 is the expected first-setup case. async fn current_room_state( client: &reqwest::Client, sender_token: &str, url: &str, ) -> Option { let resp = match client.get(url).bearer_auth(sender_token).send().await { Ok(resp) => resp, Err(e) => { tracing::debug!(error = ?e, url, "matrix: state read unreachable; writing"); return None; } }; let status = resp.status(); if !status.is_success() { if status != reqwest::StatusCode::NOT_FOUND { tracing::warn!( %status, url, "matrix: state read failed while writes still work — the \ skip-if-unchanged guard is off and every sweep will re-emit" ); } return None; } match resp.json::().await { Ok(body) => Some(body), Err(e) => { tracing::debug!(error = ?e, url, "matrix: state read body unparseable; writing"); None } } } /// PUT a state event into `room_id` using the sender token, **skipping the /// write when the room already carries identical content**. /// /// The read is not an optimisation. A PUT of identical content is a no-op /// on the room's *state*, and the homeserver still appends an event to the /// *timeline* — so "idempotent" was true one level too high. Downstream, /// an event is unread activity, which is a todo, which is a turn: a caller /// re-applying a link on a periodic sweep wakes every agent in the room on /// that sweep's cadence, forever. Measured at ~30 minutes per wake per /// agent before this guard existed. /// /// The self-healing property the re-apply exists for is unaffected: a /// missing or divergent link still gets written. async fn set_room_state( client: &reqwest::Client, sender_token: &str, room_id: &str, event_type: &str, state_key: &str, content: &serde_json::Value, ) -> Result<()> { let base = matrix_base()?; let encoded_room = encode_room_id_for_url(room_id); let encoded_key = encode_room_id_for_url(state_key); let url = format!("{base}/_matrix/client/v3/rooms/{encoded_room}/state/{event_type}/{encoded_key}"); let current = current_room_state(client, sender_token, &url).await; if !state_needs_write(current.as_ref(), content) { tracing::debug!( %room_id, event_type, "matrix: state already current, skipping PUT (no timeline event)" ); return Ok(()); } let resp = client .put(&url) .bearer_auth(sender_token) .json(content) .send() .await .with_context(|| format!("matrix: PUT state {event_type} into {room_id}"))?; let status = resp.status(); if status.is_success() { return Ok(()); } let body = resp.json::().await.unwrap_or_default(); anyhow::bail!("matrix: set state {event_type} in {room_id}: HTTP {status}, body: {body}") } /// Create (or recover) the default hive chat room and wire it as a child of /// the hive Space, persisting its room id to [`hive_chat_room_id_path()`]. /// /// The room is a normal room (not an `m.space`) named [`HIVE_CHAT_ROOM_NAME`] /// with a `restricted` join rule allowing any member of the hive Space to /// join — so the operator (who is in the Space) and every agent can chat /// without needing an explicit invite. It's linked bidirectionally to the /// Space: `m.space.child` on the Space points at the room, `m.space.parent` /// on the room points back. Joining a Space does NOT auto-join its children /// (Matrix semantics), so this gives clients a concrete room to surface + /// join instead of an empty Space. /// /// Dedup mirrors [`ensure_hive_space`]: persisted id wins, else rediscover /// by name, else create. The space-child link is re-checked on every call /// so a recovered room reconverges its hierarchy link — and written only /// when it differs, see [`set_room_state`]. /// /// # Errors /// /// Returns an error if the homeserver is unreachable, `createRoom` fails, /// or the room-ID file cannot be written. A failure wiring the space-child /// link is logged but not fatal (the room still exists + is joinable). pub async fn ensure_hive_chat_room( client: &reqwest::Client, sender_token: &str, space_room_id: &str, server_name: &str, ) -> Result { let base = matrix_base()?; // 1. Stored room id wins (fast path). 2. Rediscover by name before // creating (prevents duplicates after a state wipe). 3. Create. let room_id = if let Some(id) = std::fs::read_to_string(hive_chat_room_id_path()) .ok() .map(|s| s.trim().to_owned()) .filter(|s| !s.is_empty()) { tracing::debug!(room_id = %id, "matrix: hive chat room already provisioned"); id } else if let Some(id) = find_chat_room_by_name(client, sender_token).await { persist_chat_room_id(&id)?; tracing::info!(room_id = %id, "matrix: recovered hive chat room by name"); id } else { // `initial_state` is applied after the preset-derived state, so the // restricted join rule overrides private_chat's invite-only default. let url = format!("{base}/_matrix/client/v3/createRoom"); let body = serde_json::json!({ "name": HIVE_CHAT_ROOM_NAME, "topic": HIVE_CHAT_ROOM_TOPIC, // Pin the room version explicitly: the `restricted` join rule // below needs room version >= 8. tuwunel's default is currently // higher, but pinning makes the dependency explicit so a future // homeserver-default change can't silently invalidate the // restricted rule (which would quietly fall back to invite-only // and break the "operator joins from the Space" path). "room_version": "10", "preset": "private_chat", "visibility": "private", "initial_state": [ { "type": "m.room.join_rules", "state_key": "", "content": { "join_rule": "restricted", "allow": [ { "type": "m.room_membership", "room_id": space_room_id } ] } }, { "type": "m.space.parent", "state_key": space_room_id, "content": { "via": [server_name], "canonical": true } } ] }); let resp = client .post(&url) .bearer_auth(sender_token) .json(&body) .send() .await .context("matrix: POST /createRoom (hive chat room)")?; let status = resp.status(); let json = resp .json::() .await .context("matrix: parse /createRoom response (chat room)")?; if !status.is_success() { anyhow::bail!("matrix: createRoom (chat) HTTP {status}, body: {json}"); } let id = json["room_id"] .as_str() .with_context(|| format!("matrix: createRoom (chat) missing room_id: {json}"))? .to_owned(); persist_chat_room_id(&id)?; tracing::info!(room_id = %id, "matrix: created hive chat room"); id }; // Wire the Space → room child link. Without an `m.space.child` carrying // a `via`, the room won't surface in the Space hierarchy. `suggested` // hints clients to surface it prominently. let child_content = serde_json::json!({ "via": [server_name], "suggested": true, }); if let Err(e) = set_room_state( client, sender_token, space_room_id, "m.space.child", &room_id, &child_content, ) .await { tracing::warn!(error = ?e, "matrix: set m.space.child on hive space failed"); } Ok(room_id) } /// Invite `@{localpart}:{server_name}` to `room_id` using the admin /// account. Idempotent — treats already-member responses as success. async fn invite_to_room( client: &reqwest::Client, sender_token: &str, room_id: &str, localpart: &str, server_name: &str, ) -> Result<()> { let user_id = format!("@{localpart}:{server_name}"); invite_user_id(client, sender_token, room_id, &user_id).await } /// Fetch a user's current membership in a room via the sender token, or /// `None` if there is no membership event (never invited) or the lookup /// fails. Returns the raw membership string (`invite`, `join`, `leave`, …). async fn room_membership( client: &reqwest::Client, sender_token: &str, encoded_room_id: &str, user_id: &str, ) -> Option { let base = matrix_http()?; // `:` must be percent-encoded in both the room-id and user-id path // segments; `@` and `!` are permitted path characters per RFC 3986. let encoded_user = user_id.replace(':', "%3A"); let url = format!( "{base}/_matrix/client/v3/rooms/{encoded_room_id}/state/m.room.member/{encoded_user}" ); let resp = client .get(&url) .bearer_auth(sender_token) .send() .await .ok()?; if !resp.status().is_success() { // 404 = no membership event yet; anything else we treat as "unknown" // and let the caller fall through to the invite attempt. return None; } let body = resp.json::().await.ok()?; body["membership"].as_str().map(ToOwned::to_owned) } /// 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). 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, room_id: &str, user_id: &str, ) -> Result<()> { let base = matrix_base()?; // `:` must be percent-encoded in the room-id path segment; `!` is // permitted in URL path characters per RFC 3986. let encoded_room_id = room_id.replace(':', "%3A"); // Skip the invite entirely when the user is already invited or joined. // Re-POSTing an invite to a pending member re-sends the invite event, // which re-notifies the agent on every provisioning sweep. if let Some(membership) = room_membership(client, sender_token, &encoded_room_id, user_id).await && matches!(membership.as_str(), "invite" | "join") { tracing::debug!(%user_id, %room_id, %membership, "matrix: invite skipped (already a member/invited)"); return Ok(()); } let url = format!("{base}/_matrix/client/v3/rooms/{encoded_room_id}/invite"); let resp = client .post(&url) .bearer_auth(sender_token) .json(&serde_json::json!({ "user_id": user_id })) .send() .await .with_context(|| format!("matrix: POST /rooms/.../invite for {user_id}"))?; let status = resp.status(); if status.is_success() { tracing::debug!(%user_id, %room_id, "matrix: invited to room"); return Ok(()); } // 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}") } /// Invite an arbitrary Matrix user to the hive Space (default) or an /// explicit `room_id` / alias. `user` may be a fully-qualified id /// (`@name:server`) or a bare localpart, which is qualified with the /// homeserver's `server_name`. Returns the resolved room id the invite /// targeted. Used by `hivectl matrix invite`. /// /// # Errors /// /// Returns an error if the sender token or `server_name` can't be read, /// the target room can't be resolved (no `--room` and no persisted hive /// space), or the invite POST fails for a reason other than the user /// already being a member / invited. pub async fn invite_user( client: &reqwest::Client, sender_token: &str, user: &str, room_override: Option<&str>, server_name: &str, ) -> Result { // Qualify a bare localpart to a full user id on the hive homeserver. let user_id = if user.starts_with('@') { user.to_owned() } else { format!("@{user}:{server_name}") }; // Resolve the room: explicit override (id or #alias) wins; otherwise // the persisted hive Space. let room_id = match room_override { Some(r) if r.starts_with('#') => resolve_room_alias(client, sender_token, r).await?, Some(r) if r.starts_with('!') => r.to_owned(), Some(r) => anyhow::bail!( "matrix: --room {r:?} is neither a room id nor an alias; \ prefix with '!' for a room id (!abc:server) or '#' for an \ alias (#name:server)" ), None => std::fs::read_to_string(hive_space_room_id_path()) .map(|s| s.trim().to_owned()) .context( "matrix: no --room given and no persisted hive space \ (run the hive-c0re matrix sweep first)", )?, }; invite_user_id(client, sender_token, &room_id, &user_id).await?; Ok(room_id) } /// Resolve a `#alias:server` to its room id via the directory API. async fn resolve_room_alias( client: &reqwest::Client, sender_token: &str, alias: &str, ) -> Result { let base = matrix_base()?; let encoded = alias.replace('#', "%23").replace(':', "%3A"); let url = format!("{base}/_matrix/client/v3/directory/room/{encoded}"); let resp = client .get(&url) .bearer_auth(sender_token) .send() .await .with_context(|| format!("matrix: GET directory for {alias}"))?; let status = resp.status(); let json = resp.json::().await.unwrap_or_default(); if !status.is_success() { anyhow::bail!("matrix: resolve alias {alias}: HTTP {status}, body: {json}"); } json["room_id"] .as_str() .map(str::to_owned) .ok_or_else(|| anyhow::anyhow!("matrix: alias {alias} response missing room_id: {json}")) } /// Make sure the hive's own `@hive-:` account exists, then provision /// the hive Space and chat room and invite every agent container to both. /// Agents' own accounts are not created here: `swarm-controller` mints them /// with the swarm's appservice token. Runs only as a `MatrixSweep` job node /// (see `job_queue::templates::matrix_sweep`). No-op when the /// hive-matrix container isn't running. Per-step failures are logged /// but don't abort the sweep. /// /// 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() { tracing::debug!("matrix: hive-matrix container absent, skipping user sweep"); return true; } let mut ok = true; // Loud and non-destructive: with no appservice token this sweep can // create nothing, so it does nothing. let as_token = match read_appservice_token() { Ok(t) => t, Err(e) => { tracing::warn!(error = ?e, "matrix: no appservice token; skipping the user sweep"); return false; } }; // One HTTP client for the whole sweep. let client = match reqwest::Client::builder() .timeout(std::time::Duration::from_secs(HTTP_TIMEOUT_SECS)) .build() { Ok(c) => c, Err(e) => { tracing::warn!(error = ?e, "matrix: build HTTP client failed; skipping sweep"); return false; } }; // The hive's `@hive-:` account first, because everything below provisions // THROUGH it (the Space, the chat room and every invite are sent with // its token) — as an ordinary user that created those rooms, not as a // homeserver admin. if let Err(e) = ensure_hive_user(&client, &as_token).await { tracing::warn!(error = ?e, "matrix: ensure_hive_user failed"); ok = false; } let Ok(containers) = crate::lifecycle::list().await else { tracing::warn!("matrix: nixos-container list failed; skipping user sweep"); return false; }; let mut agent_names: Vec = Vec::new(); for c in &containers { let Some(name) = c.strip_prefix(crate::lifecycle::AGENT_PREFIX) else { continue; }; agent_names.push(name.to_owned()); } if !provision_space(&client, &agent_names).await { 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. let sender_token = match read_sender_token() { Ok(t) => t, Err(e) => { tracing::warn!(error = ?e, "matrix: skipping hive space provisioning (no matrix sender token)"); return false; } }; // server_name first — the agent invites need it (fully-qualified user ids). let server_name = match discover_server_name(client).await { Ok(s) => s, Err(e) => { tracing::warn!(error = ?e, "matrix: discover_server_name failed; skipping space provisioning"); return false; } }; let room_id = match ensure_hive_space(client, &sender_token).await { Ok(id) => id, Err(e) => { tracing::warn!(error = ?e, "matrix: ensure_hive_space failed"); return false; } }; // This hive's own account first, then all agents. let localpart = match hive_localpart() { Ok(l) => l, Err(e) => { tracing::warn!(error = ?e, "matrix: skipping hive space provisioning"); return false; } }; if let Err(e) = invite_to_room(client, &sender_token, &room_id, &localpart, &server_name).await { tracing::warn!(error = ?e, %localpart, "matrix: invite the hive account to the space failed"); ok = false; } for name in agent_names { if let Err(e) = invite_to_room(client, &sender_token, &room_id, name, &server_name).await { tracing::warn!(%name, error = ?e, "matrix: invite agent to space failed"); ok = false; } } // Provision the default hive chat room as an m.space.child of the Space // and invite `@hive-:` + every agent. Joining a Space alone surfaces no // rooms to chat in (Matrix semantics — children aren't auto-joined), so // without this the Space is empty. The restricted join rule additionally // lets the operator (a Space member) join from the Space hierarchy. match ensure_hive_chat_room(client, &sender_token, &room_id, &server_name).await { Ok(chat_room_id) => { if let Err(e) = invite_to_room( client, &sender_token, &chat_room_id, &localpart, &server_name, ) .await { tracing::warn!(error = ?e, %localpart, "matrix: invite the hive account to the chat room failed"); ok = false; } for name in agent_names { if let Err(e) = invite_to_room(client, &sender_token, &chat_room_id, name, &server_name).await { tracing::warn!(%name, error = ?e, "matrix: invite agent to chat room failed"); ok = false; } } } Err(e) => { tracing::warn!(error = ?e, "matrix: ensure_hive_chat_room failed"); ok = false; } } ok } #[cfg(test)] mod tests { use super::*; #[test] fn random_hex_is_well_formed_and_correct_length() { let h = random_hex(16).expect("/dev/urandom readable"); assert_eq!(h.len(), 32); assert!(h.chars().all(|c| c.is_ascii_hexdigit())); } /// The steady state. This is the whole point of the guard: the sweep /// runs forever, and every write it makes is a wake for every agent /// in the room. #[test] fn state_matching_the_room_is_not_rewritten() { let desired = serde_json::json!({ "via": ["example.org"], "suggested": true }); let current = desired.clone(); assert!(!state_needs_write(Some(¤t), &desired)); } /// Key order in the homeserver's response must not force a write — /// the guard leans on `serde_json::Value` comparing objects as maps, /// and a byte- or order-sensitive comparison would silently degrade /// to writing every time while still looking correct. #[test] fn state_matching_but_reordered_is_not_rewritten() { let desired = serde_json::json!({ "via": ["example.org"], "suggested": true }); let current = serde_json::json!({ "suggested": true, "via": ["example.org"] }); assert!(!state_needs_write(Some(¤t), &desired)); } #[test] fn diverged_state_is_rewritten() { let desired = serde_json::json!({ "via": ["example.org"], "suggested": true }); let current = serde_json::json!({ "via": ["old.example.org"], "suggested": true }); assert!(state_needs_write(Some(¤t), &desired)); } /// Fail-open, and deliberately so: an unreadable current state must /// write. The re-apply exists to repair a missing link, so "we don't /// know" has to behave like "it's missing", not like "it's fine". #[test] fn unknown_state_is_written() { let desired = serde_json::json!({ "via": ["example.org"], "suggested": true }); assert!(state_needs_write(None, &desired)); } #[test] fn random_hex_two_calls_differ() { // Sanity check — not a statistical claim, just guards // against ever accidentally returning a constant. let a = random_hex(16).expect("/dev/urandom readable"); let b = random_hex(16).expect("/dev/urandom readable"); assert_ne!(a, b); } #[test] fn extract_access_token_pulls_from_success_body() { let body = serde_json::json!({ "user_id": "@alice:matrix.example.org", "access_token": "syt_abc123", "device_id": "ABC", }); 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"); } #[test] fn extract_access_token_errors_on_missing_field() { let body = serde_json::json!({"user_id": "@alice:matrix.example.org"}); let err = extract_access_token(&body).unwrap_err(); assert!(err.to_string().contains("missing access_token")); } }