swarm-controller: mint each agent's matrix account with the swarm's token
A `MintAgentMatrixAccount` node creates the agent's account on the swarm's homeserver with the swarm appservice token, stores its token at `swarm/agents/<agent>/matrix/main`, and reads it back with whoami before reporting success. It is a root of agent creation, `after_any` into the deploy, and a five-minute backfill over every agent with a store identity queues the same node — the shape of the forge-token mint. The decision reads the stored token back rather than only checking that one is stored: the swarm and a hive both pin the device `hyperhive-<agent>`, so each login replaces the other's token. A failed read plans nothing, so an outage never rotates every agent's token. `matrixHomeserverUrl` now defaults to the swarm's `chat.` vhost, since the mint is what consults it.
This commit is contained in:
parent
9308752a09
commit
2776e121e5
7 changed files with 626 additions and 36 deletions
|
|
@ -113,6 +113,14 @@ enum SwarmNodeKind {
|
|||
/// Carries no hive: the token's store path has no hive segment, and the
|
||||
/// agent pulls it from wherever it runs.
|
||||
MintAgentForgeToken { agent: String },
|
||||
/// Make sure `agent` holds a live token for its own account on the swarm's
|
||||
/// homeserver, creating the account with the swarm's appservice token
|
||||
/// when it does not. See `matrix_account::agent_token` — including why a
|
||||
/// stored token is read back before anything is minted.
|
||||
///
|
||||
/// Carries no hive: the path it writes has no hive segment, so an agent
|
||||
/// that moves between hives keeps one matrix identity.
|
||||
MintAgentMatrixAccount { 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`.
|
||||
|
|
@ -139,6 +147,7 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
|
|||
SwarmNodeKind::InitAgentConfigRepo { .. } => "init_agent_config_repo".to_owned(),
|
||||
SwarmNodeKind::MintAgentIdentity { .. } => "mint_agent_identity".to_owned(),
|
||||
SwarmNodeKind::MintAgentForgeToken { .. } => "mint_agent_forge_token".to_owned(),
|
||||
SwarmNodeKind::MintAgentMatrixAccount { .. } => "mint_agent_matrix_account".to_owned(),
|
||||
SwarmNodeKind::SetAgentWanted { .. } => "set_agent_wanted".to_owned(),
|
||||
SwarmNodeKind::TriggerDeploy { .. } => "trigger_deploy".to_owned(),
|
||||
}
|
||||
|
|
@ -157,7 +166,8 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
|
|||
| SwarmNodeKind::CreateForgeUser { agent }
|
||||
| SwarmNodeKind::AddRepoMember { agent }
|
||||
| SwarmNodeKind::InitAgentConfigRepo { agent }
|
||||
| SwarmNodeKind::MintAgentForgeToken { agent } => {
|
||||
| SwarmNodeKind::MintAgentForgeToken { agent }
|
||||
| SwarmNodeKind::MintAgentMatrixAccount { agent } => {
|
||||
serde_json::json!({ "agent": agent })
|
||||
}
|
||||
SwarmNodeKind::TriggerDeploy { hive, agent }
|
||||
|
|
@ -209,6 +219,11 @@ struct WorkerDeps {
|
|||
/// see `wanted_writer`. `None` exactly when no swarm queue is configured
|
||||
/// on this host, same as `queue`.
|
||||
wanted: Option<std::sync::Arc<wanted::WantedWriter>>,
|
||||
/// Client-server API base of the swarm's homeserver, for the node that
|
||||
/// mints agents' accounts on it. `None` when
|
||||
/// `matrix_account::DEFAULT_HOMESERVER_ENV` is unset, which that node
|
||||
/// **fails** on rather than skips, same as every other field here.
|
||||
matrix_homeserver: Option<std::sync::Arc<str>>,
|
||||
}
|
||||
|
||||
/// What [`SwarmNodeKind::SetAgentWanted`] declares a brand-new agent to be.
|
||||
|
|
@ -304,21 +319,13 @@ async fn run_swarm_node(
|
|||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
},
|
||||
},
|
||||
SwarmNodeKind::MintAgentIdentity { hive, agent } => match deps.agent_ca {
|
||||
None => Outcome::Failed(
|
||||
"no agent certificate authority is configured on this host \
|
||||
(SWARM_CONTROLLER_AGENT_CA_FILE / SWARM_CONTROLLER_AGENT_CA_KEY_FILE unset), \
|
||||
so this agent has no identity at the swarm secret store"
|
||||
.to_owned(),
|
||||
),
|
||||
Some(authority) => {
|
||||
match agent_identity::mint_and_verify(&authority, &agent, &hive).await {
|
||||
Ok(()) => Outcome::Done,
|
||||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
}
|
||||
}
|
||||
},
|
||||
SwarmNodeKind::MintAgentIdentity { hive, agent } => {
|
||||
mint_identity(deps.agent_ca.as_deref(), &agent, &hive).await
|
||||
}
|
||||
SwarmNodeKind::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await,
|
||||
SwarmNodeKind::MintAgentMatrixAccount { agent } => {
|
||||
mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await
|
||||
}
|
||||
SwarmNodeKind::SetAgentWanted { hive, agent } => match deps.wanted {
|
||||
None => Outcome::Failed(
|
||||
"no swarm queue is configured on this host, so no wanted-state \
|
||||
|
|
@ -338,6 +345,29 @@ async fn run_swarm_node(
|
|||
(builder, outcome)
|
||||
}
|
||||
|
||||
/// The `MintAgentIdentity` arm, lifted out so `run_swarm_node` stays under
|
||||
/// `clippy::too_many_lines`.
|
||||
async fn mint_identity(
|
||||
authority: Option<&agent_identity::Authority>,
|
||||
agent: &str,
|
||||
hive: &str,
|
||||
) -> hive_jobq::scheduler::Outcome {
|
||||
use hive_jobq::scheduler::Outcome;
|
||||
|
||||
let Some(authority) = authority else {
|
||||
return Outcome::Failed(
|
||||
"no agent certificate authority is configured on this host \
|
||||
(SWARM_CONTROLLER_AGENT_CA_FILE / SWARM_CONTROLLER_AGENT_CA_KEY_FILE unset), \
|
||||
so this agent has no identity at the swarm secret store"
|
||||
.to_owned(),
|
||||
);
|
||||
};
|
||||
match agent_identity::mint_and_verify(authority, agent, hive).await {
|
||||
Ok(()) => Outcome::Done,
|
||||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
}
|
||||
}
|
||||
|
||||
/// The `MintAgentForgeToken` arm, lifted out so `run_swarm_node` stays under
|
||||
/// `clippy::too_many_lines`.
|
||||
async fn mint_forge_token(
|
||||
|
|
@ -359,6 +389,28 @@ async fn mint_forge_token(
|
|||
}
|
||||
}
|
||||
|
||||
/// The `MintAgentMatrixAccount` arm, lifted out for the same reason.
|
||||
async fn mint_matrix_account(
|
||||
homeserver: Option<&str>,
|
||||
agent: &str,
|
||||
) -> hive_jobq::scheduler::Outcome {
|
||||
use hive_jobq::scheduler::Outcome;
|
||||
|
||||
// Loud on purpose: an agent quietly created without a matrix account is
|
||||
// an agent nothing can talk to, and no hive mints one any more.
|
||||
let Some(base) = homeserver else {
|
||||
return Outcome::Failed(format!(
|
||||
"no matrix homeserver configured on this host ({} unset), so no agent \
|
||||
matrix account can be minted",
|
||||
matrix_account::DEFAULT_HOMESERVER_ENV
|
||||
));
|
||||
};
|
||||
match matrix_account::agent_token::ensure_agent_matrix_account(base, 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
|
||||
|
|
@ -1513,6 +1565,13 @@ fn declare_agent_job(
|
|||
agent: agent.to_owned(),
|
||||
})
|
||||
.after_ok(create_forge_user);
|
||||
// The agent's account on the swarm's homeserver. A root of its own:
|
||||
// creating it needs neither an authelia subject nor a forge user nor a
|
||||
// repo, and chaining it behind one of those would make an unrelated
|
||||
// failure look like a matrix failure.
|
||||
let mint_matrix = b.node(SwarmNodeKind::MintAgentMatrixAccount {
|
||||
agent: agent.to_owned(),
|
||||
});
|
||||
// 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
|
||||
|
|
@ -1553,6 +1612,11 @@ fn declare_agent_job(
|
|||
.after_ok(init_config)
|
||||
.after_any(mint_identity)
|
||||
.after_any(mint_forge_token)
|
||||
// `after_any` for the same reason: the container reads this token
|
||||
// rather than producing it, so the deploy must not overtake the mint —
|
||||
// but a host with no homeserver configured must still create agents,
|
||||
// and only this node fails, by name.
|
||||
.after_any(mint_matrix)
|
||||
.after_ok(set_wanted);
|
||||
vec![create_identity.guid()]
|
||||
}
|
||||
|
|
@ -1795,6 +1859,57 @@ fn queue_forge_token_mints(
|
|||
Ok(ids)
|
||||
}
|
||||
|
||||
/// Insert one `MintAgentMatrixAccount` job per agent and return the nodes'
|
||||
/// ids. `matrix_account::agent_token::spawn`'s periodic pass comes through
|
||||
/// here, so a backfilled mint is the same node a new agent gets.
|
||||
fn queue_matrix_account_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::MintAgentMatrixAccount { agent })
|
||||
.guid(),
|
||||
]
|
||||
})
|
||||
.map_err(|e| anyhow::anyhow!("{e}"))?;
|
||||
ids.extend(queued);
|
||||
}
|
||||
Ok(ids)
|
||||
}
|
||||
|
||||
/// The homeserver agents' own matrix accounts are minted on. Unset is this
|
||||
/// controller creating agents with no matrix account, which the mint node
|
||||
/// reports by name rather than `main` refusing to start.
|
||||
fn configured_matrix_homeserver() -> Option<Arc<str>> {
|
||||
matrix_account::configured_default_homeserver()
|
||||
.filter(|v| !v.is_empty())
|
||||
.map(Arc::from)
|
||||
}
|
||||
|
||||
/// Start the matrix-account backfill when a homeserver is configured. Lifted
|
||||
/// out of `main` for `clippy::too_many_lines`.
|
||||
fn spawn_matrix_account_backfill(
|
||||
jobq: &Arc<Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>>,
|
||||
homeserver: Option<Arc<str>>,
|
||||
) {
|
||||
let Some(base) = homeserver else {
|
||||
return;
|
||||
};
|
||||
let sched = Arc::clone(jobq);
|
||||
matrix_account::agent_token::spawn(base, move |agents| {
|
||||
if let Err(e) = queue_matrix_account_mints(&sched, agents) {
|
||||
tracing::warn!(error = %format!("{e:#}"), "agent matrix accounts: queueing failed");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// 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
|
||||
|
|
@ -2167,12 +2282,14 @@ async fn main() -> Result<()> {
|
|||
queue: status.as_ref().map(|s| s.queue_client()),
|
||||
agent_ca: load_agent_authority(),
|
||||
wanted: wanted_writer(status.as_ref()),
|
||||
matrix_homeserver: configured_matrix_homeserver(),
|
||||
};
|
||||
|
||||
let jobq = Arc::new(Mutex::new(hive_jobq::scheduler::Scheduler::new(
|
||||
hive_jobq::Graph::new(),
|
||||
hive_jobq::resources::ResourceTable::new(),
|
||||
)));
|
||||
spawn_matrix_account_backfill(&jobq, deps.matrix_homeserver.clone());
|
||||
spawn_jobq_worker(Arc::clone(&jobq), deps);
|
||||
// Bound to a named variable, not `_` — dropping the provider stops its
|
||||
// `PeriodicReader`, so it must live as long as `main` does (which it
|
||||
|
|
@ -2844,6 +2961,7 @@ mod tests {
|
|||
queue: None,
|
||||
agent_ca: None,
|
||||
wanted: None,
|
||||
matrix_homeserver: None,
|
||||
};
|
||||
let runner =
|
||||
hive_jobq::scheduler::Scheduler::claim_next(&sched, move |id, kind, builder| {
|
||||
|
|
@ -2898,6 +3016,7 @@ mod tests {
|
|||
queue: None,
|
||||
agent_ca: None,
|
||||
wanted: None,
|
||||
matrix_homeserver: None,
|
||||
};
|
||||
let runner =
|
||||
hive_jobq::scheduler::Scheduler::claim_next(&sched, move |id, kind, builder| {
|
||||
|
|
@ -2952,6 +3071,7 @@ mod tests {
|
|||
queue: None,
|
||||
agent_ca: None,
|
||||
wanted: None,
|
||||
matrix_homeserver: None,
|
||||
};
|
||||
let runner =
|
||||
hive_jobq::scheduler::Scheduler::claim_next(&sched, move |id, kind, builder| {
|
||||
|
|
@ -3004,6 +3124,7 @@ mod tests {
|
|||
queue: None,
|
||||
agent_ca: None,
|
||||
wanted: None,
|
||||
matrix_homeserver: None,
|
||||
};
|
||||
let runner =
|
||||
hive_jobq::scheduler::Scheduler::claim_next(&sched, move |id, kind, builder| {
|
||||
|
|
@ -3242,6 +3363,92 @@ mod tests {
|
|||
assert_eq!(kind.data(1)["agent"], "atlas");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_matrix_account_node_renders_the_agent() {
|
||||
use hive_jobq_wire::WireNode as _;
|
||||
|
||||
let kind = SwarmNodeKind::MintAgentMatrixAccount {
|
||||
agent: "atlas".to_owned(),
|
||||
};
|
||||
assert_eq!(kind.label(), "mint_agent_matrix_account");
|
||||
assert_eq!(kind.data(1)["agent"], "atlas");
|
||||
}
|
||||
|
||||
/// Agent creation mints the matrix account as a root of its own, and the
|
||||
/// deploy waits for it without being cancelled by it: a host with no
|
||||
/// homeserver configured must still deploy the agent.
|
||||
#[tokio::test]
|
||||
async fn the_matrix_account_mint_is_a_root_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 find = |label: &str| {
|
||||
graph
|
||||
.nodes()
|
||||
.find(|n| n.payload.label() == label)
|
||||
.unwrap_or_else(|| panic!("the graph holds a {label} node"))
|
||||
};
|
||||
let mint = find("mint_agent_matrix_account");
|
||||
assert!(mint.deps.is_empty(), "the mint waits for nothing");
|
||||
let when = find("trigger_deploy")
|
||||
.deps
|
||||
.iter()
|
||||
.find_map(|d| match d {
|
||||
hive_jobq::Dep::Node { id, when } if *id == mint.id => Some(*when),
|
||||
_ => None,
|
||||
})
|
||||
.expect("the deploy waits for the matrix mint");
|
||||
assert!(
|
||||
when.accepts(hive_jobq::TerminalState::Failed),
|
||||
"a host with no homeserver must not cancel the deploy; this edge \
|
||||
has to be `after_any`, not `after_ok`"
|
||||
);
|
||||
}
|
||||
|
||||
/// The backfill's queueing: one mint node per agent, nothing else.
|
||||
#[test]
|
||||
fn a_queued_matrix_mint_is_one_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_matrix_account_mints(&sched, vec!["a".to_owned(), "b".to_owned()])
|
||||
.expect("two 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_matrix_account");
|
||||
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 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 pair of nodes agent creation inserts: the forge user,
|
||||
|
|
|
|||
Loading…
Reference in a new issue