Watch
0
0
Fork
You've already forked hyperhive
0

swarm-controller: sweep agents' legacy hyperhive-* forge tokens at start

Before b5d07d4d, hive-c0re minted a new `hyperhive-<unix-seconds>` token
for an agent on every spawn and rebuild and never revoked one, so every
live agent's forge user carries a pile of write-scoped tokens nothing
holds. Nothing in the tree lists or deletes them.

On each start, swarm-controller now walks the store's hive-agent-*
roster (the one the swarm-agent mint pass walks), and for every agent
whose swarm-agent token that pass would keep, deletes each token named
exactly `hyperhive-<digits>`. It logs the count per agent and a total.

- `core` is refused by name in both the roster filter and the per-user
  delete: hive-c0re still names core's live admin token
  `hyperhive-<unix-seconds>`.
- An agent whose swarm-agent token is not current is skipped, because
  consumers fall back to `<state>/forge-token`, the last hyperhive-*
  token, until the swarm token is fetched.
- A failed list or delete is logged and skipped; the sweep does not
  retry and never blocks startup. A second start deletes nothing.

Tokens on forge users of agents no longer on the store roster (already
destroyed) are not reached.

Closes #4644
This commit is contained in:
atlas 2026-09-26 15:44:38 +02:00 • committed by mara
commit 23e0c313b8
4 changed files with 368 additions and 5 deletions

View file

@ -34,6 +34,7 @@ use utoipa::ToSchema;
use crate::webhook::DeliveryKind;
pub mod agent_token;
pub mod legacy_tokens;
pub mod objects;
pub mod site_admin;

View file

@ -27,7 +27,7 @@ use super::Client;
/// The forge's name for every agent token this module mints.
///
/// Deliberately not `hyperhive-…`: that prefix is what `hive-c0re` named each
/// of its re-mints, and a later sweep of those must never match this one.
/// of its re-mints, and [`super::legacy_tokens`] deletes those.
pub const AGENT_TOKEN_NAME: &str = "swarm-agent";
/// The scopes an agent token carries. Byte-identical to the `TOKEN_SCOPES`
@ -167,7 +167,7 @@ pub(super) fn is_not_found(e: &ForgejoError) -> bool {
impl Client {
/// Every access token the forge lists for `agent`, or `None` when the
/// forge has no such user.
async fn list_agent_tokens(&self, agent: &str) -> Result<Option<Vec<AccessToken>>> {
pub(super) async fn list_agent_tokens(&self, agent: &str) -> Result<Option<Vec<AccessToken>>> {
match self.api.admin_list_user_access_tokens(agent).all().await {
Ok(tokens) => Ok(Some(tokens)),
Err(e) if is_not_found(&e) => Ok(None),
@ -270,7 +270,7 @@ impl Client {
}
/// What the forge and the store say about `agent`'s token.
async fn observe_agent(
pub(super) async fn observe_agent(
&self,
store: &swarm_secret_client::SecretStore,
agent: &str,

View file

@ -0,0 +1,361 @@
//! A startup sweep of the `hyperhive-<unix-seconds>` access tokens `hive-c0re`
//! minted for agents on every spawn and rebuild, each one live on the forge
//! with write scopes.
//!
//! Two filters decide what goes, and both must pass:
//!
//! - **Whose tokens.** Only agents on the store's `hive-agent-*` roster, the
//! one [`super::agent_token`]'s pass walks, and only those whose
//! [`AGENT_TOKEN_NAME`] token that pass would [`Decision::Keep`]: until the
//! replacement is stored and live, the agent may still be reading its last
//! `hyperhive-*` token from `<state>/forge-token`. [`CORE_USER`] is refused
//! by name at every step. `hive-c0re` still names `core`'s live admin token
//! `hyperhive-<unix-seconds>`, so nothing about the token tells them apart.
//! - **Which tokens.** A name that is exactly `hyperhive-` followed by one or
//! more ASCII digits ([`is_legacy_token_name`]). [`AGENT_TOKEN_NAME`] cannot
//! match.
//!
//! One pass per daemon start. A failed read or delete is logged and skipped;
//! the next start picks up whatever is left, and a start with nothing left
//! deletes nothing.
use std::sync::Arc;
use anyhow::{Context, Result};
use swarm_secret_client::{client::DEFAULT_CERT_MOUNT, policy};
use super::Client;
use super::agent_token::{AGENT_TOKEN_NAME, Decision, Observed};
/// The forge user `hive-c0re` runs as. Its admin token carries a legacy-shaped
/// name and is in use.
pub const CORE_USER: &str = "core";
/// The prefix `hive-c0re`'s `mint_token` put before the unix seconds.
const LEGACY_TOKEN_PREFIX: &str = "hyperhive-";
/// Whether `name` is exactly `hyperhive-<digits>`, the shape of every token
/// `hive-c0re` minted.
pub fn is_legacy_token_name(name: &str) -> bool {
name.strip_prefix(LEGACY_TOKEN_PREFIX)
.is_some_and(|secs| !secs.is_empty() && secs.bytes().all(|b| b.is_ascii_digit()))
}
/// The agents a sweep may touch: those whose [`AGENT_TOKEN_NAME`] token is
/// current, never [`CORE_USER`].
pub fn sweepable(observed: &[(String, Observed)]) -> Vec<String> {
observed
.iter()
.filter(|(agent, o)| agent != CORE_USER && *o == Observed::Decided(Decision::Keep))
.map(|(agent, _)| agent.clone())
.collect()
}
impl Client {
/// Delete every legacy token the forge lists for `user`, and return how
/// many went. A delete the forge refuses is logged and skipped.
async fn sweep_user(&self, user: &str) -> Result<usize> {
if user == CORE_USER {
return Ok(0);
}
let Some(listed) = self.list_agent_tokens(user).await? else {
return Ok(0);
};
let mut removed = 0;
for name in listed
.iter()
.filter_map(|t| t.name.as_deref())
.filter(|n| is_legacy_token_name(n))
{
match self.api.admin_delete_user_access_token(user, name).await {
Ok(()) => removed += 1,
Err(e) if super::agent_token::is_not_found(&e) => {}
Err(e) => tracing::warn!(
agent = user,
token = name,
error = %e,
"legacy forge token sweep: delete failed; skipped"
),
}
}
Ok(removed)
}
/// Sweep each of `agents` and log what went, per agent and in total.
/// Returns the total.
async fn sweep_legacy_tokens(&self, agents: &[String]) -> usize {
let mut total = 0;
for agent in agents {
match self.sweep_user(agent).await {
Ok(removed) => {
tracing::info!(agent, removed, "legacy forge token sweep: agent done");
total += removed;
}
Err(e) => tracing::warn!(
agent,
error = %format!("{e:#}"),
"legacy forge token sweep: listing failed; agent skipped"
),
}
}
tracing::info!(
agents = agents.len(),
removed = total,
"legacy forge token sweep: done"
);
total
}
/// The roster agents [`sweepable`] allows, observed the way the mint pass
/// observes them.
async fn sweep_roster(&self) -> Result<Vec<String>> {
let store = crate::store::connect()
.await
.context("logging in to the swarm secret store")?;
let roles = store
.list_cert_roles(DEFAULT_CERT_MOUNT)
.await
.context("listing the store's cert-auth roles")?;
let mut observed = Vec::new();
for agent in policy::agents_from_role_names(&roles) {
let o = self.observe_agent(&store, &agent).await;
if o != Observed::Decided(Decision::Keep) {
tracing::info!(
agent,
?o,
"legacy forge token sweep: {AGENT_TOKEN_NAME} token not current; agent skipped"
);
}
observed.push((agent, o));
}
Ok(sweepable(&observed))
}
}
/// Run one sweep in the background. A failure is logged and not retried.
pub fn spawn(client: Arc<Client>) {
tokio::spawn(async move {
match client.sweep_roster().await {
Ok(agents) => {
client.sweep_legacy_tokens(&agents).await;
}
Err(e) => tracing::warn!(
error = %format!("{e:#}"),
"legacy forge token sweep: roster unavailable; not run"
),
}
});
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::sync::Mutex;
use forgejo_api::{Auth, Forgejo};
use super::super::agent_token::MintReason;
use super::*;
#[test]
fn only_hyperhive_digits_is_legacy() {
for name in ["hyperhive-1", "hyperhive-1700000000"] {
assert!(is_legacy_token_name(name), "{name}");
}
for name in [
AGENT_TOKEN_NAME,
"hyperhive",
"hyperhive-",
"hyperhive-abc",
"hyperhive-123abc",
"hyperhive-12 3",
"hyperhive-+123",
"xhyperhive-123",
"hyperhive-123-x",
"swarm-controller-boot-1700000000",
] {
assert!(!is_legacy_token_name(name), "{name}");
}
}
#[test]
fn only_current_non_core_agents_are_sweepable() {
let observed = [
("a".to_owned(), Observed::Decided(Decision::Keep)),
(CORE_USER.to_owned(), Observed::Decided(Decision::Keep)),
(
"b".to_owned(),
Observed::Decided(Decision::Mint(MintReason::NotOnForge)),
),
("c".to_owned(), Observed::Unknown),
("d".to_owned(), Observed::NoForgeUser),
("e".to_owned(), Observed::Decided(Decision::Conflict)),
("f".to_owned(), Observed::Decided(Decision::Keep)),
];
assert_eq!(sweepable(&observed), ["a", "f"]);
}
/// Each user's token names, as the stub forge holds them.
type Tokens = Arc<Mutex<BTreeMap<String, Vec<String>>>>;
/// Every `DELETE` the stub forge was sent, as `(user, token)`.
type Deletes = Arc<Mutex<Vec<(String, String)>>>;
/// A stand-in forge on a loopback port answering the admin token list and
/// delete from `tokens`. A delete of `fail` answers 500 and removes
/// nothing.
async fn stub_forge(tokens: Tokens, fail: (&'static str, &'static str)) -> (Client, Deletes) {
use axum::http::{Method, StatusCode, Uri, header::CONTENT_TYPE};
let deletes = Deletes::default();
let recorded = deletes.clone();
let app = axum::Router::new().fallback(move |method: Method, uri: Uri| {
let tokens = tokens.clone();
let recorded = recorded.clone();
async move {
let parts: Vec<&str> = uri.path().trim_matches('/').split('/').collect();
let (status, body) = match (method, parts.as_slice()) {
(Method::GET, ["api", "v1", "admin", "users", user, "tokens"]) => {
match tokens.lock().unwrap().get(*user) {
Some(held) => {
let list: Vec<serde_json::Value> = held
.iter()
.enumerate()
.map(|(i, name)| serde_json::json!({ "id": i, "name": name }))
.collect();
(StatusCode::OK, serde_json::Value::from(list))
}
None => (
StatusCode::NOT_FOUND,
serde_json::json!({ "message": "user does not exist" }),
),
}
}
(Method::DELETE, ["api", "v1", "admin", "users", user, "tokens", token]) => {
recorded
.lock()
.unwrap()
.push(((*user).to_owned(), (*token).to_owned()));
if (*user, *token) == fail {
(StatusCode::INTERNAL_SERVER_ERROR, serde_json::json!({}))
} else {
let mut all = tokens.lock().unwrap();
let held = all.entry((*user).to_owned()).or_default();
held.retain(|n| n != token);
(StatusCode::NO_CONTENT, serde_json::Value::Null)
}
}
_ => (
StatusCode::NOT_FOUND,
serde_json::json!({ "message": "not stubbed" }),
),
};
let count = body.as_array().map_or(0, Vec::len).to_string();
let body = if status == StatusCode::NO_CONTENT {
String::new()
} else {
body.to_string()
};
(
status,
[
(CONTENT_TYPE.as_str(), "application/json".to_owned()),
("x-total-count", count),
],
body,
)
}
});
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = url::Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap();
tokio::spawn(async move { axum::serve(listener, app).await });
let api = Forgejo::new(Auth::Token("stub-token"), url.clone()).unwrap();
(Client { api, url }, deletes)
}
const MIX: [&str; 5] = [
"hyperhive-123",
"hyperhive-1700000000",
"hyperhive-abc",
"hyperhive",
AGENT_TOKEN_NAME,
];
fn forge_with_core_and_two_agents() -> Tokens {
let held = |names: &[&str]| names.iter().map(|n| (*n).to_owned()).collect();
Arc::new(Mutex::new(BTreeMap::from([
(CORE_USER.to_owned(), held(&MIX)),
("alice".to_owned(), held(&MIX)),
("bob".to_owned(), held(&MIX)),
])))
}
fn sorted(deletes: &Deletes) -> Vec<(String, String)> {
let mut d = deletes.lock().unwrap().clone();
d.sort();
d
}
fn pair(user: &str, token: &str) -> (String, String) {
(user.to_owned(), token.to_owned())
}
/// Core is in the list handed to the sweep, as a stand-in for any path
/// that lets it through [`sweepable`]; its tokens must survive anyway.
#[tokio::test]
async fn only_agents_legacy_tokens_go_and_a_rerun_removes_nothing() {
let tokens = forge_with_core_and_two_agents();
let (client, deletes) = stub_forge(tokens.clone(), ("", "")).await;
let users = [CORE_USER.to_owned(), "alice".to_owned(), "bob".to_owned()];
assert_eq!(client.sweep_legacy_tokens(&users).await, 4);
assert_eq!(
sorted(&deletes),
[
pair("alice", "hyperhive-123"),
pair("alice", "hyperhive-1700000000"),
pair("bob", "hyperhive-123"),
pair("bob", "hyperhive-1700000000"),
]
);
let left = tokens.lock().unwrap().clone();
assert_eq!(left[CORE_USER], MIX);
for agent in ["alice", "bob"] {
assert_eq!(
left[agent],
["hyperhive-abc", "hyperhive", AGENT_TOKEN_NAME]
);
}
deletes.lock().unwrap().clear();
assert_eq!(client.sweep_legacy_tokens(&users).await, 0);
assert!(deletes.lock().unwrap().is_empty());
}
#[tokio::test]
async fn a_refused_delete_moves_on_to_the_next() {
let tokens = forge_with_core_and_two_agents();
let (client, deletes) = stub_forge(tokens.clone(), ("alice", "hyperhive-123")).await;
let users = ["alice".to_owned(), "bob".to_owned()];
assert_eq!(client.sweep_legacy_tokens(&users).await, 3);
assert_eq!(deletes.lock().unwrap().len(), 4);
assert_eq!(
tokens.lock().unwrap()["alice"],
[
"hyperhive-123",
"hyperhive-abc",
"hyperhive",
AGENT_TOKEN_NAME
]
);
}
#[tokio::test]
async fn a_user_the_forge_does_not_have_removes_nothing() {
let tokens = forge_with_core_and_two_agents();
let (client, deletes) = stub_forge(tokens, ("", "")).await;
assert_eq!(client.sweep_legacy_tokens(&["ruth".to_owned()]).await, 0);
assert!(deletes.lock().unwrap().is_empty());
}
}

View file

@ -2159,8 +2159,8 @@ fn spawn_matrix_account_backfill(
}
/// Start the forge's periodic passes — the swarm-wide objects and agents'
/// tokens — when a forge is configured. Lifted out of `main` for
/// `clippy::too_many_lines`.
/// tokens — and the one-shot legacy token sweep, when a forge is configured.
/// Lifted out of `main` for `clippy::too_many_lines`.
fn spawn_forge_workers(
jobq: &Arc<Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>>,
forge_client: Option<Arc<forge::Client>>,
@ -2169,6 +2169,7 @@ fn spawn_forge_workers(
return;
};
forge::objects::spawn(Arc::clone(&client));
forge::legacy_tokens::spawn(Arc::clone(&client));
let sched = Arc::clone(jobq);
forge::agent_token::spawn(client, move |agents| {
if let Err(e) = queue_forge_token_mints(&sched, agents) {