feat(#3124): publish the agent set the swarm declares for each hive
The hive-side loop landed without anything to converge to: nothing wrote `$KV.hive-wanted.<hive>`, so in production only the "no key" branch ran. This is the writer. `WantedWriter` mirrors `StatusReader` — that module reads what hives report, this one writes what they are told, so it holds a client rather than a bucket handle and resolves the store on first use. It shares the status reader's connection: the controller has exactly one by design, and a second connect would double the auth-callout traffic and give the two paths independent reconnect state. The value under a hive's key is the map of every agent on that hive, so a plain `put` of a single-agent change would drop a concurrent change to a different agent, with only one revision of history to not recover from. Writes are read-modify-write against the entry revision, and only `WrongLastRevision` / `AlreadyExists` count as a lost race — every other error returns immediately rather than spinning the retry loop and then blaming a concurrent writer that never existed. `apply` is split out and tested because it holds the invariant: declaring one agent preserves the rest, and a current value that will not decode is an error rather than a fresh start. Overwriting a document nobody can read discards every other agent's declaration. Two routes, no swarmctl verb and no jobq node: `create_agent` needs a graph because it is multi-step, and one CAS'd write is not. `build_app` is extracted from `main` in the same change because `main` sat at exactly the `too_many_lines` limit, so adding an endpoint tripped a lint about the startup sequence. The route list is the part that grows.
This commit is contained in:
parent
066f31b58b
commit
76d5871d20
3 changed files with 424 additions and 9 deletions
216
swarm-controller/src/wanted.rs
Normal file
216
swarm-controller/src/wanted.rs
Normal file
|
|
@ -0,0 +1,216 @@
|
|||
//! 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.
|
||||
//! The lifecycle is deliberately identical — a NATS client rather than a
|
||||
//! bucket handle, resolved on first use and cached, so a controller that
|
||||
//! starts before the bucket exists picks it up without a restart.
|
||||
//!
|
||||
//! 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 anyhow::{Context, Result};
|
||||
use async_nats::jetstream::kv::{CreateErrorKind, UpdateErrorKind};
|
||||
use swarm_queue_client::wanted::{AgentState, AgentWanted, HiveWanted};
|
||||
|
||||
/// 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),
|
||||
}
|
||||
|
||||
/// 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.
|
||||
fn apply(current: Option<&[u8]>, agent: &str, state: AgentState) -> 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(),
|
||||
};
|
||||
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,
|
||||
store: tokio::sync::OnceCell<async_nats::jetstream::kv::Store>,
|
||||
}
|
||||
|
||||
impl WantedWriter {
|
||||
#[must_use]
|
||||
pub fn new(client: async_nats::Client) -> Self {
|
||||
Self {
|
||||
client,
|
||||
store: tokio::sync::OnceCell::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// The 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 this one; a hive opens it read-only.
|
||||
async fn store(
|
||||
&self,
|
||||
) -> std::result::Result<&async_nats::jetstream::kv::Store, swarm_queue_client::Error> {
|
||||
self.store
|
||||
.get_or_try_init(|| swarm_queue_client::wanted::open_or_create(&self.client))
|
||||
.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().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"))
|
||||
}
|
||||
|
||||
/// Declare `agent` on `hive` to be in `state`, and return the whole
|
||||
/// declaration as published.
|
||||
///
|
||||
/// 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.
|
||||
pub async fn set(&self, hive: &str, agent: &str, state: AgentState) -> Result<HiveWanted> {
|
||||
swarm_queue_client::ensure_connected(&self.client)?;
|
||||
let store = self.store().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)?;
|
||||
|
||||
// `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::apply;
|
||||
use swarm_queue_client::wanted::AgentState;
|
||||
|
||||
#[test]
|
||||
fn declaring_one_agent_preserves_every_other() {
|
||||
let current = br#"{"agents":{"iris":{"state":"up"},"argus":{"state":"offline"}}}"#;
|
||||
let (declaration, _) = apply(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, _) = apply(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, _) = apply(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 = apply(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!(apply(Some(current), "atlas", AgentState::Up).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn the_encoded_form_round_trips() {
|
||||
let (_, encoded) = apply(None, "atlas", AgentState::Offline).unwrap();
|
||||
let (again, _) = apply(Some(&encoded), "iris", AgentState::Up).unwrap();
|
||||
assert_eq!(again.agents["atlas"].state, AgentState::Offline);
|
||||
assert_eq!(again.agents["iris"].state, AgentState::Up);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue