hive-matrix-daemon now learns which external matrix accounts it has from the swarm secret store, under the agent's own certificate, and the hive push chain for matrix is gone. The daemon lists swarm/agents/<agent>/matrix/ (the `list` its policy grants on its own metadata subtree), reads each account's homeserver from its credential, and brings the accounts up with their tokens from the store. Every two minutes it lists again and exits with 75 when the set of linked accounts changed; the unit restarts on 75 without counting a failure. A listed name whose credential reads as absent is skipped and logged once. At start it removes the matrix-token-<a> / matrix-account-<a>.json pairs a hive delivered (a sidecar marks a pair as delivered; a declared tokenFile keeps its token). Removed: CredentialNotice and the $SWARM.credential.* subject and NATS grant, the controller's publish and its queue precondition on the PUT route, hive-c0re's credential subscription arm and workers/credential.rs, priv_client::write_agent_matrix_token, hive-priv's WriteAgentMatrixToken and its helpers, and the daemon's state-dir account discovery. Kept: WriteAgentGithubToken and the external-forge path (WriteAgentExtraForgeAccount, extra_forges.rs) are untouched, and a declared matrixAccounts tokenFile is still read when the store has no token for that account. Refs #4348
484 lines
21 KiB
Rust
484 lines
21 KiB
Rust
//! `hive-matrix-daemon` binary — long-running matrix-sdk Client + sync
|
|
//! loop per matrix account. Bridges incoming room events to hyperhive
|
|
//! wake signals and serves its MCP tools directly over streamable-http
|
|
//! on `--http <addr>` — no stdio bridge, no separate bin claude has to
|
|
//! respawn every turn.
|
|
//!
|
|
//! Lifecycle:
|
|
//! 1. Read the account list: `accounts::configured()` (declared) plus
|
|
//! `accounts::discover_linked` (linked in the swarm secret store).
|
|
//! 2. For each account: resolve its access token (`account_token` — the
|
|
//! swarm secret store first, read by this agent as itself, then a declared
|
|
//! token file) → whoami probe → recover `user_id` + `device_id`
|
|
//! → restore matrix-sdk session (no login flow), install the
|
|
//! message-event handler, and spawn its own sync loop.
|
|
//! 3. Serve the MCP tools against an account→Client registry; each tool call
|
|
//! routes to its `account` arg (the primary account when omitted).
|
|
//! 4. Every [`LINKED_REFRESH`], re-read the linked accounts; exit with
|
|
//! [`ACCOUNTS_CHANGED_EXIT`] when the set changed, so systemd restarts us.
|
|
//!
|
|
//! Standalone-degraded boot: the PRIMARY account having no token →
|
|
//! exit 0 cleanly; systemd restarts us once one exists (the path-watcher for
|
|
//! a token file, the store re-check timer for a stored token). A SECONDARY
|
|
//! account missing its token is skipped (the daemon still serves the others).
|
|
//!
|
|
//! Stale-token recovery (`M_UNKNOWN_TOKEN`): handled in
|
|
//! `client::build_and_restore`, and it mirrors the missing-token policy
|
|
//! above — a rejected PRIMARY token drops the stale token + sdk state and
|
|
//! exits 0 until a fresh token exists, while a rejected SECONDARY token
|
|
//! is removed and that one account is skipped so the daemon keeps serving
|
|
//! the primary and any other healthy account.
|
|
|
|
use std::collections::{BTreeSet, HashSet};
|
|
use std::process::ExitCode;
|
|
use std::sync::Arc;
|
|
|
|
use anyhow::{Context, Result};
|
|
use clap::Parser;
|
|
use matrix_sdk::{Client, config::SyncSettings};
|
|
|
|
use hive_matrix_mcp::accounts::{AccountCfg, Registry};
|
|
use hive_matrix_mcp::client::PermanentBringUpError;
|
|
use hive_matrix_mcp::credential::Token;
|
|
use hive_matrix_mcp::{accounts, client, credential, mcp, paths, timeline, wake};
|
|
|
|
#[derive(Parser)]
|
|
#[command(name = "hive-matrix-daemon", about = "matrix-sdk client + MCP daemon")]
|
|
struct Cli {
|
|
/// Serve the MCP tools over streamable-http on this address (e.g.
|
|
/// `127.0.0.1:8792`). Bind loopback only.
|
|
#[arg(long)]
|
|
http: std::net::SocketAddr,
|
|
}
|
|
|
|
/// A per-account sync loop, boxed so loops for N accounts can be driven
|
|
/// concurrently on the main task. Deliberately NOT `Send`: matrix-sdk's
|
|
/// sync future isn't `Send`, so these run on the current task (via
|
|
/// `select_all`) rather than `tokio::spawn`/`JoinSet`.
|
|
type SyncLoop = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>>>>;
|
|
|
|
/// How often the daemon rewrites the accounts snapshot to advance its
|
|
/// mtime, so the dashboard's `as_of` tracks daemon liveness. 30s keeps
|
|
/// the staleness window small while the rewrite cost (one tmp + rename of
|
|
/// a tiny file) is negligible.
|
|
const ACCOUNTS_HEARTBEAT_SECS: u64 = 30;
|
|
|
|
/// Retry delays (seconds) for transient secondary-account failures.
|
|
/// After the initial attempt we back off through these before giving up
|
|
/// and skipping the account for the rest of the daemon lifetime.
|
|
/// Transient = anything that is NOT a `PermanentBringUpError` (bad/expired
|
|
/// token). A down homeserver at DNS-not-ready boot time is the typical
|
|
/// case; the total wait is ~52s before we give up.
|
|
const SECONDARY_RETRY_DELAYS_SECS: &[u64] = &[2, 5, 15, 30];
|
|
|
|
/// How often the daemon re-reads the store's linked accounts. A newly linked
|
|
/// account comes up within this long of the swarm UI storing it.
|
|
const LINKED_REFRESH: std::time::Duration = std::time::Duration::from_mins(2);
|
|
|
|
/// The exit status the daemon leaves with when the store's linked accounts
|
|
/// changed. `nix/agent-modules/matrix.nix` names the same number in the unit's
|
|
/// `RestartForceExitStatus` and `SuccessExitStatus`, so systemd restarts it
|
|
/// without counting a failure.
|
|
const ACCOUNTS_CHANGED_EXIT: u8 = 75;
|
|
|
|
#[tokio::main]
|
|
async fn main() -> Result<ExitCode> {
|
|
tracing_subscriber::fmt()
|
|
.with_env_filter(
|
|
tracing_subscriber::EnvFilter::try_from_default_env()
|
|
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")),
|
|
)
|
|
.with_writer(std::io::stderr)
|
|
// This is a systemd-managed daemon — stderr always goes to journald,
|
|
// never a human terminal, and journald doesn't strip ANSI escapes:
|
|
// they land in victorialogs as raw byte-array spam otherwise.
|
|
.with_ansi(false)
|
|
.init();
|
|
|
|
let cli = Cli::parse();
|
|
let declared = accounts::configured().context("read matrix account config")?;
|
|
accounts::remove_delivered_files(&declared);
|
|
let mut warned = HashSet::new();
|
|
let linked = linked_now(&declared, &mut warned).await.unwrap_or_default();
|
|
let booted = account_names(&declared, &linked);
|
|
let mut cfgs = declared.clone();
|
|
cfgs.extend(linked);
|
|
let multi = cfgs.len() > 1;
|
|
let primary = cfgs[0].name.clone();
|
|
let mut registry = Registry::new(primary);
|
|
let mut sync_loops: Vec<SyncLoop> = Vec::new();
|
|
|
|
for (idx, cfg) in cfgs.into_iter().enumerate() {
|
|
let is_primary = idx == 0;
|
|
// Account-tag the todos only in multi-account mode so single-
|
|
// account todo summaries stay byte-identical to the legacy format.
|
|
let tag = multi.then(|| cfg.name.clone());
|
|
match bring_up_account(&cfg, tag.clone(), is_primary).await {
|
|
Ok(Some((client, sync_loop))) => {
|
|
registry.insert(cfg.name, client);
|
|
sync_loops.push(sync_loop);
|
|
}
|
|
// No token yet: for the primary that means the daemon isn't
|
|
// useful — exit 0, and systemd restarts us once a token exists
|
|
// (the path-watcher for a file, the re-check timer for the store).
|
|
Ok(None) if is_primary => {
|
|
tracing::warn!(
|
|
account = %cfg.name,
|
|
"primary matrix account has no token yet; exiting cleanly \
|
|
(systemd restarts us once one exists)"
|
|
);
|
|
return Ok(ExitCode::SUCCESS);
|
|
}
|
|
Ok(None) => {
|
|
tracing::warn!(account = %cfg.name, "secondary matrix account has no token; skipping");
|
|
}
|
|
// The primary failing to restore is fatal (propagate so
|
|
// systemd retries on a transient blip — matches legacy
|
|
// behaviour); a secondary failing is retried with backoff
|
|
// before being skipped for this daemon lifetime.
|
|
Err(e) if is_primary => return Err(e.context("bring up primary matrix account")),
|
|
Err(e) => {
|
|
if let Some((client, sync_loop)) = bring_up_secondary_with_retry(&cfg, tag, e).await
|
|
{
|
|
registry.insert(cfg.name, client);
|
|
sync_loops.push(sync_loop);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
if registry.is_empty() {
|
|
tracing::warn!("no matrix accounts restored; exiting cleanly");
|
|
return Ok(ExitCode::SUCCESS);
|
|
}
|
|
|
|
// Publish the live-account snapshot (BE-4) for the dashboard: the
|
|
// accounts that restored, with effective homeserver + user id. Lands
|
|
// in the host-visible state dir so hive-c0re reads it to show live
|
|
// up/down + backfill homeserver. Best-effort — a failed write must
|
|
// not stop the daemon from serving.
|
|
if let Err(e) = registry.write_snapshot(&paths::accounts_file()) {
|
|
tracing::warn!(error = %format!("{e:#}"), "failed to write matrix-accounts snapshot");
|
|
}
|
|
|
|
// Serve the MCP tools against the registry. Spawned before driving
|
|
// the sync loops so claude can reach the stable http URL as soon as
|
|
// the first turn fires.
|
|
let registry = Arc::new(registry);
|
|
|
|
// Heartbeat: periodically rewrite the accounts snapshot so its mtime
|
|
// advances while the daemon lives. The dashboard derives `as_of` from
|
|
// the file mtime, so a stalled mtime now means the daemon is down —
|
|
// which lets the dashboard dim accounts whose snapshot is stale rather
|
|
// than reporting the boot-time set forever (BE-4 follow-up).
|
|
let hb_registry = Arc::clone(®istry);
|
|
tokio::spawn(async move {
|
|
let mut tick =
|
|
tokio::time::interval(std::time::Duration::from_secs(ACCOUNTS_HEARTBEAT_SECS));
|
|
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
|
let path = paths::accounts_file();
|
|
loop {
|
|
tick.tick().await;
|
|
if let Err(e) = hb_registry.heartbeat_snapshot(&path) {
|
|
tracing::warn!(error = %format!("{e:#}"), "matrix-accounts heartbeat write failed");
|
|
}
|
|
}
|
|
});
|
|
|
|
let http_addr = cli.http;
|
|
let mcp_registry = Arc::clone(®istry);
|
|
tokio::spawn(async move {
|
|
if let Err(e) = mcp::serve_http(http_addr, mcp_registry).await {
|
|
tracing::error!(error = %e, "mcp http server exited");
|
|
}
|
|
});
|
|
|
|
Ok(drive(sync_loops, wait_for_linked_change(declared, booted, warned)).await)
|
|
}
|
|
|
|
/// Drive all per-account sync loops concurrently on this task (they aren't
|
|
/// `Send`, so no `tokio::spawn`) until one returns or `changed` resolves.
|
|
/// matrix-sdk reconnects internally, so any loop returning is exceptional —
|
|
/// log it and exit so systemd restarts the whole daemon cleanly.
|
|
async fn drive(
|
|
sync_loops: Vec<SyncLoop>,
|
|
changed: impl std::future::Future<Output = ()>,
|
|
) -> ExitCode {
|
|
tokio::select! {
|
|
(result, idx, _rest) = futures_util::future::select_all(sync_loops) => {
|
|
match result {
|
|
Ok(()) => tracing::warn!(
|
|
account_index = idx,
|
|
"a matrix sync loop exited cleanly; restarting daemon"
|
|
),
|
|
Err(e) => {
|
|
tracing::error!(account_index = idx, error = %format!("{e:#}"), "a matrix sync loop errored; restarting daemon");
|
|
}
|
|
}
|
|
ExitCode::SUCCESS
|
|
}
|
|
() = changed => {
|
|
tracing::info!("the store's linked matrix accounts changed; exiting to be restarted onto them");
|
|
ExitCode::from(ACCOUNTS_CHANGED_EXIT)
|
|
}
|
|
}
|
|
}
|
|
|
|
/// The accounts the store links to this agent beyond `declared`, or `None`
|
|
/// when the store could not be read.
|
|
///
|
|
/// Each linked name that cannot be brought up is logged the first time it is
|
|
/// seen and recorded in `warned`, so a refresh does not repeat it.
|
|
async fn linked_now(
|
|
declared: &[AccountCfg],
|
|
warned: &mut HashSet<String>,
|
|
) -> Option<Vec<AccountCfg>> {
|
|
let found = match accounts::discover_linked(declared).await {
|
|
Ok(found) => found,
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
error = %format!("{e:#}"),
|
|
"could not read this agent's linked matrix accounts from the store"
|
|
);
|
|
return None;
|
|
}
|
|
};
|
|
for name in found.missing {
|
|
if warned.insert(name.clone()) {
|
|
tracing::warn!(
|
|
account = %name,
|
|
"the store lists this matrix account but holds no credential for it; skipping"
|
|
);
|
|
}
|
|
}
|
|
for name in found.no_homeserver {
|
|
if warned.insert(name.clone()) {
|
|
tracing::warn!(
|
|
account = %name,
|
|
"this linked matrix account was stored without a homeserver; skipping \
|
|
(re-link it from the swarm UI with one)"
|
|
);
|
|
}
|
|
}
|
|
Some(found.accounts)
|
|
}
|
|
|
|
/// Every account name the daemon would bring up from `declared` + `linked`.
|
|
fn account_names(declared: &[AccountCfg], linked: &[AccountCfg]) -> BTreeSet<String> {
|
|
declared
|
|
.iter()
|
|
.chain(linked)
|
|
.map(|a| a.name.clone())
|
|
.collect()
|
|
}
|
|
|
|
/// Re-read the store's linked accounts every [`LINKED_REFRESH`] and resolve
|
|
/// once they name a different set than `booted`. A failed read is logged and
|
|
/// skipped rather than counted as a change.
|
|
async fn wait_for_linked_change(
|
|
declared: Vec<AccountCfg>,
|
|
booted: BTreeSet<String>,
|
|
mut warned: HashSet<String>,
|
|
) {
|
|
let mut tick = tokio::time::interval(LINKED_REFRESH);
|
|
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
|
|
// The first tick is immediate, and boot has just read the store.
|
|
tick.tick().await;
|
|
loop {
|
|
tick.tick().await;
|
|
let Some(linked) = linked_now(&declared, &mut warned).await else {
|
|
continue;
|
|
};
|
|
if account_names(&declared, &linked) != booted {
|
|
return;
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Try to bring up a secondary account, retrying with exponential backoff
|
|
/// on transient failures (anything that is NOT a `PermanentBringUpError`).
|
|
/// Returns `Some((client, sync_loop))` on success, or `None` to signal
|
|
/// that the account should be skipped for this daemon lifetime (permanent
|
|
/// failure, no token, or all retries exhausted).
|
|
async fn bring_up_secondary_with_retry(
|
|
cfg: &AccountCfg,
|
|
tag: Option<String>,
|
|
first_error: anyhow::Error,
|
|
) -> Option<(Client, SyncLoop)> {
|
|
// Permanent failure: the token was invalid/expired and has already been
|
|
// removed from disk. Retrying won't help — skip immediately.
|
|
if first_error
|
|
.downcast_ref::<PermanentBringUpError>()
|
|
.is_some()
|
|
{
|
|
tracing::error!(
|
|
account = %cfg.name,
|
|
error = %format!("{first_error:#}"),
|
|
"secondary matrix account: permanent failure; skipping"
|
|
);
|
|
return None;
|
|
}
|
|
// Transient failure (network/DNS/homeserver 5xx): retry with backoff.
|
|
tracing::warn!(
|
|
account = %cfg.name,
|
|
error = %format!("{first_error:#}"),
|
|
"secondary matrix account bring-up failed (transient); will retry"
|
|
);
|
|
for &delay in SECONDARY_RETRY_DELAYS_SECS {
|
|
tracing::info!(
|
|
account = %cfg.name,
|
|
delay_s = delay,
|
|
"retrying secondary account bring-up after backoff"
|
|
);
|
|
tokio::time::sleep(std::time::Duration::from_secs(delay)).await;
|
|
match bring_up_account(cfg, tag.clone(), false).await {
|
|
Ok(Some((client, sync_loop))) => {
|
|
tracing::info!(account = %cfg.name, "secondary matrix account recovered");
|
|
return Some((client, sync_loop));
|
|
}
|
|
Ok(None) => {
|
|
tracing::warn!(account = %cfg.name, "secondary matrix account has no token; skipping");
|
|
return None;
|
|
}
|
|
Err(re) if re.downcast_ref::<PermanentBringUpError>().is_some() => {
|
|
tracing::error!(
|
|
account = %cfg.name,
|
|
error = %format!("{re:#}"),
|
|
"secondary matrix account: permanent failure on retry; skipping"
|
|
);
|
|
return None;
|
|
}
|
|
Err(re) => {
|
|
tracing::warn!(
|
|
account = %cfg.name,
|
|
error = %format!("{re:#}"),
|
|
"secondary matrix account bring-up still failing (transient)"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
tracing::error!(
|
|
account = %cfg.name,
|
|
retries = SECONDARY_RETRY_DELAYS_SECS.len(),
|
|
"secondary matrix account failed after all retries; skipping for this daemon lifetime"
|
|
);
|
|
None
|
|
}
|
|
|
|
/// Resolve `cfg`'s access token, preferring the copy this agent can fetch
|
|
/// itself.
|
|
///
|
|
/// 🏛️ The store read is the point: the credential at
|
|
/// `swarm/agents/<agent>/matrix/<account>` is **this agent's**, and it is read
|
|
/// from inside this agent's container under this agent's own certificate,
|
|
/// never by the hive on the agent's behalf. It is held in memory from here to
|
|
/// `restore_session` and is written nowhere.
|
|
///
|
|
/// `main` included: `swarm-controller` mints it with the swarm's appservice
|
|
/// token and stores it at `swarm/agents/<agent>/matrix/main`, and no hive
|
|
/// mints it any more. An agent whose harness forwards no
|
|
/// [`credential::ENV_AGENT`] cannot name its own subtree, so it reads the file.
|
|
///
|
|
/// The file is what remains when the store answers nothing: a `tokenFile` an
|
|
/// operator declared, or a `main` token a hive minted before the swarm did. An
|
|
/// account the store links has no file. A store read that *fails* falls back
|
|
/// too, and says why.
|
|
///
|
|
/// # Errors
|
|
/// As [`credential::Token::from_file`]: a token file that exists but cannot be
|
|
/// read, or is empty.
|
|
async fn account_token(cfg: &AccountCfg) -> Result<Option<Token>> {
|
|
if let Some(agent) = credential::agent_name() {
|
|
match Token::from_store(&agent, &cfg.name).await {
|
|
Ok(Some(token)) => return Ok(Some(token)),
|
|
// No store in this deployment: the declared file is the whole
|
|
// mechanism here.
|
|
Ok(None) => {}
|
|
Err(e) => tracing::warn!(
|
|
account = %cfg.name,
|
|
error = %format!("{e:#}"),
|
|
"could not read this account's credential from the store as this agent; \
|
|
falling back to its declared token file"
|
|
),
|
|
}
|
|
}
|
|
match &cfg.token_file {
|
|
Some(path) => Token::from_file(path).await,
|
|
None => Ok(None),
|
|
}
|
|
}
|
|
|
|
/// Restore one account's client (when its token exists), install its
|
|
/// message handler, and build its sync loop. Returns
|
|
/// `Ok(Some((client, sync_loop)))` when the account came up, `Ok(None)`
|
|
/// when its token file is absent (caller decides primary-vs-secondary
|
|
/// handling). The returned sync loop is driven by the caller.
|
|
async fn bring_up_account(
|
|
cfg: &AccountCfg,
|
|
tag: Option<String>,
|
|
is_primary: bool,
|
|
) -> Result<Option<(Client, SyncLoop)>> {
|
|
let Some(homeserver) = cfg.homeserver() else {
|
|
// No homeserver for this account: the hive has none to offer (no
|
|
// matrix vhost) or this agent's `services.hyperhive.agent.matrix.url` is null. Same
|
|
// no-op as a missing token — an absent integration, not a guess at
|
|
// one.
|
|
tracing::info!(
|
|
account = %cfg.name,
|
|
"no homeserver configured (HIVE_MATRIX_URL unset); skipping account"
|
|
);
|
|
return Ok(None);
|
|
};
|
|
let Some(token) = account_token(cfg).await? else {
|
|
return Ok(None);
|
|
};
|
|
tracing::info!(
|
|
account = %cfg.name,
|
|
homeserver,
|
|
credential = %token.origin(),
|
|
state_dir = %cfg.state_dir.display(),
|
|
"bringing up matrix account"
|
|
);
|
|
let client = client::build_and_restore(&homeserver, &token, &cfg.state_dir, is_primary)
|
|
.await
|
|
.with_context(|| format!("build matrix client for account {}", cfg.name))?;
|
|
// Best-effort: sync the agent icon to this account's matrix avatar over
|
|
// the live (authenticated, correct-homeserver) Client. Replaces the old
|
|
// curl oneshot; failures are swallowed inside sync_avatar.
|
|
client::sync_avatar(&client, &cfg.state_dir, &cfg.name).await;
|
|
// Best-effort: bootstrap cross-signing so the account's device isn't
|
|
// flagged "unverified" in other users' clients. Idempotent + swallows
|
|
// failures (a UIAA-requiring homeserver) inside ensure_cross_signing.
|
|
client::ensure_cross_signing(&client, &cfg.name).await;
|
|
|
|
let sync_client = client.clone();
|
|
let cb_client = client.clone();
|
|
// Separate dedup sets: one tracks invites already pushed as todos, the
|
|
// other unread-message rooms. Both are pruned to their current state
|
|
// each sweep (see the sweep fns) so re-invites / new messages re-push.
|
|
let invite_notified = Arc::new(tokio::sync::Mutex::new(std::collections::HashSet::new()));
|
|
let unread_notified = Arc::new(tokio::sync::Mutex::new(std::collections::HashSet::new()));
|
|
// Startup cancel-and-recreate (loose-ends v2): wipe this agent's
|
|
// matrix todos so stale ones (rooms read / invites resolved while the
|
|
// daemon was down) don't linger, then let the first sweep rebuild the
|
|
// set to match current reality. Best-effort; the sweep converges.
|
|
let _ = wake::send_todo_clear(None, true).await;
|
|
let sync_loop: SyncLoop = Box::pin(async move {
|
|
sync_client
|
|
.sync_with_callback(SyncSettings::default(), move |_response| {
|
|
let client = cb_client.clone();
|
|
let invite_notified = invite_notified.clone();
|
|
let unread_notified = unread_notified.clone();
|
|
let tag = tag.clone();
|
|
async move {
|
|
timeline::sweep_invites(&client, &invite_notified, tag.as_deref()).await;
|
|
timeline::sweep_unread(&client, &unread_notified, tag.as_deref()).await;
|
|
matrix_sdk::LoopCtrl::Continue
|
|
}
|
|
})
|
|
.await
|
|
.context("matrix-sdk sync loop exited")?;
|
|
Ok(())
|
|
});
|
|
Ok(Some((client, sync_loop)))
|
|
}
|