swarm-controller: read queued placements before published ones
create_agent read the other hives' wanted state first and snapshotted the queued SetAgentWanted nodes second. A node finishing between the two was in neither — not yet declared at the first read, already terminal at the second — so a second hive could get the same name. The queue is now snapshotted first: a SetAgentWanted node only turns terminal after its declaration is written, so anything terminal by then is visible to the read that follows. The order lives in `placements_elsewhere`, and a test that finishes a node between the two reads fails with them swapped. Refs #4396
This commit is contained in:
parent
5785c0024c
commit
efed43b672
1 changed files with 75 additions and 7 deletions
|
|
@ -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<E, Fut>(
|
||||
agent: &str,
|
||||
hive: &str,
|
||||
snapshot_queued: impl FnOnce() -> Vec<(String, String)>,
|
||||
read_declared: impl FnOnce() -> Fut,
|
||||
) -> Result<Vec<String>, E>
|
||||
where
|
||||
Fut: std::future::Future<
|
||||
Output = Result<Vec<(String, swarm_queue_client::wanted::HiveWanted)>, 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;
|
||||
|
|
|
|||
Loading…
Reference in a new issue