diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 053ca140..0efb130d 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -1434,13 +1434,21 @@ async fn create_agent( // path and matrix id all carry the bare name. The same name on the same // hive is that agent being re-created, and goes through. let _gate = state.create_gate.lock().await; - let declared = declarations_elsewhere(&state, &agent, &hive).await?; - let mut sched = state - .jobq - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - let queued = queued_placements(sched.graph()); - let elsewhere = placed_elsewhere(&agent, &hive, &declared, &queued); + let elsewhere = placements_elsewhere( + &agent, + &hive, + || { + queued_placements( + state + .jobq + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .graph(), + ) + }, + || declarations_elsewhere(&state, &agent, &hive), + ) + .await?; if !elsewhere.is_empty() { let detail = format!( "agent name {agent:?} is already placed on hive {} — agent names are unique across \ @@ -1449,6 +1457,10 @@ async fn create_agent( ); return Err(error_problem(axum::http::StatusCode::CONFLICT, &detail)); } + let mut sched = state + .jobq + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); let ids = sched .insert_job(None, |b| declare_agent_job(b, &agent, &hive)) .map_err(|e| { @@ -1466,6 +1478,30 @@ async fn create_agent( })) } +/// [`placed_elsewhere`] over a snapshot of the queued placements and a read +/// of the published ones, **in that order**. +/// +/// A `SetAgentWanted` node writes its declaration before it turns terminal. +/// Snapshotting the queue first means a node already terminal there has +/// written, so the read that follows sees its declaration; a node still +/// queued is in the snapshot. Reading first leaves a gap: a node finishing +/// between the two is in neither. +async fn placements_elsewhere( + agent: &str, + hive: &str, + snapshot_queued: impl FnOnce() -> Vec<(String, String)>, + read_declared: impl FnOnce() -> Fut, +) -> Result, E> +where + Fut: std::future::Future< + Output = Result, E>, + >, +{ + let queued = snapshot_queued(); + let declared = read_declared().await?; + Ok(placed_elsewhere(agent, hive, &declared, &queued)) +} + /// Every hive but `hive`'s published wanted state, for [`placed_elsewhere`]. /// /// No queue wired up means no hive has been sent a declaration, so there is @@ -3070,6 +3106,38 @@ mod tests { (hive.to_owned(), declaration) } + /// A `SetAgentWanted` node that finishes between the two reads: whichever + /// read runs first sees it queued and undeclared, the second sees it + /// terminal and declared. Only the queue-first order catches it. + #[tokio::test] + async fn a_placement_finishing_between_the_two_reads_is_still_caught() { + use swarm_queue_client::wanted::AgentState; + let reads = std::cell::Cell::new(0); + let snapshot_queued = || { + reads.set(reads.get() + 1); + if reads.get() == 1 { + vec![("sec0nd".to_owned(), "atlas".to_owned())] + } else { + Vec::new() + } + }; + let read_declared = || { + reads.set(reads.get() + 1); + let declared = if reads.get() == 1 { + Vec::new() + } else { + vec![declaring("sec0nd", "atlas", AgentState::Paused)] + }; + std::future::ready(Ok::<_, ()>(declared)) + }; + + let elsewhere = + super::placements_elsewhere("atlas", "pr1ma", snapshot_queued, read_declared) + .await + .expect("both reads succeed"); + assert_eq!(elsewhere, ["sec0nd"]); + } + #[test] fn a_live_declaration_on_another_hive_is_a_placement() { use swarm_queue_client::wanted::AgentState;