swarm-controller: let creating an agent reuse a destroyed name
Review caught a path the new declaration node breaks: an operator may
already declare an agent `Destroyed` over the per-agent state endpoint,
and `apply` then refuses any transition off that state. Before the
wanted-state node existed the refusal was inert at creation time, but
now creating an agent under a previously-destroyed name builds its
identity, repo and config, fails the declaration, and silently cancels
the deploy — while the caller sees a 200 and a job id.
`AgentState::Destroyed`'s own doc comment already sanctions this case
("no state that brings a destroyed agent back short of a fresh deploy");
nothing implemented it. Give `apply` an `Intent`, keep the refusal for
redeclares, and add `WantedWriter::create` for the one caller that is a
fresh deploy. A separate method rather than a parameter on `set`, so no
other caller can reach the override by passing an argument wrong.
This commit is contained in:
parent
bfd019fe8d
commit
48aa7a1e79
2 changed files with 152 additions and 25 deletions
|
|
@ -340,7 +340,12 @@ async fn declare_new_agent(
|
||||||
) -> hive_jobq::scheduler::Outcome {
|
) -> hive_jobq::scheduler::Outcome {
|
||||||
use 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,
|
Ok(_) => Outcome::Done,
|
||||||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||||
}
|
}
|
||||||
|
|
@ -1387,11 +1392,11 @@ fn declare_agent_job(
|
||||||
// unconfigured option into an agent nobody runs.
|
// unconfigured option into an agent nobody runs.
|
||||||
//
|
//
|
||||||
// `after_ok` on the declaration, though: an agent deployed without its
|
// `after_ok` on the declaration, though: an agent deployed without its
|
||||||
// pause landing first is the race this node exists to close, and the
|
// pause landing first is the race this node exists to close, so a
|
||||||
// only way the declaration fails without a live queue is a host that has
|
// declaration that did not land must not be deployed past. A host with no
|
||||||
// no queue at all — on which `TriggerDeploy` has nothing to publish to
|
// queue at all fails this node — but `TriggerDeploy` has nothing to
|
||||||
// either, so the strict edge cancels a node that could not have
|
// publish to on such a host either, so the strict edge cancels a node
|
||||||
// succeeded.
|
// that could not have succeeded anyway.
|
||||||
let _trigger_deploy = b
|
let _trigger_deploy = b
|
||||||
.node(SwarmNodeKind::TriggerDeploy {
|
.node(SwarmNodeKind::TriggerDeploy {
|
||||||
hive: hive.to_owned(),
|
hive: hive.to_owned(),
|
||||||
|
|
|
||||||
|
|
@ -57,6 +57,28 @@ enum Wrote {
|
||||||
Failed(anyhow::Error),
|
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.
|
/// 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
|
/// 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
|
/// value is an **error**, never treated as absent: overwriting a document
|
||||||
/// nobody can read discards the declarations of every other agent on that
|
/// nobody can read discards the declarations of every other agent on that
|
||||||
/// hive, which is exactly what a fresh-start fallback would do quietly.
|
/// 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>)> {
|
///
|
||||||
|
/// `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 {
|
let mut declaration = match current {
|
||||||
Some(raw) => serde_json::from_slice::<HiveWanted>(raw)
|
Some(raw) => serde_json::from_slice::<HiveWanted>(raw)
|
||||||
.context("the hive's current declaration is not decodable")?,
|
.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
|
// 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
|
// client that already sees the agent as destroyed and asks again should
|
||||||
// not be punished for it. Anything else moving off `Destroyed` is the
|
// not be punished for it. Anything else moving off `Destroyed` is the
|
||||||
// one transition this module exists to refuse.
|
// one transition this module exists to refuse — unless the caller is
|
||||||
if let Some(existing) = declaration.agents.get(agent)
|
// 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
|
&& existing.state == AgentState::Destroyed
|
||||||
&& state != AgentState::Destroyed
|
&& state != AgentState::Destroyed
|
||||||
{
|
{
|
||||||
|
|
@ -149,15 +182,47 @@ impl WantedWriter {
|
||||||
.with_context(|| format!("the declaration for {hive} is not decodable"))
|
.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.
|
/// 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.
|
||||||
|
///
|
||||||
|
/// 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
|
/// 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
|
/// `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
|
/// 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
|
/// one revision of history, so a lost write is not recoverable after the
|
||||||
/// fact — the conflict has to be caught here.
|
/// fact — the conflict has to be caught here.
|
||||||
pub async fn set(&self, hive: &str, agent: &str, state: AgentState) -> Result<HiveWanted> {
|
async fn write(
|
||||||
|
&self,
|
||||||
|
hive: &str,
|
||||||
|
agent: &str,
|
||||||
|
state: AgentState,
|
||||||
|
intent: Intent,
|
||||||
|
) -> Result<HiveWanted> {
|
||||||
swarm_queue_client::ensure_connected(&self.client)?;
|
swarm_queue_client::ensure_connected(&self.client)?;
|
||||||
let store = self.store(hive).await?;
|
let store = self.store(hive).await?;
|
||||||
|
|
||||||
|
|
@ -167,8 +232,12 @@ impl WantedWriter {
|
||||||
.await
|
.await
|
||||||
.with_context(|| format!("reading the declaration for {hive}"))?;
|
.with_context(|| format!("reading the declaration for {hive}"))?;
|
||||||
let revision = entry.as_ref().map(|e| e.revision);
|
let revision = entry.as_ref().map(|e| e.revision);
|
||||||
let (declaration, encoded) =
|
let (declaration, encoded) = apply(
|
||||||
apply(entry.as_ref().map(|e| e.value.as_ref()), agent, state)?;
|
entry.as_ref().map(|e| e.value.as_ref()),
|
||||||
|
agent,
|
||||||
|
state,
|
||||||
|
intent,
|
||||||
|
)?;
|
||||||
|
|
||||||
// `update` and `create` have separate error types, and only one
|
// `update` and `create` have separate error types, and only one
|
||||||
// variant of each means "someone else got there first". Every
|
// variant of each means "someone else got there first". Every
|
||||||
|
|
@ -209,13 +278,24 @@ impl WantedWriter {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::apply;
|
use super::{Intent, apply};
|
||||||
use swarm_queue_client::wanted::AgentState;
|
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]
|
#[test]
|
||||||
fn declaring_one_agent_preserves_every_other() {
|
fn declaring_one_agent_preserves_every_other() {
|
||||||
let current = br#"{"agents":{"iris":{"state":"up"},"argus":{"state":"offline"}}}"#;
|
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.len(), 3);
|
||||||
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
|
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
|
||||||
assert_eq!(declaration.agents["argus"].state, AgentState::Offline);
|
assert_eq!(declaration.agents["argus"].state, AgentState::Offline);
|
||||||
|
|
@ -225,7 +305,7 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn redeclaring_an_agent_replaces_only_its_own_state() {
|
fn redeclaring_an_agent_replaces_only_its_own_state() {
|
||||||
let current = br#"{"agents":{"iris":{"state":"up"},"atlas":{"state":"up"}}}"#;
|
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.len(), 2);
|
||||||
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
|
assert_eq!(declaration.agents["iris"].state, AgentState::Up);
|
||||||
assert_eq!(declaration.agents["atlas"].state, AgentState::Offline);
|
assert_eq!(declaration.agents["atlas"].state, AgentState::Offline);
|
||||||
|
|
@ -233,7 +313,7 @@ mod tests {
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn a_hive_with_no_declaration_yet_gets_a_one_agent_one() {
|
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.len(), 1);
|
||||||
assert_eq!(declaration.agents["atlas"].state, AgentState::Up);
|
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.
|
// here would publish a one-agent document over a hive's whole set.
|
||||||
#[test]
|
#[test]
|
||||||
fn an_undecodable_declaration_is_an_error_not_a_fresh_start() {
|
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!(
|
assert!(
|
||||||
err.to_string().contains("not decodable"),
|
err.to_string().contains("not decodable"),
|
||||||
"unexpected error: {err}"
|
"unexpected error: {err}"
|
||||||
|
|
@ -252,13 +332,13 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn an_unknown_state_in_the_current_value_is_also_an_error() {
|
fn an_unknown_state_in_the_current_value_is_also_an_error() {
|
||||||
let current = br#"{"agents":{"iris":{"state":"sideways"}}}"#;
|
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]
|
#[test]
|
||||||
fn the_encoded_form_round_trips() {
|
fn the_encoded_form_round_trips() {
|
||||||
let (_, encoded) = apply(None, "atlas", AgentState::Offline).unwrap();
|
let (_, encoded) = redeclare(None, "atlas", AgentState::Offline).unwrap();
|
||||||
let (again, _) = apply(Some(&encoded), "iris", AgentState::Up).unwrap();
|
let (again, _) = redeclare(Some(&encoded), "iris", AgentState::Up).unwrap();
|
||||||
assert_eq!(again.agents["atlas"].state, AgentState::Offline);
|
assert_eq!(again.agents["atlas"].state, AgentState::Offline);
|
||||||
assert_eq!(again.agents["iris"].state, AgentState::Up);
|
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() {
|
fn moving_a_destroyed_agent_to_any_other_state_is_refused() {
|
||||||
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
|
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
|
||||||
for state in [AgentState::Up, AgentState::Offline, AgentState::Paused] {
|
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!(
|
assert!(
|
||||||
err.downcast_ref::<super::TerminalStateError>().is_some(),
|
err.downcast_ref::<super::TerminalStateError>().is_some(),
|
||||||
"expected a TerminalStateError for {state:?}, got: {err}"
|
"expected a TerminalStateError for {state:?}, got: {err}"
|
||||||
|
|
@ -284,7 +364,7 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn redeclaring_destroyed_as_destroyed_is_not_a_transition() {
|
fn redeclaring_destroyed_as_destroyed_is_not_a_transition() {
|
||||||
let current = br#"{"agents":{"atlas":{"state":"destroyed"}}}"#;
|
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);
|
assert_eq!(declaration.agents["atlas"].state, AgentState::Destroyed);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -293,8 +373,50 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn a_destroyed_agent_does_not_block_declaring_a_different_one() {
|
fn a_destroyed_agent_does_not_block_declaring_a_different_one() {
|
||||||
let current = br#"{"agents":{"atlas":{"state":"destroyed"},"iris":{"state":"up"}}}"#;
|
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["atlas"].state, AgentState::Destroyed);
|
||||||
assert_eq!(declaration.agents["iris"].state, AgentState::Offline);
|
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}"
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue