diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 609fd7ab..f394f011 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -340,7 +340,12 @@ async fn declare_new_agent( ) -> hive_jobq::scheduler::Outcome { use hive_jobq::scheduler::Outcome; - match writer.set(hive, agent, NEW_AGENT_WANTED_STATE).await { + // `create`, not `set`: a name that was destroyed earlier still carries a + // terminal entry, and recreating the agent is the fresh deploy that + // entry's own doc comment names as the way back. `set` would refuse it, + // cancelling the deploy of an agent whose identity, repo and config the + // swarm has just built. + match writer.create(hive, agent, NEW_AGENT_WANTED_STATE).await { Ok(_) => Outcome::Done, Err(e) => Outcome::Failed(format!("{e:#}")), } @@ -1387,11 +1392,11 @@ fn declare_agent_job( // unconfigured option into an agent nobody runs. // // `after_ok` on the declaration, though: an agent deployed without its - // pause landing first is the race this node exists to close, and the - // only way the declaration fails without a live queue is a host that has - // no queue at all — on which `TriggerDeploy` has nothing to publish to - // either, so the strict edge cancels a node that could not have - // succeeded. + // pause landing first is the race this node exists to close, so a + // declaration that did not land must not be deployed past. A host with no + // queue at all fails this node — but `TriggerDeploy` has nothing to + // publish to on such a host either, so the strict edge cancels a node + // that could not have succeeded anyway. let _trigger_deploy = b .node(SwarmNodeKind::TriggerDeploy { hive: hive.to_owned(), diff --git a/swarm-controller/src/wanted.rs b/swarm-controller/src/wanted.rs index 6a4d0863..4d5def90 100644 --- a/swarm-controller/src/wanted.rs +++ b/swarm-controller/src/wanted.rs @@ -57,6 +57,28 @@ enum Wrote { 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 @@ -75,7 +97,15 @@ const MAX_ATTEMPTS: usize = 5; /// 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)> { +/// +/// `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)> { let mut declaration = match current { Some(raw) => serde_json::from_slice::(raw) .context("the hive's current declaration is not decodable")?, @@ -84,8 +114,11 @@ fn apply(current: Option<&[u8]>, agent: &str, state: AgentState) -> Result<(Hive // 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. - if let Some(existing) = declaration.agents.get(agent) + // 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 { @@ -149,15 +182,47 @@ impl WantedWriter { .with_context(|| format!("the declaration for {hive} is not decodable")) } - /// Declare `agent` on `hive` to be in `state`, and return the whole + /// 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 { + 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. + /// + /// 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 { + 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. - pub async fn set(&self, hive: &str, agent: &str, state: AgentState) -> Result { + async fn write( + &self, + hive: &str, + agent: &str, + state: AgentState, + intent: Intent, + ) -> Result { swarm_queue_client::ensure_connected(&self.client)?; let store = self.store(hive).await?; @@ -167,8 +232,12 @@ impl WantedWriter { .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)?; + let (declaration, encoded) = apply( + entry.as_ref().map(|e| e.value.as_ref()), + agent, + state, + intent, + )?; // `update` and `create` have separate error types, and only one // variant of each means "someone else got there first". Every @@ -209,13 +278,24 @@ impl WantedWriter { #[cfg(test)] mod tests { - use super::apply; - use swarm_queue_client::wanted::AgentState; + 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)> { + 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, _) = apply(Some(current), "atlas", AgentState::Up).unwrap(); + 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); @@ -225,7 +305,7 @@ mod tests { #[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(); + 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); @@ -233,7 +313,7 @@ mod tests { #[test] fn a_hive_with_no_declaration_yet_gets_a_one_agent_one() { - let (declaration, _) = apply(None, "atlas", AgentState::Up).unwrap(); + let (declaration, _) = redeclare(None, "atlas", AgentState::Up).unwrap(); assert_eq!(declaration.agents.len(), 1); assert_eq!(declaration.agents["atlas"].state, AgentState::Up); } @@ -242,7 +322,7 @@ mod tests { // 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(); + let err = redeclare(Some(b"{not json"), "atlas", AgentState::Up).unwrap_err(); assert!( err.to_string().contains("not decodable"), "unexpected error: {err}" @@ -252,13 +332,13 @@ mod tests { #[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()); + assert!(redeclare(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(); + 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); } @@ -270,7 +350,7 @@ mod tests { 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 = apply(Some(current), "atlas", state).unwrap_err(); + let err = redeclare(Some(current), "atlas", state).unwrap_err(); assert!( err.downcast_ref::().is_some(), "expected a TerminalStateError for {state:?}, got: {err}" @@ -284,7 +364,7 @@ mod tests { #[test] fn redeclaring_destroyed_as_destroyed_is_not_a_transition() { let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#; - let (declaration, _) = apply(Some(current), "atlas", AgentState::Destroyed).unwrap(); + let (declaration, _) = redeclare(Some(current), "atlas", AgentState::Destroyed).unwrap(); assert_eq!(declaration.agents["atlas"].state, AgentState::Destroyed); } @@ -293,8 +373,50 @@ mod tests { #[test] fn a_destroyed_agent_does_not_block_declaring_a_different_one() { let current = br#"{"agents":{"atlas":{"state":"destroyed"},"iris":{"state":"up"}}}"#; - let (declaration, _) = apply(Some(current), "iris", AgentState::Offline).unwrap(); + 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); + } + + /// `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}" + ); + } }