hyperhive/swarmctl/src/agent.rs
atlas d56d8f2b36 swarmctl: forge make-admin
Calls POST /api/forge/users/{name}/admin and prints what it found. It
fails with the controller's message when the user has not logged in via
SSO yet, and when the name is an agent's.

Refs #3782
2026-09-25 08:29:56 +02:00

351 lines
14 KiB
Rust

//! `swarmctl agent create`, `swarmctl agent mint-identity` and `swarmctl agent
//! mint-forge-token` — queue work on the swarm-controller's job graph.
//!
//! Each POSTs and returns as soon as the work is *inserted*. Both verbs
//! print the queued node id and stop: creation's last node only *publishes*
//! a deploy, after which the hive converges on its own clock, so even a
//! settled graph would not mean the agent is up and there is nothing here
//! that could be waited on honestly.
//!
//! Transport is a bare `hyper` HTTP/1.1 client handshaked onto a
//! [`tokio::net::UnixStream`] via [`hyper_util::rt::TokioIo`] — the
//! controller has no TCP port, so this crate speaks plain HTTP directly
//! over the socket rather than pulling in a full client stack for one
//! POST.
//!
//! The two shapes below **mirror** the controller's own private
//! `CreateAgentRequest` / `CreateAgentResponse` rather than sharing them:
//! this crate does not link `swarm-controller` and there is no wire
//! crate between them. The seam is narrow and both ends validate, so
//! drift surfaces as a 400 naming the field.
//!
//! Long-form rationale for all three: `swarmctl/README.md`.
use std::path::Path;
use anyhow::{Context as _, Result, bail};
use http_body_util::{BodyExt as _, Full};
use hyper_util::rt::TokioIo;
use serde::{Deserialize, Serialize};
use tokio::net::UnixStream;
/// Body of `POST /api/agents`. Mirrors the controller's own
/// `CreateAgentRequest` — see this module's doc comment.
#[derive(Serialize)]
struct CreateAgentRequest<'a> {
name: &'a str,
/// Where the agent's deploy message is addressed. Required by the
/// endpoint and not defaulted anywhere, because it is only knowable
/// from the operator making the choice.
hive: &'a str,
}
/// Success body of `POST /api/agents`.
#[derive(Deserialize)]
struct CreateAgentResponse {
node_id: u64,
/// Name collisions the controller **allowed through** — it is the
/// "warn now, refuse later" step, so an empty list is the normal case
/// and the field is omitted from the JSON entirely when it is empty.
#[serde(default)]
warnings: Vec<String>,
}
/// Body of `POST /api/agents/{name}/identity`. Mirrors the controller's own
/// `MintAgentIdentityRequest` — see this module's doc comment.
#[derive(Serialize)]
struct MintIdentityRequest<'a> {
hive: &'a str,
}
/// Success body of `POST /api/agents/{name}/identity`.
#[derive(Deserialize)]
struct MintIdentityResponse {
node_id: u64,
}
/// Success body of `POST /api/agents/{name}/forge-token`, which takes no body.
#[derive(Deserialize)]
struct MintForgeTokenResponse {
node_id: u64,
}
/// Run `swarmctl agent create`.
///
/// Synchronous on purpose: every other verb in this crate is, and this is
/// the only one that needs a reactor at all. A current-thread runtime
/// built here keeps that cost inside the one arm that incurs it, the same
/// way `PathArgs::resolve` is called only in the arms that need the
/// deployment's env vars.
pub(crate) fn create(socket: &Path, name: &str, hive: &str) -> Result<()> {
// Client-side before the round trip, so a typo is an immediate local
// error rather than a 400 the operator has to wait for. The controller
// validates both again — that is the gate that matters, this is the
// one that is fast.
let name = parse_ident(name, "agent name")?;
let hive = parse_ident(hive, "hive")?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
.build()
.context("starting a tokio runtime for the controller request")?;
let resp = rt.block_on(post_create(socket, &name, &hive))?;
println!("queued: job node {}", resp.node_id);
println!(
"agent {name:?} will be deployed to hive {hive:?} once the job graph runs; \
`swarmctl` does not wait for it"
);
// stderr, and after the id: the id is the result, these are asides the
// operator should still not have to go find in the daemon's journal.
for warning in &resp.warnings {
eprintln!("warning: {warning}");
}
Ok(())
}
/// Parse a CLI-supplied name into a validated identifier, naming which
/// argument was wrong — `invalid hive` and `invalid agent name` send the
/// operator to different flags.
pub(crate) fn parse_ident(value: &str, what: &str) -> Result<String> {
hive_types::Ident::parse(value)
.map(hive_types::Ident::into_string)
.map_err(|reason| anyhow::anyhow!("invalid {what} {value:?}: {reason}"))
}
/// Run `swarmctl agent mint-identity`.
///
/// Synchronous for the same reason [`create`] is, and built on the same
/// round trip.
pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()> {
// Client-side first, so a typo is a local error rather than a 400 the
// operator waits for. The controller validates both again.
let name = parse_ident(name, "agent name")?;
let hive = parse_ident(hive, "hive")?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
.build()
.context("starting a tokio runtime for the controller request")?;
let resp: MintIdentityResponse = rt.block_on(post(
socket,
&format!("/api/agents/{name}/identity"),
&MintIdentityRequest { hive: &hive },
"mint-identity",
))?;
println!("queued: job node {}", resp.node_id);
println!(
"agent {name:?} will have its identity re-minted on hive {hive:?} once the job graph \
runs; `swarmctl` does not wait for it. An existing queue secret is kept as it is; the \
store certificate is re-minted and the agent picks the new one up on its next boot"
);
Ok(())
}
/// Run `swarmctl agent mint-forge-token`.
///
/// Synchronous for the same reason [`create`] is, and built on the same
/// round trip. No `--hive`: the token's store path has no hive in it.
pub(crate) fn mint_forge_token(socket: &Path, name: &str) -> Result<()> {
let name = parse_ident(name, "agent name")?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
.build()
.context("starting a tokio runtime for the controller request")?;
let resp: MintForgeTokenResponse = rt.block_on(post(
socket,
&format!("/api/agents/{name}/forge-token"),
&serde_json::json!({}),
"mint-forge-token",
))?;
println!("queued: job node {}", resp.node_id);
println!(
"agent {name:?}'s forge token will be checked once the job graph runs, and minted \
only if it is missing or stale; `swarmctl` does not wait for it. The agent picks a \
new token up on its next fetch"
);
Ok(())
}
/// One `POST /api/agents` round trip over the controller's unix socket.
async fn post_create(socket: &Path, name: &str, hive: &str) -> Result<CreateAgentResponse> {
post(
socket,
"/api/agents",
&CreateAgentRequest { name, hive },
"create-agent",
)
.await
}
/// One JSON `POST` to `uri` over the controller's unix socket.
///
/// `what` names the request in error messages — the operator needs to know
/// which call failed, and every other part of this function is identical
/// between the two verbs.
pub(crate) async fn post<Req: Serialize, Resp: serde::de::DeserializeOwned>(
socket: &Path,
uri: &str,
request: &Req,
what: &str,
) -> Result<Resp> {
let body =
serde_json::to_vec(request).with_context(|| format!("serialising the {what} request"))?;
let stream = UnixStream::connect(socket)
.await
.with_context(|| connect_hint(socket))?;
let (mut sender, conn) = hyper::client::conn::http1::handshake(TokioIo::new(stream))
.await
.with_context(|| format!("HTTP handshake with {}", socket.display()))?;
// Drives the connection's I/O; a one-shot POST has nothing to do with
// the join handle, and a connection that dies surfaces as an error on
// `send_request` or on reading the body.
tokio::spawn(conn);
let req = hyper::Request::builder()
.method(hyper::Method::POST)
.uri(uri)
// A unix socket has no authority of its own, but HTTP/1.1 requires
// the header; the controller routes on the path alone.
.header(hyper::header::HOST, "localhost")
.header(hyper::header::CONTENT_TYPE, "application/json")
.body(Full::new(bytes::Bytes::from(body)))
.with_context(|| format!("building the {what} request"))?;
let resp = sender
.send_request(req)
.await
.with_context(|| format!("sending the {what} request to swarm-controller"))?;
let status = resp.status();
let body = resp
.into_body()
.collect()
.await
.context("reading swarm-controller's response")?
.to_bytes();
if !status.is_success() {
bail!("{}", describe_error(status, &body));
}
serde_json::from_slice(&body).with_context(|| {
format!(
"swarm-controller answered {status} but its body is not a {what} response: {}",
String::from_utf8_lossy(&body).trim()
)
})
}
/// Turn a failed response into one line the operator can act on.
///
/// The controller answers errors as RFC 9457 `application/problem+json`,
/// where the actionable sentence is `detail` — a 400 for an unknown hive
/// names the hives that *would* have worked. Printing the raw body instead
/// would bury that in JSON punctuation, and printing only the status would
/// throw it away. Falls back through `title` to the raw body so a response
/// from something that is *not* the controller is still shown rather than
/// reduced to a bare number.
fn describe_error(status: hyper::StatusCode, body: &[u8]) -> String {
let text = String::from_utf8_lossy(body);
let text = text.trim();
let field = |v: &serde_json::Value, k: &str| {
v.get(k)
.and_then(serde_json::Value::as_str)
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_owned)
};
let detail = serde_json::from_str::<serde_json::Value>(text)
.ok()
.and_then(|v| field(&v, "detail").or_else(|| field(&v, "title")));
match detail {
Some(detail) => format!("swarm-controller refused the request ({status}): {detail}"),
None if text.is_empty() => format!("swarm-controller refused the request ({status})"),
None => format!("swarm-controller refused the request ({status}): {text}"),
}
}
/// Message for a socket that could not be dialled.
///
/// Worth spelling out because the three ways this fails need three
/// different fixes and `No such file or directory (os error 2)` names
/// none of them: the daemon is off, it binds elsewhere, or this host is
/// not the one running it at all. No group-membership case to mention
/// here — `swarmctl` already runs as root.
fn connect_hint(socket: &Path) -> String {
format!(
"could not connect to swarm-controller at {} — is it running on this host? \
(`systemctl status swarm-controller`). The path comes from \
`--controller-socket` or `SWARM_CONTROLLER_SOCKET`, which the \
swarm-controller nix module sets from `socketPath`",
socket.display()
)
}
#[cfg(test)]
mod tests {
use super::{CreateAgentResponse, describe_error, parse_ident};
/// The endpoint omits `warnings` entirely when there are none, so the
/// normal response must still decode — a missing field here would make
/// every successful creation fail to parse.
#[test]
fn a_response_without_warnings_decodes() {
let resp: CreateAgentResponse =
serde_json::from_str(r#"{"node_id":7}"#).expect("the normal shape decodes");
assert_eq!(resp.node_id, 7);
assert!(resp.warnings.is_empty());
}
#[test]
fn warnings_are_carried_through() {
let resp: CreateAgentResponse =
serde_json::from_str(r#"{"node_id":9,"warnings":["reserved name"]}"#).expect("decodes");
assert_eq!(resp.warnings, vec!["reserved name".to_owned()]);
}
/// The whole point of parsing problem+json: the 400 for an unknown
/// hive carries the roster of hives that would have worked, and that
/// sentence has to reach the operator's terminal intact.
#[test]
fn a_problem_json_detail_reaches_the_operator() {
let body = r#"{"status":400,"title":"Bad Request","detail":"hive \"tyop\" is not in this swarm — known hives: alpha, beta"}"#;
let msg = describe_error(hyper::StatusCode::BAD_REQUEST, body.as_bytes());
assert!(msg.contains("known hives: alpha, beta"), "{msg}");
assert!(msg.contains("400"), "{msg}");
// No JSON punctuation: the operator gets prose, not a body dump.
assert!(!msg.contains('{'), "{msg}");
}
/// A response from something that is not the controller (an nginx
/// error page, a proxy) must still be shown rather than swallowed.
#[test]
fn a_non_json_body_is_shown_verbatim() {
let msg = describe_error(hyper::StatusCode::BAD_GATEWAY, b" 502 Bad Gateway ");
assert!(msg.contains("502 Bad Gateway"), "{msg}");
}
#[test]
fn an_empty_body_still_names_the_status() {
let msg = describe_error(hyper::StatusCode::INTERNAL_SERVER_ERROR, b"");
assert!(msg.contains("500"), "{msg}");
}
/// Which argument was wrong has to be in the message: `agent create`
/// takes two identifiers and they are fixed in different places.
#[test]
fn an_invalid_ident_names_the_argument_it_came_from() {
let err = parse_ident("Not An Ident", "hive").expect_err("rejected");
let msg = err.to_string();
assert!(msg.contains("invalid hive"), "{msg}");
assert!(msg.contains("Not An Ident"), "{msg}");
assert!(
parse_ident("scribe-01", "agent name").is_ok(),
"a normal name must pass"
);
}
}