forge-token.nix fetches swarm/agents/<agent>/forge-token under the agent's own store identity into /run/hive-agent-forge-token/token, and re-fetches on a timer so a rotation lands. hive-forge, the git credential helper, hive-forge-notify, forge-avatar-sync and the web UI read that file first and fall back to <state>/forge-token. tea-login is deleted: it copied the token into ~/.config/tea, which docs/swarm/credentials.md forbids for a store secret. hive-forge covers the same verbs. swarmctl gains agent mint-forge-token. Refs #3782
351 lines
14 KiB
Rust
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.
|
|
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.
|
|
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"
|
|
);
|
|
}
|
|
}
|