From 22f0acfd6d8b76f333c73ac32b30c1b13d2d1210 Mon Sep 17 00:00:00 2001 From: atlas Date: Thu, 24 Sep 2026 17:37:09 +0200 Subject: [PATCH] swarm-controller: the forge-token backfill creates a missing forge user An agent with a hive-agent-* store identity but no forge user was observed as NoForgeUser and dropped by plan(), so it never got a token. plan() now keeps it, and queue_forge_token_mints inserts CreateForgeUser ahead of MintAgentForgeToken with after_ok, the edge declare_agent_job already uses. ensure_agent_user folds an existing user into success, so the extra node is a no-op for agents that have one. Refs #3782 --- swarm-controller/src/forge/agent_token.rs | 44 +++++++--- swarm-controller/src/main.rs | 98 ++++++++++++++++++----- 2 files changed, 112 insertions(+), 30 deletions(-) diff --git a/swarm-controller/src/forge/agent_token.rs b/swarm-controller/src/forge/agent_token.rs index ee8faec9..19a322d8 100644 --- a/swarm-controller/src/forge/agent_token.rs +++ b/swarm-controller/src/forge/agent_token.rs @@ -126,21 +126,30 @@ pub fn classify(stored: Option<&forge::Credential>, listed: &[AccessToken]) -> D pub enum Observed { /// Both reads worked, and this is what [`classify`] made of them. Decided(Decision), - /// The forge has no user by this agent's name. Creating one is agent - /// creation's job, not this module's. + /// The forge has no user by this agent's name. Planned like a mint: the + /// queued job creates the user first (`CreateForgeUser`, idempotent), so + /// an agent that holds a store identity but was never given a forge + /// account still ends up with both. NoForgeUser, /// A read failed. Nothing is known, so nothing is done. Unknown, } -/// The agents a pass mints for. +/// The agents a pass mints for: a [`Decision::Mint`], or an agent with no +/// forge user yet ([`Observed::NoForgeUser`]), whose job creates the user +/// before it mints. /// -/// Only [`Decision::Mint`]. [`Observed::Unknown`] in particular never mints: a -/// forge or store outage must not turn into a rotation of every agent's token. +/// [`Observed::Unknown`] in particular never mints: a forge or store outage +/// must not turn into a rotation of every agent's token. pub fn plan(observed: &[(String, Observed)]) -> Vec { observed .iter() - .filter(|(_, o)| matches!(o, Observed::Decided(Decision::Mint(_)))) + .filter(|(_, o)| { + matches!( + o, + Observed::Decided(Decision::Mint(_)) | Observed::NoForgeUser + ) + }) .map(|(agent, _)| agent.clone()) .collect() } @@ -191,7 +200,10 @@ impl Client { .await .with_context(|| format!("reading {path}"))?; let Some(listed) = self.list_agent_tokens(agent).await? else { - bail!("the forge has no user {agent:?}, so there is nothing to mint a token for"); + bail!( + "the forge has no user {agent:?} even after this job's create_forge_user \ + node, so there is nothing to mint a token for" + ); }; let reason = match classify(stored.as_ref(), &listed) { Decision::Keep => { @@ -299,7 +311,10 @@ impl Client { for agent in policy::agents_from_role_names(&roles) { let o = self.observe_agent(&store, &agent).await; if o == Observed::NoForgeUser { - tracing::debug!(agent, "agent forge token: no forge user; skipped"); + tracing::info!( + agent, + "agent forge token: no forge user; creating one first" + ); } observed.push((agent, o)); } @@ -497,7 +512,7 @@ mod tests { } #[test] - fn only_a_mint_decision_is_planned() { + fn a_mint_or_a_missing_user_is_planned() { let observed = [ ("a".to_owned(), Observed::Decided(Decision::Keep)), ( @@ -512,7 +527,16 @@ mod tests { Observed::Decided(Decision::Mint(MintReason::ScopeMismatch)), ), ]; - assert_eq!(plan(&observed), ["b", "f"]); + assert_eq!(plan(&observed), ["b", "d", "f"]); + } + + /// An agent with a store identity and no forge user used to be skipped + /// on every pass, forever. It is planned now; the job it gets creates + /// the user before it mints (see `queue_forge_token_mints`). + #[test] + fn a_missing_forge_user_is_planned() { + let observed = [("ruth".to_owned(), Observed::NoForgeUser)]; + assert_eq!(plan(&observed), ["ruth"]); } #[test] diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index fe785c87..e43fcb0e 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -1756,11 +1756,16 @@ struct MintAgentForgeTokenResponse { node_id: u64, } -/// Insert one `MintAgentForgeToken` node per agent, each its own job. +/// Insert one job per agent: `CreateForgeUser`, then `MintAgentForgeToken` +/// `after_ok` it — the same edge [`declare_agent_job`] uses. Returns the mint +/// nodes' ids. /// /// The one place that node is queued outside agent creation: both the manual /// route below and `forge::agent_token::spawn`'s periodic pass come through -/// here, so a backfilled mint is the same node a new agent gets. +/// here, so a backfilled mint is the same node a new agent gets. The user +/// node is what lets the pass reach an agent that has a store identity but +/// no forge account; for one that has an account it is a no-op, because +/// `forge::Client::ensure_agent_user` folds "already exists" into success. fn queue_forge_token_mints( sched: &Mutex>, agents: Vec, @@ -1772,7 +1777,14 @@ fn queue_forge_token_mints( for agent in agents { let queued = sched .insert_job(None, |b| { - vec![b.node(SwarmNodeKind::MintAgentForgeToken { agent }).guid()] + let create_forge_user = b.node(SwarmNodeKind::CreateForgeUser { + agent: agent.clone(), + }); + vec![ + b.node(SwarmNodeKind::MintAgentForgeToken { agent }) + .after_ok(create_forge_user) + .guid(), + ] }) .map_err(|e| anyhow::anyhow!("{e}"))?; ids.extend(queued); @@ -3033,9 +3045,10 @@ mod tests { /// The manual route and the periodic pass both go through /// `queue_forge_token_mints`, so this is the assertion that a backfilled - /// mint is the same node agent creation inserts. + /// mint is the same pair of nodes agent creation inserts: the forge user, + /// then the mint `after_ok` it, per agent. #[test] - fn a_queued_mint_is_one_forge_token_node_per_agent() { + fn a_queued_mint_creates_the_forge_user_first() { use hive_jobq_wire::WireNode as _; let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new( @@ -3043,24 +3056,69 @@ mod tests { hive_jobq::resources::ResourceTable::new(), )); let ids = super::queue_forge_token_mints(&sched, vec!["a".to_owned(), "b".to_owned()]) - .expect("two single-node jobs insert"); - assert_eq!(ids.len(), 2); + .expect("two jobs insert"); + assert_eq!(ids.len(), 2, "one returned id per agent: its mint node"); let guard = sched .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); - let mut agents: Vec = guard - .graph() - .nodes() - .map(|n| { - assert_eq!(n.payload.label(), "mint_agent_forge_token"); - n.payload.data(n.id.get())["agent"] - .as_str() - .expect("agent is a string") - .to_owned() - }) - .collect(); - agents.sort(); - assert_eq!(agents, ["a", "b"]); + let graph = guard.graph(); + let agent_of = |n: &hive_jobq::Node| { + n.payload.data(n.id.get())["agent"] + .as_str() + .expect("agent is a string") + .to_owned() + }; + for agent in ["a", "b"] { + let find = |label: &str| { + graph + .nodes() + .find(|n| n.payload.label() == label && agent_of(n) == agent) + .unwrap_or_else(|| panic!("{agent} has a {label} node")) + }; + let user = find("create_forge_user").id; + let mint = find("mint_agent_forge_token"); + assert!(ids.contains(&mint.id), "the returned id is {agent}'s mint"); + let when = mint + .deps + .iter() + .find_map(|d| match d { + hive_jobq::Dep::Node { id, when } if *id == user => Some(*when), + _ => None, + }) + .unwrap_or_else(|| panic!("{agent}'s mint waits for its own forge user")); + assert!( + !when.accepts(hive_jobq::TerminalState::Failed), + "a token cannot be minted for a user that was never created; this \ + edge has to be `after_ok`" + ); + } + assert_eq!( + graph.nodes().count(), + 4, + "two nodes per agent, nothing else" + ); + } + + /// The backfill end to end, minus the IO: an agent the pass observed with + /// no forge user is planned, and the job it gets creates that user before + /// it mints. + #[test] + fn an_agent_with_no_forge_user_gets_a_user_and_a_mint() { + use crate::forge::agent_token::{Observed, plan}; + use hive_jobq_wire::WireNode as _; + + let agents = plan(&[("ruth".to_owned(), Observed::NoForgeUser)]); + let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new( + hive_jobq::Graph::new(), + hive_jobq::resources::ResourceTable::new(), + )); + super::queue_forge_token_mints(&sched, agents).expect("one job inserts"); + let guard = sched + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let mut labels: Vec = guard.graph().nodes().map(|n| n.payload.label()).collect(); + labels.sort(); + assert_eq!(labels, ["create_forge_user", "mint_agent_forge_token"]); } /// The socket must not share a directory with anything else, because