swarm-controller: mint each agent's forge token and store it in bao
A MintAgentForgeToken node mints a fixed-name swarm-agent token with the admin API, keeps it when the stored value's last eight and the normalised scopes match the forge's list, and otherwise deletes and re-creates it. The token is stored at swarm/agents/<agent>/forge-token. Agent creation inserts the node, and a pass at start and every five minutes inserts it for every agent holding a store identity whose token is missing or stale. Refs #3782
This commit is contained in:
parent
2fda529ca8
commit
52c8c0b0de
7 changed files with 977 additions and 4 deletions
|
|
@ -33,6 +33,8 @@ use utoipa::ToSchema;
|
|||
|
||||
use crate::webhook::DeliveryKind;
|
||||
|
||||
pub mod agent_token;
|
||||
|
||||
/// An agent's open config-PR, as [`Client::list_open_config_prs`] reports it
|
||||
/// and `GET /api/agents/{name}/config-pr` serves it.
|
||||
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, ToSchema)]
|
||||
|
|
@ -135,6 +137,10 @@ const TOKEN_FILE_ENV: &str = "SWARM_CONTROLLER_FORGE_TOKEN_FILE";
|
|||
/// Typed forgejo client, built once at startup from env.
|
||||
pub struct Client {
|
||||
api: Forgejo,
|
||||
/// The forge's base URL, kept so a client authenticated as someone
|
||||
/// else can be built against the same forge (see
|
||||
/// [`agent_token`]'s read-back).
|
||||
url: url::Url,
|
||||
}
|
||||
|
||||
impl Client {
|
||||
|
|
@ -168,8 +174,12 @@ impl Client {
|
|||
}
|
||||
let parsed_url =
|
||||
url::Url::parse(&url).with_context(|| format!("parse {URL_ENV} ({url}) as a URL"))?;
|
||||
let api = Forgejo::new(Auth::Token(token), parsed_url).context("build forgejo client")?;
|
||||
Ok(Some(Self { api }))
|
||||
let api =
|
||||
Forgejo::new(Auth::Token(token), parsed_url.clone()).context("build forgejo client")?;
|
||||
Ok(Some(Self {
|
||||
api,
|
||||
url: parsed_url,
|
||||
}))
|
||||
}
|
||||
|
||||
/// Creation options for an empty repo defaulting to `main`.
|
||||
|
|
|
|||
539
swarm-controller/src/forge/agent_token.rs
Normal file
539
swarm-controller/src/forge/agent_token.rs
Normal file
|
|
@ -0,0 +1,539 @@
|
|||
//! Each agent's own forge access token: minted here with the admin API, stored
|
||||
//! at `swarm/agents/<agent>/forge-token`, and pulled from there by the agent
|
||||
//! container itself (`nix/agent-modules/forge-token.nix`).
|
||||
//!
|
||||
//! One token per agent, under the fixed name [`AGENT_TOKEN_NAME`]. Forgejo
|
||||
//! refuses a second token with a name the user already has, so a rotation is a
|
||||
//! delete then a create, and there is never more than one of these per agent.
|
||||
//! The `hyperhive-<unix-seconds>` tokens `hive-c0re` used to mint are never
|
||||
//! touched here: the name cannot match them.
|
||||
//!
|
||||
//! Two callers insert the same `MintAgentForgeToken` job node: agent creation,
|
||||
//! and [`spawn`]'s pass at start and every [`RECONCILE_INTERVAL`] over every
|
||||
//! agent that holds a store identity. The decision is pure ([`classify`] and
|
||||
//! [`plan`]), so the tests pin it; the IO on either side only reads or acts.
|
||||
|
||||
use std::collections::BTreeSet;
|
||||
use std::sync::Arc;
|
||||
|
||||
use anyhow::{Context, Result, bail};
|
||||
use forgejo_api::structs::{AccessToken, CreateAccessTokenOption};
|
||||
use forgejo_api::{ApiErrorKind, Auth, Forgejo, ForgejoError};
|
||||
use reqwest::StatusCode;
|
||||
use swarm_secret_client::{client::DEFAULT_CERT_MOUNT, forge, policy};
|
||||
|
||||
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.
|
||||
pub const AGENT_TOKEN_NAME: &str = "swarm-agent";
|
||||
|
||||
/// The scopes an agent token carries. Byte-identical to `hive-c0re`'s
|
||||
/// `forge::users::TOKEN_SCOPES`, so the move from hive to swarm changes
|
||||
/// nothing an agent can do. Narrowing it is a separate decision.
|
||||
///
|
||||
/// ⚠️ Duplicated across the crate boundary until the last `hive-c0re` caller
|
||||
/// of that constant goes. A test here pins the literal.
|
||||
pub const AGENT_TOKEN_SCOPES: &str = "read:user,write:user,read:notification,write:notification,write:repository,write:issue,write:organization,write:misc";
|
||||
|
||||
/// How often [`spawn`] re-checks every agent's token.
|
||||
const RECONCILE_INTERVAL: std::time::Duration = std::time::Duration::from_mins(5);
|
||||
|
||||
/// Why a token has to be (re)minted.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MintReason {
|
||||
/// The store holds nothing for this agent.
|
||||
NotStored,
|
||||
/// The forge lists no [`AGENT_TOKEN_NAME`] token for this agent.
|
||||
NotOnForge,
|
||||
/// The forge's token is not the one the store holds.
|
||||
LastEightMismatch,
|
||||
/// The forge's token carries other scopes than [`AGENT_TOKEN_SCOPES`].
|
||||
ScopeMismatch,
|
||||
}
|
||||
|
||||
/// What to do about one agent's token.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum Decision {
|
||||
/// The stored token is the forge's token, with the right scopes.
|
||||
Keep,
|
||||
/// Delete any [`AGENT_TOKEN_NAME`] token, create a new one, store it.
|
||||
Mint(MintReason),
|
||||
/// The forge lists more than one [`AGENT_TOKEN_NAME`] token. Forgejo
|
||||
/// refuses a duplicate name, so this is a forge in a state nothing here
|
||||
/// can reason about; a delete by name would 422 on it anyway.
|
||||
Conflict,
|
||||
}
|
||||
|
||||
/// A scope list reduced to what Forgejo would store for it.
|
||||
///
|
||||
/// Forgejo normalises a new token's scopes before storing them
|
||||
/// (`models/auth/access_token_scope.go`, `Normalize` → `toScope`): a `read:X`
|
||||
/// beside `write:X` is dropped, because the write bit implies the read one.
|
||||
/// The admin list then returns the stored form. Comparing
|
||||
/// [`AGENT_TOKEN_SCOPES`] to that list as written would never match, and every
|
||||
/// pass would rotate every agent's token.
|
||||
fn normalized_scopes<'a>(scopes: impl IntoIterator<Item = &'a str>) -> BTreeSet<String> {
|
||||
let all: BTreeSet<&str> = scopes.into_iter().map(str::trim).collect();
|
||||
all.iter()
|
||||
.filter(|s| {
|
||||
s.strip_prefix("read:")
|
||||
.is_none_or(|area| !all.contains(format!("write:{area}").as_str()))
|
||||
})
|
||||
.map(|s| (*s).to_owned())
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The last eight characters of `token`, as Forgejo reports them for a token
|
||||
/// it holds (`token_last_eight`). `None` for a token shorter than that, which
|
||||
/// no Forgejo token is.
|
||||
fn last_eight(token: &str) -> Option<&str> {
|
||||
token.len().checked_sub(8).and_then(|at| token.get(at..))
|
||||
}
|
||||
|
||||
/// Decide what to do about an agent's token, from what the store holds and
|
||||
/// what the forge lists for the agent.
|
||||
///
|
||||
/// Every token not named [`AGENT_TOKEN_NAME`] is ignored, so the
|
||||
/// `hyperhive-*` tokens `hive-c0re` left behind never affect the decision.
|
||||
pub fn classify(stored: Option<&forge::Credential>, listed: &[AccessToken]) -> Decision {
|
||||
let ours: Vec<&AccessToken> = listed
|
||||
.iter()
|
||||
.filter(|t| t.name.as_deref() == Some(AGENT_TOKEN_NAME))
|
||||
.collect();
|
||||
let token = match ours[..] {
|
||||
[] => return Decision::Mint(MintReason::NotOnForge),
|
||||
[token] => token,
|
||||
_ => return Decision::Conflict,
|
||||
};
|
||||
let Some(stored) = stored else {
|
||||
return Decision::Mint(MintReason::NotStored);
|
||||
};
|
||||
if token.token_last_eight.as_deref() != last_eight(&stored.value) {
|
||||
return Decision::Mint(MintReason::LastEightMismatch);
|
||||
}
|
||||
let listed_scopes = normalized_scopes(token.scopes.iter().flatten().map(String::as_str));
|
||||
if listed_scopes != normalized_scopes(AGENT_TOKEN_SCOPES.split(',')) {
|
||||
return Decision::Mint(MintReason::ScopeMismatch);
|
||||
}
|
||||
Decision::Keep
|
||||
}
|
||||
|
||||
/// What one pass found for one agent.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
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.
|
||||
NoForgeUser,
|
||||
/// A read failed. Nothing is known, so nothing is done.
|
||||
Unknown,
|
||||
}
|
||||
|
||||
/// The agents a pass mints for.
|
||||
///
|
||||
/// 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.
|
||||
pub fn plan(observed: &[(String, Observed)]) -> Vec<String> {
|
||||
observed
|
||||
.iter()
|
||||
.filter(|(_, o)| matches!(o, Observed::Decided(Decision::Mint(_))))
|
||||
.map(|(agent, _)| agent.clone())
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Whether a forge error is a 404, whichever of its two shapes it came in.
|
||||
fn is_not_found(e: &ForgejoError) -> bool {
|
||||
match e {
|
||||
ForgejoError::ApiError(api) => {
|
||||
matches!(api.error_kind(), ApiErrorKind::NotFound { .. })
|
||||
|| matches!(api.error_kind(), ApiErrorKind::Other(s) if *s == StatusCode::NOT_FOUND)
|
||||
}
|
||||
ForgejoError::UnexpectedStatusCode(s) => *s == StatusCode::NOT_FOUND,
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
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>>> {
|
||||
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),
|
||||
Err(e) => Err(e).with_context(|| format!("listing {agent}'s forge tokens")),
|
||||
}
|
||||
}
|
||||
|
||||
/// Make sure `agent`'s stored token is a live [`AGENT_TOKEN_NAME`] token
|
||||
/// with [`AGENT_TOKEN_SCOPES`], and mint one if it is not. The whole job
|
||||
/// of the `MintAgentForgeToken` node.
|
||||
///
|
||||
/// A rotation deletes the old token first, then creates, stores, and uses
|
||||
/// the new one: it logs in with it and asks the forge who it is, so the
|
||||
/// node does not report success on a write nobody has read. A crash
|
||||
/// between any two steps leaves a state the next run classifies as
|
||||
/// [`Decision::Mint`] again.
|
||||
///
|
||||
/// # Errors
|
||||
/// When the store or the forge refuses a step, when the agent has no forge
|
||||
/// user, or when the new token authenticates as someone else.
|
||||
pub async fn ensure_agent_forge_token(&self, agent: &str) -> Result<()> {
|
||||
let path = forge::agent_token_path(agent)?;
|
||||
let store = crate::store::connect()
|
||||
.await
|
||||
.context("logging in to the swarm secret store")?;
|
||||
let stored: Option<forge::Credential> = store
|
||||
.read_optional(&path)
|
||||
.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");
|
||||
};
|
||||
let reason = match classify(stored.as_ref(), &listed) {
|
||||
Decision::Keep => {
|
||||
tracing::debug!(agent, "agent forge token is current; left as it is");
|
||||
return Ok(());
|
||||
}
|
||||
Decision::Conflict => bail!(
|
||||
"the forge lists more than one {AGENT_TOKEN_NAME:?} token for {agent:?}; \
|
||||
delete them by id and re-run"
|
||||
),
|
||||
Decision::Mint(reason) => reason,
|
||||
};
|
||||
|
||||
if listed
|
||||
.iter()
|
||||
.any(|t| t.name.as_deref() == Some(AGENT_TOKEN_NAME))
|
||||
{
|
||||
match self
|
||||
.api
|
||||
.admin_delete_user_access_token(agent, AGENT_TOKEN_NAME)
|
||||
.await
|
||||
{
|
||||
Ok(()) => {}
|
||||
Err(e) if is_not_found(&e) => {}
|
||||
Err(e) => {
|
||||
return Err(e)
|
||||
.with_context(|| format!("deleting {agent}'s {AGENT_TOKEN_NAME} token"));
|
||||
}
|
||||
}
|
||||
}
|
||||
let created = self
|
||||
.api
|
||||
.admin_create_user_access_token(
|
||||
agent,
|
||||
CreateAccessTokenOption {
|
||||
name: AGENT_TOKEN_NAME.to_owned(),
|
||||
repositories: None,
|
||||
scopes: Some(AGENT_TOKEN_SCOPES.split(',').map(str::to_owned).collect()),
|
||||
},
|
||||
)
|
||||
.await
|
||||
.with_context(|| format!("creating {agent}'s {AGENT_TOKEN_NAME} token"))?;
|
||||
let Some(value) = created.sha1.filter(|v| !v.is_empty()) else {
|
||||
bail!("the forge created {agent}'s token but returned no value for it");
|
||||
};
|
||||
let credential = forge::Credential {
|
||||
value,
|
||||
name: AGENT_TOKEN_NAME.to_owned(),
|
||||
};
|
||||
store
|
||||
.write(&path, &credential)
|
||||
.await
|
||||
.with_context(|| format!("storing {agent}'s forge token at {path}"))?;
|
||||
|
||||
let as_agent = Forgejo::new(Auth::Token(&credential.value), self.url.clone())
|
||||
.context("building a forge client with the new token")?;
|
||||
let me = as_agent
|
||||
.user_get_current()
|
||||
.await
|
||||
.with_context(|| format!("using {agent}'s new forge token"))?;
|
||||
if me.login.as_deref() != Some(agent) {
|
||||
bail!("{agent}'s new forge token authenticates as someone else");
|
||||
}
|
||||
tracing::info!(agent, ?reason, %path, "agent forge token minted and stored");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// What the forge and the store say about `agent`'s token.
|
||||
async fn observe_agent(
|
||||
&self,
|
||||
store: &swarm_secret_client::SecretStore,
|
||||
agent: &str,
|
||||
) -> Observed {
|
||||
let stored = match forge::agent_token_path(agent) {
|
||||
Ok(path) => store.read_optional::<forge::Credential>(&path).await,
|
||||
Err(e) => Err(e),
|
||||
};
|
||||
let stored = match stored {
|
||||
Ok(stored) => stored,
|
||||
Err(e) => {
|
||||
tracing::warn!(agent, error = %e, "agent forge token: store read failed");
|
||||
return Observed::Unknown;
|
||||
}
|
||||
};
|
||||
match self.list_agent_tokens(agent).await {
|
||||
Ok(Some(listed)) => Observed::Decided(classify(stored.as_ref(), &listed)),
|
||||
Ok(None) => Observed::NoForgeUser,
|
||||
Err(e) => {
|
||||
tracing::warn!(agent, error = %format!("{e:#}"), "agent forge token: forge read failed");
|
||||
Observed::Unknown
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// One pass: every agent holding a store identity, observed.
|
||||
async fn observe_all(&self) -> Result<Vec<(String, Observed)>> {
|
||||
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::NoForgeUser {
|
||||
tracing::debug!(agent, "agent forge token: no forge user; skipped");
|
||||
}
|
||||
observed.push((agent, o));
|
||||
}
|
||||
Ok(observed)
|
||||
}
|
||||
}
|
||||
|
||||
/// Check every agent's token now and every [`RECONCILE_INTERVAL`] after, and
|
||||
/// hand the agents that need one to `enqueue`, which inserts a
|
||||
/// `MintAgentForgeToken` node for each.
|
||||
///
|
||||
/// The roster is the store's `hive-agent-*` cert-auth roles: exactly the
|
||||
/// agents that can read what the node writes. The first tick fires
|
||||
/// immediately (`tokio::time::interval`'s default). A pass that fails is
|
||||
/// logged and retried on the next tick; it never stops the daemon.
|
||||
pub fn spawn(client: Arc<Client>, enqueue: impl Fn(Vec<String>) + Send + 'static) {
|
||||
tokio::spawn(async move {
|
||||
let mut ticker = tokio::time::interval(RECONCILE_INTERVAL);
|
||||
loop {
|
||||
ticker.tick().await;
|
||||
match client.observe_all().await {
|
||||
Ok(observed) => {
|
||||
let agents = plan(&observed);
|
||||
if agents.is_empty() {
|
||||
tracing::debug!(
|
||||
checked = observed.len(),
|
||||
"agent forge tokens: all current"
|
||||
);
|
||||
} else {
|
||||
tracing::info!(
|
||||
checked = observed.len(),
|
||||
minting = agents.len(),
|
||||
"agent forge tokens: queueing mints"
|
||||
);
|
||||
enqueue(agents);
|
||||
}
|
||||
}
|
||||
Err(e) => tracing::warn!(
|
||||
error = %format!("{e:#}"),
|
||||
retry_in_s = RECONCILE_INTERVAL.as_secs(),
|
||||
"agent forge tokens: pass failed; retrying next tick"
|
||||
),
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
const VALUE: &str = "0123456789abcdef0123456789abcdef01234567";
|
||||
|
||||
fn stored() -> forge::Credential {
|
||||
forge::Credential {
|
||||
value: VALUE.to_owned(),
|
||||
name: AGENT_TOKEN_NAME.to_owned(),
|
||||
}
|
||||
}
|
||||
|
||||
fn token(name: &str, last_eight: &str, scopes: &[&str]) -> AccessToken {
|
||||
AccessToken {
|
||||
created_at: None,
|
||||
id: Some(1),
|
||||
name: Some(name.to_owned()),
|
||||
repositories: None,
|
||||
scopes: Some(scopes.iter().map(|s| (*s).to_owned()).collect()),
|
||||
sha1: None,
|
||||
token_last_eight: Some(last_eight.to_owned()),
|
||||
}
|
||||
}
|
||||
|
||||
/// What Forgejo v16 stores and lists for [`AGENT_TOKEN_SCOPES`]: the
|
||||
/// `toScope` order, with each `read:X` folded into its `write:X`.
|
||||
const LISTED_SCOPES: [&str; 6] = [
|
||||
"write:misc",
|
||||
"write:notification",
|
||||
"write:organization",
|
||||
"write:issue",
|
||||
"write:repository",
|
||||
"write:user",
|
||||
];
|
||||
|
||||
fn ours() -> AccessToken {
|
||||
token(AGENT_TOKEN_NAME, "01234567", &LISTED_SCOPES)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_scopes_are_hive_c0res_byte_for_byte() {
|
||||
assert_eq!(
|
||||
AGENT_TOKEN_SCOPES,
|
||||
"read:user,write:user,read:notification,write:notification,write:repository,write:issue,write:organization,write:misc"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_current_token_is_kept() {
|
||||
assert_eq!(classify(Some(&stored()), &[ours()]), Decision::Keep);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scopes_as_written_are_kept_too() {
|
||||
// A forge that listed the scopes un-normalised must not rotate either.
|
||||
let t = token(
|
||||
AGENT_TOKEN_NAME,
|
||||
"01234567",
|
||||
&AGENT_TOKEN_SCOPES.split(',').collect::<Vec<_>>(),
|
||||
);
|
||||
assert_eq!(classify(Some(&stored()), &[t]), Decision::Keep);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn nothing_stored_mints() {
|
||||
assert_eq!(
|
||||
classify(None, &[ours()]),
|
||||
Decision::Mint(MintReason::NotStored)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn nothing_on_the_forge_mints() {
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &[]),
|
||||
Decision::Mint(MintReason::NotOnForge)
|
||||
);
|
||||
assert_eq!(classify(None, &[]), Decision::Mint(MintReason::NotOnForge));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_different_last_eight_mints() {
|
||||
let t = token(AGENT_TOKEN_NAME, "ffffffff", &LISTED_SCOPES);
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &[t]),
|
||||
Decision::Mint(MintReason::LastEightMismatch)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_missing_scope_mints() {
|
||||
let t = token(AGENT_TOKEN_NAME, "01234567", &LISTED_SCOPES[1..]);
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &[t]),
|
||||
Decision::Mint(MintReason::ScopeMismatch)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_extra_scope_mints() {
|
||||
let mut scopes = LISTED_SCOPES.to_vec();
|
||||
scopes.push("write:admin");
|
||||
let t = token(AGENT_TOKEN_NAME, "01234567", &scopes);
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &[t]),
|
||||
Decision::Mint(MintReason::ScopeMismatch)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn scope_order_does_not_matter() {
|
||||
let mut scopes = LISTED_SCOPES.to_vec();
|
||||
scopes.reverse();
|
||||
let t = token(AGENT_TOKEN_NAME, "01234567", &scopes);
|
||||
assert_eq!(classify(Some(&stored()), &[t]), Decision::Keep);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legacy_tokens_are_ignored() {
|
||||
// A pile of `hyperhive-*` tokens, one of them even matching the
|
||||
// stored value's last eight, is still "no swarm token".
|
||||
let legacy = [
|
||||
token("hyperhive-1700000000", "01234567", &LISTED_SCOPES),
|
||||
token("hyperhive-1700000001", "aaaaaaaa", &LISTED_SCOPES),
|
||||
];
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &legacy),
|
||||
Decision::Mint(MintReason::NotOnForge)
|
||||
);
|
||||
let mut with_ours = legacy.to_vec();
|
||||
with_ours.push(ours());
|
||||
assert_eq!(classify(Some(&stored()), &with_ours), Decision::Keep);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn two_swarm_tokens_are_a_conflict_not_a_mint() {
|
||||
assert_eq!(
|
||||
classify(Some(&stored()), &[ours(), ours()]),
|
||||
Decision::Conflict
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn last_eight_is_the_tail() {
|
||||
assert_eq!(last_eight(VALUE), Some("01234567"));
|
||||
assert_eq!(last_eight("short"), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_a_mint_decision_is_planned() {
|
||||
let observed = [
|
||||
("a".to_owned(), Observed::Decided(Decision::Keep)),
|
||||
(
|
||||
"b".to_owned(),
|
||||
Observed::Decided(Decision::Mint(MintReason::NotStored)),
|
||||
),
|
||||
("c".to_owned(), Observed::Unknown),
|
||||
("d".to_owned(), Observed::NoForgeUser),
|
||||
("e".to_owned(), Observed::Decided(Decision::Conflict)),
|
||||
(
|
||||
"f".to_owned(),
|
||||
Observed::Decided(Decision::Mint(MintReason::ScopeMismatch)),
|
||||
),
|
||||
];
|
||||
assert_eq!(plan(&observed), ["b", "f"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_outage_plans_nothing() {
|
||||
let observed = [
|
||||
("a".to_owned(), Observed::Unknown),
|
||||
("b".to_owned(), Observed::Unknown),
|
||||
];
|
||||
assert!(plan(&observed).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn both_404_shapes_are_not_found() {
|
||||
assert!(is_not_found(&ForgejoError::UnexpectedStatusCode(
|
||||
StatusCode::NOT_FOUND
|
||||
)));
|
||||
assert!(is_not_found(&ForgejoError::ApiError(
|
||||
ApiErrorKind::NotFound { errors: None }.into()
|
||||
)));
|
||||
assert!(!is_not_found(&ForgejoError::UnexpectedStatusCode(
|
||||
StatusCode::FORBIDDEN
|
||||
)));
|
||||
}
|
||||
}
|
||||
|
|
@ -105,6 +105,14 @@ enum SwarmNodeKind {
|
|||
/// its hive's shared queue credential, so the document cannot be rendered
|
||||
/// without knowing which hive the agent belongs to.
|
||||
MintAgentIdentity { hive: String, agent: String },
|
||||
/// Make sure `agent` holds a live forge access token in the swarm secret
|
||||
/// store, minting one with the forge's admin API when it does not. See
|
||||
/// `forge::agent_token` — including why a rotation is a delete then a
|
||||
/// create.
|
||||
///
|
||||
/// Carries no hive: the token's store path has no hive segment, and the
|
||||
/// agent pulls it from wherever it runs.
|
||||
MintAgentForgeToken { agent: String },
|
||||
/// Declare `agent` on `hive` as `Paused` in the swarm's wanted-state
|
||||
/// store, so a freshly created agent does not start driving turns the
|
||||
/// moment it's deployed — the operator has to explicitly flip it to `Up`.
|
||||
|
|
@ -130,6 +138,7 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
|
|||
SwarmNodeKind::AddRepoMember { .. } => "add_repo_member".to_owned(),
|
||||
SwarmNodeKind::InitAgentConfigRepo { .. } => "init_agent_config_repo".to_owned(),
|
||||
SwarmNodeKind::MintAgentIdentity { .. } => "mint_agent_identity".to_owned(),
|
||||
SwarmNodeKind::MintAgentForgeToken { .. } => "mint_agent_forge_token".to_owned(),
|
||||
SwarmNodeKind::SetAgentWanted { .. } => "set_agent_wanted".to_owned(),
|
||||
SwarmNodeKind::TriggerDeploy { .. } => "trigger_deploy".to_owned(),
|
||||
}
|
||||
|
|
@ -147,7 +156,8 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
|
|||
| SwarmNodeKind::CreateRepo { agent }
|
||||
| SwarmNodeKind::CreateForgeUser { agent }
|
||||
| SwarmNodeKind::AddRepoMember { agent }
|
||||
| SwarmNodeKind::InitAgentConfigRepo { agent } => {
|
||||
| SwarmNodeKind::InitAgentConfigRepo { agent }
|
||||
| SwarmNodeKind::MintAgentForgeToken { agent } => {
|
||||
serde_json::json!({ "agent": agent })
|
||||
}
|
||||
SwarmNodeKind::TriggerDeploy { hive, agent }
|
||||
|
|
@ -308,6 +318,7 @@ async fn run_swarm_node(
|
|||
}
|
||||
}
|
||||
},
|
||||
SwarmNodeKind::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await,
|
||||
SwarmNodeKind::SetAgentWanted { hive, agent } => match deps.wanted {
|
||||
None => Outcome::Failed(
|
||||
"no swarm queue is configured on this host, so no wanted-state \
|
||||
|
|
@ -327,6 +338,27 @@ async fn run_swarm_node(
|
|||
(builder, outcome)
|
||||
}
|
||||
|
||||
/// The `MintAgentForgeToken` arm, lifted out so `run_swarm_node` stays under
|
||||
/// `clippy::too_many_lines`.
|
||||
async fn mint_forge_token(
|
||||
forge: Option<Arc<forge::Client>>,
|
||||
agent: &str,
|
||||
) -> hive_jobq::scheduler::Outcome {
|
||||
use hive_jobq::scheduler::Outcome;
|
||||
|
||||
let Some(client) = forge else {
|
||||
return Outcome::Failed(
|
||||
"no forge configured on this host (SWARM_CONTROLLER_FORGE_URL / \
|
||||
SWARM_CONTROLLER_FORGE_TOKEN_FILE unset), so no agent forge token can be minted"
|
||||
.to_owned(),
|
||||
);
|
||||
};
|
||||
match client.ensure_agent_forge_token(agent).await {
|
||||
Ok(()) => Outcome::Done,
|
||||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
}
|
||||
}
|
||||
|
||||
/// Declare a brand-new agent at [`NEW_AGENT_WANTED_STATE`] — unless it turns
|
||||
/// out not to be new: an agent that already has a declaration (other than
|
||||
/// `Destroyed`, which this treats as reusable) is left alone, so a retried
|
||||
|
|
@ -1470,6 +1502,14 @@ fn declare_agent_job(
|
|||
agent: agent.to_owned(),
|
||||
})
|
||||
.after_ok(create_identity);
|
||||
// The agent's own forge token, in the store before the hive deploys
|
||||
// the container that pulls it. After the forge user, because the token
|
||||
// is minted for that user.
|
||||
let mint_forge_token = b
|
||||
.node(SwarmNodeKind::MintAgentForgeToken {
|
||||
agent: agent.to_owned(),
|
||||
})
|
||||
.after_ok(create_forge_user);
|
||||
// Declared before the deploy trigger so the pause is visible in the
|
||||
// wanted-state store before the hive brings the container up — see the
|
||||
// node's own doc comment for why "before", not just "eventually". Needs
|
||||
|
|
@ -1491,6 +1531,11 @@ fn declare_agent_job(
|
|||
// create agents exactly as it does today. `after_ok` there would turn an
|
||||
// unconfigured option into an agent nobody runs.
|
||||
//
|
||||
// The forge-token mint gets `after_any` for the same reason: a host with
|
||||
// no forge or no store must still create agents; only that node fails,
|
||||
// by name, and the agent's container fetches the token on its next
|
||||
// timer tick once one is minted.
|
||||
//
|
||||
// `after_ok` on the declaration, though: an agent deployed without its
|
||||
// pause landing first is the race this node exists to close, so a
|
||||
// declaration that did not land must not be deployed past. A host with no
|
||||
|
|
@ -1504,6 +1549,7 @@ fn declare_agent_job(
|
|||
})
|
||||
.after_ok(init_config)
|
||||
.after_any(mint_identity)
|
||||
.after_any(mint_forge_token)
|
||||
.after_ok(set_wanted);
|
||||
vec![create_identity.guid()]
|
||||
}
|
||||
|
|
@ -1703,6 +1749,73 @@ async fn mint_agent_identity(
|
|||
Ok(Json(MintAgentIdentityResponse { node_id: id.get() }))
|
||||
}
|
||||
|
||||
/// Success body of `POST /api/agents/{name}/forge-token`.
|
||||
#[derive(Clone, Debug, Serialize, ToSchema)]
|
||||
struct MintAgentForgeTokenResponse {
|
||||
/// The queued node, so a caller can follow it in the job view.
|
||||
node_id: u64,
|
||||
}
|
||||
|
||||
/// Insert one `MintAgentForgeToken` node per agent, each its own job.
|
||||
///
|
||||
/// 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.
|
||||
fn queue_forge_token_mints(
|
||||
sched: &Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>,
|
||||
agents: Vec<String>,
|
||||
) -> Result<Vec<hive_jobq::NodeId>> {
|
||||
let mut sched = sched
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
let mut ids = Vec::with_capacity(agents.len());
|
||||
for agent in agents {
|
||||
let queued = sched
|
||||
.insert_job(None, |b| {
|
||||
vec![b.node(SwarmNodeKind::MintAgentForgeToken { agent }).guid()]
|
||||
})
|
||||
.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
ids.extend(queued);
|
||||
}
|
||||
Ok(ids)
|
||||
}
|
||||
|
||||
/// Check an agent's forge token now, and mint one if it is missing or stale.
|
||||
///
|
||||
/// The periodic pass (`forge::agent_token::spawn`) does the same every five
|
||||
/// minutes for every agent with a store identity; this is for an operator who
|
||||
/// does not want to wait, or for an agent the pass skipped. Queues the same
|
||||
/// node agent creation does. Idempotent: a current token is left alone.
|
||||
#[utoipa::path(
|
||||
post,
|
||||
path = "/api/agents/{name}/forge-token",
|
||||
params(("name" = String, Path, description = "agent name")),
|
||||
responses(
|
||||
(status = 200, description = "mint queued", body = MintAgentForgeTokenResponse),
|
||||
(status = 400, description = "`name` is not a valid identifier (problem+json)", body = String),
|
||||
(status = 500, description = "the job could not be queued (problem+json)", body = String),
|
||||
),
|
||||
tag = "agents"
|
||||
)]
|
||||
async fn mint_agent_forge_token(
|
||||
State(state): State<AppState>,
|
||||
Path(name): Path<String>,
|
||||
) -> Result<Json<MintAgentForgeTokenResponse>, problem_details::ProblemDetails> {
|
||||
let agent = hive_types::Ident::parse(&name)
|
||||
.map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))?
|
||||
.into_string();
|
||||
let ids = queue_forge_token_mints(&state.jobq, vec![agent]).map_err(|e| {
|
||||
error_problem(
|
||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
||||
&e.to_string(),
|
||||
)
|
||||
})?;
|
||||
let [id] = ids[..] else {
|
||||
unreachable!("exactly one handle was asked for");
|
||||
};
|
||||
Ok(Json(MintAgentForgeTokenResponse { node_id: id.get() }))
|
||||
}
|
||||
|
||||
/// Every agent with an open config PR, in one response — the bulk
|
||||
/// counterpart to [`get_agent_config_pr`]. swarm-ui's config-PR table needs
|
||||
/// every agent's status to render, and fetching them one at a time doesn't
|
||||
|
|
@ -1990,6 +2103,14 @@ async fn main() -> Result<()> {
|
|||
);
|
||||
}
|
||||
|
||||
if let Some(client) = forge_client.clone() {
|
||||
let sched = Arc::clone(&jobq);
|
||||
forge::agent_token::spawn(client, move |agents| {
|
||||
if let Err(e) = queue_forge_token_mints(&sched, agents) {
|
||||
tracing::warn!(error = %format!("{e:#}"), "agent forge tokens: queueing failed");
|
||||
}
|
||||
});
|
||||
}
|
||||
let config_prs = forge_client.clone().map(config_pr::spawn);
|
||||
let state_forge = keep_forge_for_state(forge_client, webhook_secret.clone());
|
||||
|
||||
|
|
@ -2041,6 +2162,7 @@ fn build_app(state: AppState) -> axum::Router {
|
|||
.routes(routes!(get_config_prs))
|
||||
.routes(routes!(create_agent))
|
||||
.routes(routes!(mint_agent_identity))
|
||||
.routes(routes!(mint_agent_forge_token))
|
||||
.routes(routes!(get_agents))
|
||||
.routes(routes!(get_agents_status))
|
||||
.routes(routes!(set_agent_state))
|
||||
|
|
@ -2842,6 +2964,105 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
/// The forge-token mint sits between the forge user and the deploy: it
|
||||
/// needs the user to exist, and the deploy must not overtake it — but a
|
||||
/// host with no forge or no store must still deploy the agent.
|
||||
#[tokio::test]
|
||||
async fn the_forge_token_mint_follows_the_user_and_does_not_block_the_deploy() {
|
||||
use hive_jobq_wire::WireNode as _;
|
||||
|
||||
let (state, sched) = state_with_roster();
|
||||
let _queued = super::create_agent(
|
||||
axum::extract::State(state),
|
||||
axum::Json(super::CreateAgentRequest {
|
||||
name: "atlas".to_owned(),
|
||||
hive: "pr1ma".to_owned(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.expect("a hive in the roster must be accepted");
|
||||
|
||||
let guard = sched
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
let graph = guard.graph();
|
||||
let id_of = |label: &str| {
|
||||
graph
|
||||
.nodes()
|
||||
.find(|n| n.payload.label() == label)
|
||||
.unwrap_or_else(|| panic!("the graph holds a {label} node"))
|
||||
.id
|
||||
};
|
||||
let edge = |from: hive_jobq::NodeId, to: &str| {
|
||||
graph
|
||||
.node(id_of(to))
|
||||
.expect("just found")
|
||||
.deps
|
||||
.iter()
|
||||
.find_map(|d| match d {
|
||||
hive_jobq::Dep::Node { id, when } if *id == from => Some(*when),
|
||||
_ => None,
|
||||
})
|
||||
.unwrap_or_else(|| panic!("{to} waits for node {}", from.get()))
|
||||
};
|
||||
|
||||
let user_to_mint = edge(id_of("create_forge_user"), "mint_agent_forge_token");
|
||||
assert!(
|
||||
!user_to_mint.accepts(hive_jobq::TerminalState::Failed),
|
||||
"a token for a user that was never created cannot be minted; this \
|
||||
edge has to be `after_ok`"
|
||||
);
|
||||
let mint_to_deploy = edge(id_of("mint_agent_forge_token"), "trigger_deploy");
|
||||
assert!(
|
||||
mint_to_deploy.accepts(hive_jobq::TerminalState::Failed),
|
||||
"a host with no forge or store must not cancel the deploy; this edge \
|
||||
has to be `after_any`, not `after_ok`"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_forge_token_node_renders_the_agent() {
|
||||
use hive_jobq_wire::WireNode as _;
|
||||
|
||||
let kind = SwarmNodeKind::MintAgentForgeToken {
|
||||
agent: "atlas".to_owned(),
|
||||
};
|
||||
assert_eq!(kind.label(), "mint_agent_forge_token");
|
||||
assert_eq!(kind.data(1)["agent"], "atlas");
|
||||
}
|
||||
|
||||
/// 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.
|
||||
#[test]
|
||||
fn a_queued_mint_is_one_forge_token_node_per_agent() {
|
||||
use hive_jobq_wire::WireNode as _;
|
||||
|
||||
let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new(
|
||||
hive_jobq::Graph::new(),
|
||||
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);
|
||||
let guard = sched
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
let mut agents: Vec<String> = 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"]);
|
||||
}
|
||||
|
||||
/// The socket must not share a directory with anything else, because
|
||||
/// the socket is `0666` and the directory is therefore the only access
|
||||
/// control it has. `/run/hyperhive` in particular holds hive-c0re's
|
||||
|
|
|
|||
Loading…
Reference in a new issue