hyperhive/swarm-controller/src/wanted.rs
iris 75af27abe0 swarm-controller: make agent-creation's wanted-state write idempotent
Intent::Create now no-ops when the agent already has a non-Destroyed
declaration, instead of unconditionally overwriting it to the
create path's state. Without this, re-running create against a name
that already has a wanted-state entry (an operator migrating a
pre-existing agent into this bookkeeping, or a retried request)
would silently pause an agent already running under some other
state. write() skips the network round-trip entirely when apply
returns the declaration unchanged.

mara: this does not match the expectation that agent creation is
idempotent so pre existing agents can be migrated
2026-09-18 22:59:29 +02:00

480 lines
22 KiB
Rust

//! Writes the agent set this swarm declares for each hive.
//!
//! The mirror of [`crate::status`]: that module reads what hives report,
//! this one writes what they are told, and both address the same queue.
//! Both hold a NATS client rather than a bucket handle, so a controller
//! that starts before a bucket exists picks it up without a restart.
//!
//! Where the mirror stops is the handle itself: `status` resolves one on
//! first use and caches it, while there is one wanted-state bucket **per
//! hive**, so no single handle serves them and `store` resolves per call.
//!
//! The bucket is the record. Nothing here keeps a second copy of the
//! declaration to reconcile against, because the current value can be read
//! back from the queue whenever it is needed.
use std::fmt;
use anyhow::{Context, Result};
use async_nats::jetstream::kv::{CreateErrorKind, UpdateErrorKind};
use swarm_queue_client::wanted::{AgentState, AgentWanted, HiveWanted};
/// The one failure out of [`apply`] that is the caller's mistake, not a
/// server fault — worth a distinct type so `set_agent_state` can tell it
/// apart from every other error this module raises (an undecodable
/// declaration, a lost write) and answer 400 instead of 500.
///
/// Carried through as a plain `anyhow::Error` like every other error here
/// (this crate has no typed-error convention to join), and recovered at the
/// HTTP boundary with `downcast_ref` — the shape every other "kind of error
/// that came from underneath" check in this module already uses
/// ([`Wrote`]'s `UpdateErrorKind`/`CreateErrorKind` matching), just inspecting
/// an `anyhow::Error` instead of a client library's own error enum.
#[derive(Debug)]
pub struct TerminalStateError {
agent: String,
}
impl fmt::Display for TerminalStateError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"{} is already declared destroyed — terminal, per AgentState::Destroyed's own doc \
comment: no transition brings it back short of a fresh deploy",
self.agent
)
}
}
impl std::error::Error for TerminalStateError {}
/// The three outcomes of one write attempt, which the two KV verbs report
/// through separate error types.
enum Wrote {
Ok,
/// Another writer won the race; re-read and re-apply.
LostRace,
Failed(anyhow::Error),
}
/// Why a caller is declaring a state, which is what decides whether a
/// `Destroyed` entry stands in its way.
///
/// The distinction is the one [`AgentState::Destroyed`]'s own doc comment
/// already promises — "no state that brings a destroyed agent back **short of
/// a fresh deploy**". Redeclaring is not a fresh deploy and must keep being
/// refused; creating the agent again *is* one, and the whole of its identity,
/// repo and config has just been built anew alongside this declaration.
///
/// An enum rather than a bool because the two call sites read as opposites
/// only if the argument says which one it is — `apply(.., true)` at the create
/// path would be a flag nobody can check at a glance.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Intent {
/// An operator (or anything else) changing what an existing agent should
/// be doing. Cannot move off `Destroyed`.
Redeclare,
/// The agent is being created from nothing — a fresh deploy. Replaces a
/// `Destroyed` entry left behind by an agent that had this name before.
Create,
}
/// How many times a losing writer re-reads and re-applies before giving up.
///
/// A conflict means another writer changed a *different* agent between this
/// one's read and its write, so a retry re-reads and re-applies onto the
/// winner. Bounded because an unbounded loop against a hot key is a spin,
/// and a caller that gets an error can ask again with fresh intent.
const MAX_ATTEMPTS: usize = 5;
/// Apply one agent's declared state to a hive's current declaration.
///
/// Split out from the write loop because it holds the invariant that matters:
/// the value under a hive's key is the map of **every** agent on that hive,
/// so declaring one agent must preserve the rest.
///
/// `current` is `None` when the hive has no declaration yet. An undecodable
/// value is an **error**, never treated as absent: overwriting a document
/// nobody can read discards the declarations of every other agent on that
/// hive, which is exactly what a fresh-start fallback would do quietly.
///
/// `intent` decides whether an existing `Destroyed` entry for this agent is a
/// wall or something to write over — see [`Intent`].
fn apply(
current: Option<&[u8]>,
agent: &str,
state: AgentState,
intent: Intent,
) -> Result<(HiveWanted, Vec<u8>)> {
let mut declaration = match current {
Some(raw) => serde_json::from_slice::<HiveWanted>(raw)
.context("the hive's current declaration is not decodable")?,
None => HiveWanted::default(),
};
// Redeclaring the same terminal state is a no-op, not a transition — a
// client that already sees the agent as destroyed and asks again should
// not be punished for it. Anything else moving off `Destroyed` is the
// one transition this module exists to refuse — unless the caller is
// creating the agent from nothing, which is the "fresh deploy" carve-out
// `Intent` documents.
if intent == Intent::Redeclare
&& let Some(existing) = declaration.agents.get(agent)
&& existing.state == AgentState::Destroyed
&& state != AgentState::Destroyed
{
return Err(TerminalStateError {
agent: agent.to_owned(),
}
.into());
}
// `Create` against a name that already has a *non-terminal* declaration
// is not a creation — it's the create-agent endpoint being asked to
// declare a name that turns out to already exist (a retried request, or
// an operator migrating a pre-existing agent into this bookkeeping for
// the first time). Creation is expected to be idempotent, so this is a
// no-op that leaves the existing declaration exactly as it stood rather
// than stomping it back to whatever state the create path passes —
// otherwise calling create against a name that is already `Up` would
// silently pause a running agent. `Destroyed` is excluded on purpose:
// that's the one case a fresh create *does* need to write over (see
// `Intent::Create`'s own doc comment).
if intent == Intent::Create
&& let Some(existing) = declaration.agents.get(agent)
&& existing.state != AgentState::Destroyed
{
let encoded =
serde_json::to_vec(&declaration).context("encoding the current declaration")?;
return Ok((declaration, encoded));
}
declaration
.agents
.insert(agent.to_owned(), AgentWanted { state });
let encoded = serde_json::to_vec(&declaration).context("encoding the new declaration")?;
Ok((declaration, encoded))
}
/// Writes the wanted-state bucket, and reads it back.
pub struct WantedWriter {
client: async_nats::Client,
}
impl WantedWriter {
#[must_use]
pub fn new(client: async_nats::Client) -> Self {
Self { client }
}
/// One hive's bucket handle, created on first use if nothing has made it
/// yet.
///
/// Creation lives in [`swarm_queue_client::wanted`] because a bucket is
/// described identically by everyone who may create it. Only the
/// controller creates these; a hive opens its own read-only.
///
/// Resolved per call rather than cached: there is one bucket per hive, so
/// a single cached handle cannot serve them, and the declarations this
/// writes change on operator action rather than on a loop — the extra
/// lookup is per *declaration*, not per tick. A cache here would be a map
/// whose invalidation nobody needs yet.
async fn store(
&self,
hive: &str,
) -> std::result::Result<async_nats::jetstream::kv::Store, swarm_queue_client::Error> {
swarm_queue_client::wanted::open_or_create(&self.client, hive).await
}
/// The declaration currently published for `hive`, or `None`.
pub async fn view(&self, hive: &str) -> Result<Option<HiveWanted>> {
// An unconnected client does not fail a JetStream request, it hangs
// on it — see `swarm_queue_client::ensure_connected`.
swarm_queue_client::ensure_connected(&self.client)?;
let store = self.store(hive).await?;
let Some(entry) = store
.entry(hive)
.await
.with_context(|| format!("reading the declaration for {hive}"))?
else {
return Ok(None);
};
serde_json::from_slice(&entry.value)
.map(Some)
.with_context(|| format!("the declaration for {hive} is not decodable"))
}
/// Redeclare `agent` on `hive` to be in `state`, and return the whole
/// declaration as published.
///
/// Refuses to move an agent off `Destroyed` — see
/// [`create`](Self::create) for the one caller that may.
pub async fn set(&self, hive: &str, agent: &str, state: AgentState) -> Result<HiveWanted> {
self.write(hive, agent, state, Intent::Redeclare).await
}
/// Declare a *newly created* `agent` on `hive`, and return the whole
/// declaration as published.
///
/// Same write as [`set`](Self::set) but with [`Intent::Create`], so a
/// `Destroyed` entry left by an earlier agent of this name does not
/// refuse it. Recreating a destroyed name is precisely the "fresh deploy"
/// [`AgentState::Destroyed`]'s doc comment names as the way back, and by
/// the time this runs the swarm has already built that agent's identity,
/// repo and config from nothing.
///
/// Idempotent against any other existing declaration: an agent already
/// declared in a non-terminal state (an operator migrating a pre-existing
/// agent into this bookkeeping, or a retried request) is left exactly as
/// it stood — creation must not silently pause an agent already running
/// under some other wanted state.
///
/// A separate method rather than an `Intent` parameter on `set`: the
/// caller that may overwrite a terminal state is the agent-creation job
/// node and nothing else, and an argument every other caller has to pass
/// correctly is an argument one of them eventually passes wrong.
pub async fn create(&self, hive: &str, agent: &str, state: AgentState) -> Result<HiveWanted> {
self.write(hive, agent, state, Intent::Create).await
}
/// The write both of the above are.
///
/// Read-modify-write against the entry's revision rather than a plain
/// `put`: the value is the hive's whole agent map, so a blind write
/// would drop a concurrent change to a different agent. The bucket keeps
/// one revision of history, so a lost write is not recoverable after the
/// fact — the conflict has to be caught here.
async fn write(
&self,
hive: &str,
agent: &str,
state: AgentState,
intent: Intent,
) -> Result<HiveWanted> {
swarm_queue_client::ensure_connected(&self.client)?;
let store = self.store(hive).await?;
for _ in 0..MAX_ATTEMPTS {
let entry = store
.entry(hive)
.await
.with_context(|| format!("reading the declaration for {hive}"))?;
let revision = entry.as_ref().map(|e| e.revision);
let (declaration, encoded) = apply(
entry.as_ref().map(|e| e.value.as_ref()),
agent,
state,
intent,
)?;
// `apply`'s `Intent::Create` no-op (an already-declared,
// non-terminal agent) re-encodes the declaration unchanged —
// skip the network round-trip for it entirely rather than
// writing identical bytes back.
if entry.as_ref().map(|e| e.value.as_ref()) == Some(encoded.as_slice()) {
tracing::debug!(hive, agent, "declaration already current, no write needed");
return Ok(declaration);
}
// `update` and `create` have separate error types, and only one
// variant of each means "someone else got there first". Every
// other failure returns immediately: retrying a disconnect or a
// permission error would spin the loop and then report a
// conflict, blaming a concurrent writer that never existed.
let written = match revision {
Some(revision) => match store.update(hive, encoded.into(), revision).await {
Ok(_) => Wrote::Ok,
Err(e) if matches!(e.kind(), UpdateErrorKind::WrongLastRevision) => {
Wrote::LostRace
}
Err(e) => Wrote::Failed(anyhow::Error::new(e)),
},
None => match store.create(hive, encoded.into()).await {
Ok(_) => Wrote::Ok,
Err(e) if matches!(e.kind(), CreateErrorKind::AlreadyExists) => Wrote::LostRace,
Err(e) => Wrote::Failed(anyhow::Error::new(e)),
},
};
match written {
Wrote::Ok => {
tracing::info!(hive, agent, ?state, "declared agent state");
return Ok(declaration);
}
// Re-read and re-apply onto the winner's value, not over it.
Wrote::LostRace => {
tracing::debug!(hive, agent, "declaration write lost a race, retrying");
}
Wrote::Failed(e) => {
return Err(e).with_context(|| format!("declaring {agent} on {hive}"));
}
}
}
anyhow::bail!("gave up declaring {agent} on {hive} after {MAX_ATTEMPTS} conflicting writes")
}
}
#[cfg(test)]
mod tests {
use super::{Intent, apply};
use swarm_queue_client::wanted::{AgentState, HiveWanted};
/// Every case below that predates the create carve-out is a redeclare;
/// spelling the intent out in each of them would say nothing the test
/// name doesn't. The carve-out's own tests call `apply` directly.
fn redeclare(
current: Option<&[u8]>,
agent: &str,
state: AgentState,
) -> anyhow::Result<(HiveWanted, Vec<u8>)> {
apply(current, agent, state, Intent::Redeclare)
}
#[test]
fn declaring_one_agent_preserves_every_other() {
let current = br#"{"agents":{"iris":{"state":"up"},"argus":{"state":"offline"}}}"#;
let (declaration, _) = redeclare(Some(current), "atlas", AgentState::Up).unwrap();
assert_eq!(declaration.agents.len(), 3);
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
assert_eq!(declaration.agents["argus"].state, AgentState::Offline);
assert_eq!(declaration.agents["atlas"].state, AgentState::Up);
}
#[test]
fn redeclaring_an_agent_replaces_only_its_own_state() {
let current = br#"{"agents":{"iris":{"state":"up"},"atlas":{"state":"up"}}}"#;
let (declaration, _) = redeclare(Some(current), "atlas", AgentState::Offline).unwrap();
assert_eq!(declaration.agents.len(), 2);
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
assert_eq!(declaration.agents["atlas"].state, AgentState::Offline);
}
#[test]
fn a_hive_with_no_declaration_yet_gets_a_one_agent_one() {
let (declaration, _) = redeclare(None, "atlas", AgentState::Up).unwrap();
assert_eq!(declaration.agents.len(), 1);
assert_eq!(declaration.agents["atlas"].state, AgentState::Up);
}
// The failure this function exists to prevent: a fresh-start fallback
// here would publish a one-agent document over a hive's whole set.
#[test]
fn an_undecodable_declaration_is_an_error_not_a_fresh_start() {
let err = redeclare(Some(b"{not json"), "atlas", AgentState::Up).unwrap_err();
assert!(
err.to_string().contains("not decodable"),
"unexpected error: {err}"
);
}
#[test]
fn an_unknown_state_in_the_current_value_is_also_an_error() {
let current = br#"{"agents":{"iris":{"state":"sideways"}}}"#;
assert!(redeclare(Some(current), "atlas", AgentState::Up).is_err());
}
#[test]
fn the_encoded_form_round_trips() {
let (_, encoded) = redeclare(None, "atlas", AgentState::Offline).unwrap();
let (again, _) = redeclare(Some(&encoded), "iris", AgentState::Up).unwrap();
assert_eq!(again.agents["atlas"].state, AgentState::Offline);
assert_eq!(again.agents["iris"].state, AgentState::Up);
}
/// The bug this module exists to fix: nothing used to stop a caller from
/// declaring a destroyed agent back to `Up`/`Offline`, silently
/// un-terminaling a state whose own doc comment says that is impossible.
#[test]
fn moving_a_destroyed_agent_to_any_other_state_is_refused() {
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
for state in [AgentState::Up, AgentState::Offline, AgentState::Paused] {
let err = redeclare(Some(current), "atlas", state).unwrap_err();
assert!(
err.downcast_ref::<super::TerminalStateError>().is_some(),
"expected a TerminalStateError for {state:?}, got: {err}"
);
}
}
/// Redeclaring `Destroyed` while already `Destroyed` is a no-op, not a
/// transition — a client that already sees the terminal state and asks
/// again should get success back, not an error for stating a fact.
#[test]
fn redeclaring_destroyed_as_destroyed_is_not_a_transition() {
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
let (declaration, _) = redeclare(Some(current), "atlas", AgentState::Destroyed).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Destroyed);
}
/// The refusal is per-agent — a destroyed agent on the same hive as a
/// live one must not block declaring the live one.
#[test]
fn a_destroyed_agent_does_not_block_declaring_a_different_one() {
let current = br#"{"agents":{"atlas":{"state":"destroyed"},"iris":{"state":"up"}}}"#;
let (declaration, _) = redeclare(Some(current), "iris", AgentState::Offline).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Destroyed);
assert_eq!(declaration.agents["iris"].state, AgentState::Offline);
}
/// The carve-out `AgentState::Destroyed`'s doc comment already promised
/// and nothing implemented: a *fresh deploy* does bring a destroyed name
/// back. Without this, creating an agent under a name that was destroyed
/// earlier builds its identity, repo and config and then fails the
/// declaration — cancelling the deploy while the operator sees a 200.
#[test]
fn creating_an_agent_may_reuse_a_destroyed_name() {
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
let (declaration, _) =
apply(Some(current), "atlas", AgentState::Paused, Intent::Create).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Paused);
}
/// The carve-out is scoped to the agent being created — recreating one
/// name must not disturb another destroyed agent on the same hive.
#[test]
fn creating_an_agent_leaves_other_destroyed_agents_destroyed() {
let current = br#"{"agents":{"atlas":{"state":"destroyed"},"iris":{"state":"destroyed"}}}"#;
let (declaration, _) =
apply(Some(current), "atlas", AgentState::Paused, Intent::Create).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Paused);
assert_eq!(declaration.agents["iris"].state, AgentState::Destroyed);
}
/// mara: "this does not match the expectation that agent creation is
/// idempotent so pre existing agents can be migrated" — an agent already
/// declared in some other, non-terminal state (say an operator migrating
/// a pre-existing agent into this bookkeeping) must not be stomped back
/// to whatever state the create path passes.
#[test]
fn creating_an_agent_does_not_disturb_an_existing_non_terminal_declaration() {
let current = br#"{"agents":{"atlas":{"state":"up"}}}"#;
let (declaration, _) =
apply(Some(current), "atlas", AgentState::Paused, Intent::Create).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Up);
}
/// The idempotent no-op only inspects the agent being created — it must
/// not touch any other agent's declaration on the same hive.
#[test]
fn creating_an_agent_that_already_exists_leaves_every_other_agent_alone() {
let current = br#"{"agents":{"atlas":{"state":"up"},"iris":{"state":"offline"}}}"#;
let (declaration, _) =
apply(Some(current), "atlas", AgentState::Paused, Intent::Create).unwrap();
assert_eq!(declaration.agents["atlas"].state, AgentState::Up);
assert_eq!(declaration.agents["iris"].state, AgentState::Offline);
}
/// `Intent::Create` relaxes exactly one rule. An undecodable current
/// value is still an error, for the same reason it is on a redeclare:
/// writing over it discards every other agent on the hive.
#[test]
fn creating_an_agent_still_refuses_an_undecodable_declaration() {
let err = apply(
Some(b"{not json"),
"atlas",
AgentState::Paused,
Intent::Create,
)
.unwrap_err();
assert!(
err.to_string().contains("not decodable"),
"unexpected error: {err}"
);
}
}