Watch
0
0
Fork
You've already forked hyperhive
0

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
This commit is contained in:
atlas 2026-09-24 17:37:09 +02:00 • committed by mara
commit 22f0acfd6d
2 changed files with 112 additions and 30 deletions

View file

@ -126,21 +126,30 @@ pub fn classify(stored: Option<&forge::Credential>, listed: &[AccessToken]) -> D
pub enum Observed { pub enum Observed {
/// Both reads worked, and this is what [`classify`] made of them. /// Both reads worked, and this is what [`classify`] made of them.
Decided(Decision), Decided(Decision),
/// The forge has no user by this agent's name. Creating one is agent /// The forge has no user by this agent's name. Planned like a mint: the
/// creation's job, not this module's. /// 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, NoForgeUser,
/// A read failed. Nothing is known, so nothing is done. /// A read failed. Nothing is known, so nothing is done.
Unknown, 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 /// [`Observed::Unknown`] in particular never mints: a forge or store outage
/// forge or store outage must not turn into a rotation of every agent's token. /// must not turn into a rotation of every agent's token.
pub fn plan(observed: &[(String, Observed)]) -> Vec<String> { pub fn plan(observed: &[(String, Observed)]) -> Vec<String> {
observed observed
.iter() .iter()
.filter(|(_, o)| matches!(o, Observed::Decided(Decision::Mint(_)))) .filter(|(_, o)| {
matches!(
o,
Observed::Decided(Decision::Mint(_)) | Observed::NoForgeUser
)
})
.map(|(agent, _)| agent.clone()) .map(|(agent, _)| agent.clone())
.collect() .collect()
} }
@ -191,7 +200,10 @@ impl Client {
.await .await
.with_context(|| format!("reading {path}"))?; .with_context(|| format!("reading {path}"))?;
let Some(listed) = self.list_agent_tokens(agent).await? else { 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) { let reason = match classify(stored.as_ref(), &listed) {
Decision::Keep => { Decision::Keep => {
@ -299,7 +311,10 @@ impl Client {
for agent in policy::agents_from_role_names(&roles) { for agent in policy::agents_from_role_names(&roles) {
let o = self.observe_agent(&store, &agent).await; let o = self.observe_agent(&store, &agent).await;
if o == Observed::NoForgeUser { 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)); observed.push((agent, o));
} }
@ -497,7 +512,7 @@ mod tests {
} }
#[test] #[test]
fn only_a_mint_decision_is_planned() { fn a_mint_or_a_missing_user_is_planned() {
let observed = [ let observed = [
("a".to_owned(), Observed::Decided(Decision::Keep)), ("a".to_owned(), Observed::Decided(Decision::Keep)),
( (
@ -512,7 +527,16 @@ mod tests {
Observed::Decided(Decision::Mint(MintReason::ScopeMismatch)), 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] #[test]

View file

@ -1756,11 +1756,16 @@ struct MintAgentForgeTokenResponse {
node_id: u64, 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 /// 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 /// 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( fn queue_forge_token_mints(
sched: &Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>, sched: &Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>,
agents: Vec<String>, agents: Vec<String>,
@ -1772,7 +1777,14 @@ fn queue_forge_token_mints(
for agent in agents { for agent in agents {
let queued = sched let queued = sched
.insert_job(None, |b| { .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}"))?; .map_err(|e| anyhow::anyhow!("{e}"))?;
ids.extend(queued); ids.extend(queued);
@ -3033,9 +3045,10 @@ mod tests {
/// The manual route and the periodic pass both go through /// The manual route and the periodic pass both go through
/// `queue_forge_token_mints`, so this is the assertion that a backfilled /// `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] #[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 _; use hive_jobq_wire::WireNode as _;
let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new( let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new(
@ -3043,24 +3056,69 @@ mod tests {
hive_jobq::resources::ResourceTable::new(), hive_jobq::resources::ResourceTable::new(),
)); ));
let ids = super::queue_forge_token_mints(&sched, vec!["a".to_owned(), "b".to_owned()]) let ids = super::queue_forge_token_mints(&sched, vec!["a".to_owned(), "b".to_owned()])
.expect("two single-node jobs insert"); .expect("two jobs insert");
assert_eq!(ids.len(), 2); assert_eq!(ids.len(), 2, "one returned id per agent: its mint node");
let guard = sched let guard = sched
.lock() .lock()
.unwrap_or_else(std::sync::PoisonError::into_inner); .unwrap_or_else(std::sync::PoisonError::into_inner);
let mut agents: Vec<String> = guard let graph = guard.graph();
.graph() let agent_of = |n: &hive_jobq::Node<SwarmNodeKind, super::SwarmResourceKind>| {
.nodes() n.payload.data(n.id.get())["agent"]
.map(|n| { .as_str()
assert_eq!(n.payload.label(), "mint_agent_forge_token"); .expect("agent is a string")
n.payload.data(n.id.get())["agent"] .to_owned()
.as_str() };
.expect("agent is a string") for agent in ["a", "b"] {
.to_owned() let find = |label: &str| {
}) graph
.collect(); .nodes()
agents.sort(); .find(|n| n.payload.label() == label && agent_of(n) == agent)
assert_eq!(agents, ["a", "b"]); .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<String> = 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 /// The socket must not share a directory with anything else, because