Watch
0
0
Fork
You've already forked hyperhive
0

matrix: the agent's daemon pulls its linked accounts from bao itself

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
This commit is contained in:
atlas 2026-10-01 17:41:20 +02:00
commit 97fb76ce99
22 changed files with 553 additions and 813 deletions

View file

@ -5,16 +5,17 @@
//! respawn every turn.
//!
//! Lifecycle:
//! 1. Read the configured account list (`accounts::configured()` —
//! `HIVE_MATRIX_ACCOUNTS` JSON, or the single legacy account).
//! 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 the
//! file its hive delivered) → whoami probe → recover `user_id` + `device_id`
//! 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 the account named in its `account` arg (the
//! primary account when omitted).
//! 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
@ -28,6 +29,8 @@
//! 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};
@ -68,8 +71,18 @@ const ACCOUNTS_HEARTBEAT_SECS: u64 = 30;
/// 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<()> {
async fn main() -> Result<ExitCode> {
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
@ -83,7 +96,13 @@ async fn main() -> Result<()> {
.init();
let cli = Cli::parse();
let cfgs = accounts::configured().context("read matrix account config")?;
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);
@ -108,7 +127,7 @@ async fn main() -> Result<()> {
"primary matrix account has no token yet; exiting cleanly \
(systemd restarts us once one exists)"
);
return Ok(());
return Ok(ExitCode::SUCCESS);
}
Ok(None) => {
tracing::warn!(account = %cfg.name, "secondary matrix account has no token; skipping");
@ -130,7 +149,7 @@ async fn main() -> Result<()> {
if registry.is_empty() {
tracing::warn!("no matrix accounts restored; exiting cleanly");
return Ok(());
return Ok(ExitCode::SUCCESS);
}
// Publish the live-account snapshot (BE-4) for the dashboard: the
@ -174,22 +193,106 @@ async fn main() -> Result<()> {
}
});
// Drive all per-account sync loops concurrently on this task (they
// aren't `Send`, so no `tokio::spawn`). matrix-sdk reconnects
// internally, so any loop returning is exceptional — log it and exit
// so systemd restarts the whole daemon cleanly.
let (result, idx, _rest) = futures_util::future::select_all(sync_loops).await;
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");
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)
}
}
}
Ok(())
/// 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
@ -276,10 +379,10 @@ async fn bring_up_secondary_with_retry(
/// 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: an extra account
/// the hive-side delivery loop still writes, or a `main` token a hive minted
/// before the swarm did. A store read that *fails* falls back too, and says
/// why.
/// 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
@ -288,18 +391,21 @@ 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: nothing to report, the hive-side
// delivery is the whole mechanism here.
// 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 the token the hive delivered"
falling back to its declared token file"
),
}
}
Token::from_file(&cfg.token_file).await
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