use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use std::time::Duration;
use crate::prelude::*;
use serde::{Deserialize, Serialize};
use crate::cluster::framing::SpawnRequest;
use crate::app::election::LeaderElection;
use crate::app::error::OrchestratorError;
use crate::app::node_info::{ClusterView, NodeInfo};
use crate::app::placement::{self, PlacementDecision, PlacementStrategy};
use crate::app::singleton::{
CoordinatorGenerationSource, GenerationSource, SingletonGeneration, SingletonOwnership,
SingletonPhase, SingletonSpec, resolve_owner,
};
use crate::app::spawn_sender::{SingletonStopRequest, SingletonStopSender, SpawnSender};
use crate::app::spec::{ActorSpec, CrashStrategy};
#[derive(Debug)]
pub struct Coordinator;
pub struct CoordinatorState {
pub cluster_view: ClusterView,
pub specs: BTreeMap<String, ActorSpec>,
pub placement_strategy: Box<dyn PlacementStrategy>,
pub election: Box<dyn LeaderElection>,
pub local_node_id: String,
pub pending_spawns: HashMap<u64, PendingSpawn>,
pub waiting_for_return: BTreeMap<String, WaitingSpec>,
pub singletons: BTreeMap<String, SingletonRuntime>,
pub generation_source: Arc<dyn GenerationSource>,
next_request_id: u64,
spawn_sender: Option<SpawnSender>,
singleton_stop_sender: Option<SingletonStopSender>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum OwnerKey {
Singleton(String),
}
#[derive(Debug, Clone)]
pub struct PendingSpawn {
pub spec_label: String,
pub target_node_id: String,
pub owner: Option<OwnerKey>,
}
#[derive(Debug, Clone)]
pub struct SingletonRuntime {
pub spec: SingletonSpec,
pub ownership: Option<SingletonOwnership>,
pub pending_handoff: Option<PendingHandoff>,
}
#[derive(Debug, Clone)]
pub struct PendingHandoff {
pub old_owner: String,
pub old_generation: SingletonGeneration,
pub destination: Option<HandoffDestination>,
}
#[derive(Debug, Clone)]
pub struct HandoffDestination {
pub new_owner: String,
pub new_generation: SingletonGeneration,
}
#[derive(Debug, Clone)]
pub struct WaitingSpec {
pub spec: ActorSpec,
pub failed_node_id: String,
}
impl Actor for Coordinator {
type State = CoordinatorState;
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Result<PlacementDecision, String>, remote = "coordinator::SubmitSpec")]
pub struct SubmitSpec {
pub spec: ActorSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Result<(), String>, remote = "coordinator::RemoveSpec")]
pub struct RemoveSpec {
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = ClusterViewSnapshot, remote = "coordinator::GetClusterView")]
pub struct GetClusterView;
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Vec<SpecStatus>, remote = "coordinator::GetSpecs")]
pub struct GetSpecs;
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Result<SingletonOwnership, String>, remote = "coordinator::StartSingleton")]
pub struct StartSingleton {
pub spec: SingletonSpec,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Option<SingletonOwnership>, remote = "coordinator::GetSingleton")]
pub struct GetSingleton {
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Result<(), String>, remote = "coordinator::MoveSingleton")]
pub struct MoveSingleton {
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = Result<(), String>, remote = "coordinator::StopSingleton")]
pub struct StopSingleton {
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::NotifySingletonStopped")]
pub struct NotifySingletonStopped {
pub label: String,
pub stopped_generation: SingletonGeneration,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::SingletonDrainTimeout")]
pub struct SingletonDrainTimeout {
pub label: String,
pub drained_generation: SingletonGeneration,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::NotifyNodeJoined")]
pub struct NotifyNodeJoined {
pub node_id: String,
pub info: SerializableNodeInfo,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::NotifyNodeFailed")]
pub struct NotifyNodeFailed {
pub node_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::NotifyNodeLeft")]
pub struct NotifyNodeLeft {
pub node_id: String,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::NotifySpawnAck")]
pub struct NotifySpawnAck {
pub request_id: u64,
pub success: bool,
pub error: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize, Message)]
#[message(result = (), remote = "coordinator::WaitForReturnTimeout")]
pub struct WaitForReturnTimeout {
pub label: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ClusterViewSnapshot {
pub nodes: Vec<NodeSnapshot>,
pub alive_count: usize,
pub total_count: usize,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeSnapshot {
pub node_id: String,
pub name: String,
pub class: String,
pub actor_count: usize,
pub is_alive: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SpecStatus {
pub label: String,
pub actor_type: String,
pub placed_on: Option<String>,
pub state: SpecState,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum SpecState {
Running,
Pending,
WaitingForReturn,
Unplaced,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SerializableNodeInfo {
pub name: String,
pub endpoint_id: crate::cluster::net::NodeId,
pub host: String,
pub port: u16,
pub incarnation: u64,
pub class: crate::cluster::config::NodeClass,
pub metadata: HashMap<String, String>,
}
#[handlers]
impl Coordinator {
#[handler]
fn submit_spec(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: SubmitSpec,
) -> Result<PlacementDecision, String> {
if !state.is_leader() {
return Err(OrchestratorError::NotLeader.to_string());
}
let label = msg.spec.label.clone();
if state.specs.contains_key(&label) {
return Err(OrchestratorError::SpecAlreadyExists { label }.to_string());
}
let decision = placement::select_node(
state.placement_strategy.as_ref(),
&msg.spec,
&state.cluster_view,
)
.ok_or_else(|| {
OrchestratorError::NoEligibleNodes {
reason: format!("no node satisfies constraints for {label}"),
}
.to_string()
})?;
tracing::info!(
"Placing {} on {} ({})",
label,
decision.node_id,
decision.reason
);
state.specs.insert(label.clone(), msg.spec.clone());
let request_id = state.next_request_id();
state.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: label,
target_node_id: decision.node_id.clone(),
owner: None,
},
);
state.send_spawn(&decision.node_id, request_id, &msg.spec);
Ok(decision)
}
#[handler]
fn remove_spec(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: RemoveSpec,
) -> Result<(), String> {
if !state.is_leader() {
return Err(OrchestratorError::NotLeader.to_string());
}
if state.specs.remove(&msg.label).is_none() {
return Err(OrchestratorError::SpecNotFound {
label: msg.label.clone(),
}
.to_string());
}
state
.pending_spawns
.retain(|_, ps| ps.spec_label != msg.label);
state.waiting_for_return.remove(&msg.label);
state.cluster_view.remove_actor_anywhere(&msg.label);
tracing::info!("Removed spec: {}", msg.label);
Ok(())
}
#[handler]
fn get_cluster_view(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
_msg: GetClusterView,
) -> ClusterViewSnapshot {
ClusterViewSnapshot {
alive_count: state.cluster_view.alive_count(),
total_count: state.cluster_view.total_count(),
nodes: state
.cluster_view
.nodes
.values()
.map(|n| NodeSnapshot {
node_id: n.node_id(),
name: n.identity.name.clone(),
class: n.class.to_string(),
actor_count: n.actor_count(),
is_alive: n.is_alive,
})
.collect(),
}
}
#[handler]
fn get_specs(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
_msg: GetSpecs,
) -> Vec<SpecStatus> {
state
.specs
.iter()
.map(|(label, spec)| {
let placed_on = state.cluster_view.find_actor(label).map(|s| s.to_string());
let spec_state = if state.waiting_for_return.contains_key(label) {
SpecState::WaitingForReturn
} else if state
.pending_spawns
.values()
.any(|ps| ps.spec_label == *label)
{
SpecState::Pending
} else if placed_on.is_some() {
SpecState::Running
} else {
SpecState::Unplaced
};
SpecStatus {
label: label.clone(),
actor_type: spec.actor_type_name.clone(),
placed_on,
state: spec_state,
}
})
.collect()
}
#[handler]
fn notify_node_joined(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: NotifyNodeJoined,
) {
let info = NodeInfo::new(
crate::cluster::config::NodeIdentity {
name: msg.info.name,
endpoint_id: msg.info.endpoint_id,
host: msg.info.host,
port: msg.info.port,
incarnation: msg.info.incarnation,
},
msg.info.class,
msg.info.metadata,
);
tracing::info!("Node joined: {}", msg.node_id);
state.cluster_view.upsert_node(info);
let rejoined: Vec<String> = state
.waiting_for_return
.iter()
.filter(|(_, ws)| ws.failed_node_id == msg.node_id)
.map(|(label, _)| label.clone())
.collect();
for label in rejoined {
if let Some(ws) = state.waiting_for_return.remove(&label) {
tracing::info!(
"Node {} returned — re-spawning spec {} on it",
msg.node_id,
ws.spec.label
);
let request_id = state.next_request_id();
state.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: ws.spec.label.clone(),
target_node_id: msg.node_id.clone(),
owner: None,
},
);
state.send_spawn(&msg.node_id, request_id, &ws.spec);
}
}
}
#[handler]
async fn notify_node_failed(
&mut self,
ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: NotifyNodeFailed,
) {
tracing::warn!("Node failed: {}", msg.node_id);
state.cluster_view.mark_failed(&msg.node_id);
let timers = state.handle_node_departure(&msg.node_id, false);
for (label, duration) in timers {
let endpoint = ctx.endpoint();
let runtime = ctx.receptionist().runtime().clone();
ctx.spawn(async move {
runtime.sleep(duration).await;
let _ = endpoint.send(WaitForReturnTimeout { label }).await;
});
}
state.load_singletons_from_backend().await;
state.redrive_singletons_after_loss(&msg.node_id).await;
}
#[handler]
async fn notify_node_left(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: NotifyNodeLeft,
) {
tracing::info!("Node left gracefully: {}", msg.node_id);
let _timers = state.handle_node_departure(&msg.node_id, true);
state.cluster_view.remove_node(&msg.node_id);
state.load_singletons_from_backend().await;
state.redrive_singletons_after_loss(&msg.node_id).await;
}
#[handler]
fn notify_spawn_ack(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: NotifySpawnAck,
) {
if let Some(pending) = state.pending_spawns.remove(&msg.request_id) {
if msg.success {
tracing::info!(
"Spawn confirmed: {} on {}",
pending.spec_label,
pending.target_node_id
);
state
.cluster_view
.add_actor(&pending.target_node_id, &pending.spec_label);
if let Some(OwnerKey::Singleton(label)) = &pending.owner
&& let Some(rt) = state.singletons.get_mut(label)
&& let Some(ownership) = &mut rt.ownership
{
ownership.phase = SingletonPhase::Active;
}
} else {
tracing::warn!(
"Spawn failed for {}: {}",
pending.spec_label,
msg.error.as_deref().unwrap_or("unknown error")
);
}
}
}
#[handler]
fn wait_for_return_timeout(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: WaitForReturnTimeout,
) {
if let Some(ws) = state.waiting_for_return.remove(&msg.label) {
tracing::warn!(
"WaitForReturn timeout for {} (was on {}) — falling back to Redistribute",
msg.label,
ws.failed_node_id,
);
if let Some(decision) = placement::select_node(
state.placement_strategy.as_ref(),
&ws.spec,
&state.cluster_view,
) {
let request_id = state.next_request_id();
state.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: msg.label.clone(),
target_node_id: decision.node_id.clone(),
owner: None,
},
);
state.send_spawn(&decision.node_id, request_id, &ws.spec);
tracing::info!(
"Re-placed {} on {} after timeout ({})",
msg.label,
decision.node_id,
decision.reason
);
} else {
tracing::warn!(
"No eligible node for {} after WaitForReturn timeout — spec remains unplaced",
msg.label
);
}
}
}
#[handler]
async fn start_singleton(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: StartSingleton,
) -> Result<SingletonOwnership, String> {
if !state.is_leader() {
return Err(OrchestratorError::NotLeader.to_string());
}
let label = msg.spec.label.clone();
let owner = resolve_owner(
&msg.spec.anchor,
&state.cluster_view,
state.election.as_ref(),
)
.ok_or_else(|| {
OrchestratorError::NoEligibleNodes {
reason: format!("no node satisfies anchor for singleton {label}"),
}
.to_string()
})?;
let reassert = state
.singletons
.get(&label)
.and_then(|rt| rt.ownership.as_ref())
.is_some_and(|own| {
own.owner_node_id.as_deref() == Some(owner.as_str())
&& own.phase == SingletonPhase::Active
});
if reassert {
let gen_source = state.generation_source.clone();
let generation = gen_source.claim_seq(&label).await?;
let ownership = SingletonOwnership {
label: label.clone(),
owner_node_id: Some(owner),
generation,
phase: SingletonPhase::Active,
};
if let Some(rt) = state.singletons.get_mut(&label) {
rt.ownership = Some(ownership.clone());
}
tracing::info!("Re-asserted singleton {label} (gen={generation:?})");
return Ok(ownership);
}
let gen_source = state.generation_source.clone();
let generation = gen_source.claim_term(&label, &owner).await?;
gen_source
.put_spec(&label, &msg.spec)
.await
.map_err(|e| format!("singleton {label}: put_spec failed: {e}"))?;
let ownership = SingletonOwnership {
label: label.clone(),
owner_node_id: Some(owner.clone()),
generation,
phase: SingletonPhase::Starting,
};
state.singletons.insert(
label.clone(),
SingletonRuntime {
spec: msg.spec.clone(),
ownership: Some(ownership.clone()),
pending_handoff: None,
},
);
let request_id = state.next_request_id();
state.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: label.clone(),
target_node_id: owner.clone(),
owner: Some(OwnerKey::Singleton(label.clone())),
},
);
state.send_singleton_spawn(&owner, request_id, &msg.spec, generation.packed());
tracing::info!("Starting singleton {label} on {owner} (gen={generation:?})");
Ok(ownership)
}
#[handler]
fn get_singleton(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: GetSingleton,
) -> Option<SingletonOwnership> {
state
.singletons
.get(&msg.label)
.and_then(|rt| rt.ownership.clone())
}
#[handler]
async fn move_singleton(
&mut self,
ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: MoveSingleton,
) -> Result<(), String> {
if !state.is_leader() {
return Err(OrchestratorError::NotLeader.to_string());
}
let label = msg.label;
let Some(rt) = state.singletons.get(&label) else {
return Err(format!("singleton {label} not found"));
};
if rt.pending_handoff.is_some() {
return Err(format!("singleton {label} handoff already in progress"));
}
let Some(ownership) = rt.ownership.clone() else {
return Err(format!("singleton {label} has no current owner to move"));
};
let Some(current) = ownership.owner_node_id.clone() else {
return Err(format!("singleton {label} has no current owner to move"));
};
let current_gen = ownership.generation;
let anchor = rt.spec.anchor.clone();
let drain_timeout = rt.spec.drain_timeout;
let Some(new_owner) = resolve_owner(&anchor, &state.cluster_view, state.election.as_ref())
else {
return Err(OrchestratorError::NoEligibleNodes {
reason: format!("no node satisfies anchor for singleton {label}"),
}
.to_string());
};
if new_owner == current {
return Ok(()); }
let old_alive = state
.cluster_view
.nodes
.get(¤t)
.map(|n| n.is_alive)
.unwrap_or(false);
if !old_alive {
return Err(format!(
"singleton {label} owner {current} is not alive — handled by failover, not graceful move"
));
}
let gen_source = state.generation_source.clone();
let new_gen = gen_source.claim_term(&label, &new_owner).await?;
if let Some(rt) = state.singletons.get_mut(&label) {
if let Some(own) = &mut rt.ownership {
own.phase = SingletonPhase::Draining;
}
rt.pending_handoff = Some(PendingHandoff {
old_owner: current.clone(),
old_generation: current_gen,
destination: Some(HandoffDestination {
new_owner: new_owner.clone(),
new_generation: new_gen,
}),
});
}
state.send_stop_singleton(¤t, &label, current_gen.packed());
state.schedule_drain_timeout(ctx, &label, current_gen, drain_timeout);
tracing::info!(
"Graceful move of singleton {label}: draining {current} (gen={current_gen:?}) -> {new_owner} (gen={new_gen:?})"
);
Ok(())
}
#[handler]
fn stop_singleton(
&mut self,
ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: StopSingleton,
) -> Result<(), String> {
if !state.is_leader() {
return Err(OrchestratorError::NotLeader.to_string());
}
let label = msg.label;
let Some(rt) = state.singletons.get(&label) else {
return Err(format!("singleton {label} not found"));
};
if rt.pending_handoff.is_some() {
return Err(format!("singleton {label} handoff already in progress"));
}
let current = rt.ownership.as_ref().and_then(|o| o.owner_node_id.clone());
let current_gen = rt.ownership.as_ref().map(|o| o.generation);
let drain_timeout = rt.spec.drain_timeout;
let old_alive = current
.as_ref()
.and_then(|c| state.cluster_view.nodes.get(c))
.map(|n| n.is_alive)
.unwrap_or(false);
if !old_alive {
state.singletons.remove(&label);
tracing::info!("Singleton {label} torn down (no live owner to drain)");
return Ok(());
}
let (current, current_gen) = (current.unwrap(), current_gen.unwrap());
if let Some(rt) = state.singletons.get_mut(&label) {
if let Some(own) = &mut rt.ownership {
own.phase = SingletonPhase::Draining;
}
rt.pending_handoff = Some(PendingHandoff {
old_owner: current.clone(),
old_generation: current_gen,
destination: None, });
}
state.send_stop_singleton(¤t, &label, current_gen.packed());
state.schedule_drain_timeout(ctx, &label, current_gen, drain_timeout);
tracing::info!("Tearing down singleton {label}: draining {current} (gen={current_gen:?})");
Ok(())
}
#[handler]
fn notify_singleton_stopped(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: NotifySingletonStopped,
) {
state.complete_handoff(&msg.label, msg.stopped_generation);
}
#[handler]
fn singleton_drain_timeout(
&mut self,
_ctx: &ActorContext<Self>,
state: &mut CoordinatorState,
msg: SingletonDrainTimeout,
) {
let still_draining = state
.singletons
.get(&msg.label)
.and_then(|rt| rt.pending_handoff.as_ref())
.map(|h| h.old_generation == msg.drained_generation)
.unwrap_or(false);
if still_draining {
tracing::warn!(
"Singleton {} drain timed out — force-proceeding; old owner abandoned (strictly-lower-fenced)",
msg.label
);
state.complete_handoff(&msg.label, msg.drained_generation);
}
}
}
impl CoordinatorState {
pub fn new(
local_node_id: impl Into<String>,
placement_strategy: Box<dyn PlacementStrategy>,
election: Box<dyn LeaderElection>,
) -> Self {
Self {
cluster_view: ClusterView::new(),
specs: BTreeMap::new(),
placement_strategy,
election,
local_node_id: local_node_id.into(),
pending_spawns: HashMap::new(),
waiting_for_return: BTreeMap::new(),
singletons: BTreeMap::new(),
generation_source: Arc::new(CoordinatorGenerationSource::new()),
next_request_id: 0,
spawn_sender: None,
singleton_stop_sender: None,
}
}
pub(crate) fn with_spawn_sender(mut self, sender: SpawnSender) -> Self {
self.spawn_sender = Some(sender);
self
}
pub(crate) fn with_singleton_stop_sender(mut self, sender: SingletonStopSender) -> Self {
self.singleton_stop_sender = Some(sender);
self
}
pub fn with_generation_source(mut self, source: Arc<dyn GenerationSource>) -> Self {
self.generation_source = source;
self
}
fn next_request_id(&mut self) -> u64 {
let id = self.next_request_id;
self.next_request_id += 1;
id
}
pub fn is_leader(&self) -> bool {
self.election
.elect(&self.cluster_view)
.is_some_and(|leader| leader == self.local_node_id)
}
fn send_spawn(&self, target_node_id: &str, request_id: u64, spec: &ActorSpec) {
if let Some(sender) = &self.spawn_sender {
sender.send_spawn(
target_node_id,
SpawnRequest {
request_id,
label: spec.label.clone(),
actor_type_name: spec.actor_type_name.clone(),
initial_state: spec.initial_state.clone(),
singleton_generation: None,
},
);
} else {
tracing::warn!(
"No spawn sender configured — spawn request for {} dropped",
spec.label
);
}
}
fn send_singleton_spawn(
&self,
target_node_id: &str,
request_id: u64,
spec: &SingletonSpec,
generation: u64,
) {
if let Some(sender) = &self.spawn_sender {
sender.send_spawn(
target_node_id,
SpawnRequest {
request_id,
label: spec.label.clone(),
actor_type_name: spec.actor_type_name.clone(),
initial_state: spec.initial_state.clone(),
singleton_generation: Some(generation),
},
);
} else {
tracing::warn!(
"No spawn sender configured — singleton spawn for {} dropped",
spec.label
);
}
}
async fn load_singletons_from_backend(&mut self) {
if !self.is_leader() {
return;
}
let records = match self.generation_source.list().await {
Ok(records) => records,
Err(e) => {
tracing::warn!("singleton leader-rebuild: list failed: {e}");
return;
}
};
for rec in records {
self.singletons
.entry(rec.spec.label.clone())
.or_insert_with(|| SingletonRuntime {
spec: rec.spec,
ownership: rec.ownership,
pending_handoff: None,
});
}
}
async fn redrive_singletons_after_loss(&mut self, lost_node_id: &str) {
let affected: Vec<String> = self
.singletons
.iter()
.filter(|(_, rt)| {
rt.ownership
.as_ref()
.and_then(|o| o.owner_node_id.as_deref())
== Some(lost_node_id)
})
.map(|(label, _)| label.clone())
.collect();
for label in affected {
let Some(spec) = self.singletons.get(&label).map(|rt| rt.spec.clone()) else {
continue;
};
let Some(new_owner) =
resolve_owner(&spec.anchor, &self.cluster_view, self.election.as_ref())
else {
tracing::warn!(
"Singleton {label} lost owner {lost_node_id} — no eligible node remains"
);
if let Some(rt) = self.singletons.get_mut(&label)
&& let Some(ownership) = &mut rt.ownership
{
ownership.owner_node_id = None;
}
continue;
};
let gen_source = self.generation_source.clone();
let generation = match gen_source.claim_term(&label, &new_owner).await {
Ok(generation) => generation,
Err(e) => {
tracing::error!("claim_term failed re-driving singleton {label}: {e}");
continue;
}
};
if let Some(rt) = self.singletons.get_mut(&label) {
rt.ownership = Some(SingletonOwnership {
label: label.clone(),
owner_node_id: Some(new_owner.clone()),
generation,
phase: SingletonPhase::Starting,
});
}
let request_id = self.next_request_id();
self.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: label.clone(),
target_node_id: new_owner.clone(),
owner: Some(OwnerKey::Singleton(label.clone())),
},
);
self.send_singleton_spawn(&new_owner, request_id, &spec, generation.packed());
tracing::info!(
"Re-drove singleton {label} to {new_owner} after loss of {lost_node_id} (gen={generation:?})"
);
}
}
fn send_stop_singleton(&self, target_node_id: &str, label: &str, generation: u64) {
if let Some(sender) = &self.singleton_stop_sender {
sender.send_stop(SingletonStopRequest {
target_node_id: target_node_id.to_string(),
label: label.to_string(),
generation,
});
} else {
tracing::warn!("No singleton stop sender configured — stop for {label} dropped");
}
}
fn schedule_drain_timeout(
&self,
ctx: &ActorContext<Coordinator>,
label: &str,
drained_generation: SingletonGeneration,
drain_timeout: Option<Duration>,
) {
if let Some(duration) = drain_timeout {
let endpoint = ctx.endpoint();
let label = label.to_string();
let runtime = ctx.receptionist().runtime().clone();
ctx.spawn(async move {
runtime.sleep(duration).await;
let _ = endpoint
.send(SingletonDrainTimeout {
label,
drained_generation,
})
.await;
});
}
}
fn complete_handoff(&mut self, label: &str, stopped_generation: SingletonGeneration) {
let Some(handoff) = self
.singletons
.get(label)
.and_then(|rt| rt.pending_handoff.clone())
else {
return; };
if handoff.old_generation != stopped_generation {
tracing::debug!(
"Ignoring stale stopped-ack for singleton {label}: {stopped_generation:?} != draining {:?}",
handoff.old_generation
);
return;
}
match handoff.destination {
None => {
self.singletons.remove(label);
tracing::info!("Singleton {label} torn down (owner drained)");
}
Some(dest) => {
let spec = self.singletons.get(label).map(|rt| rt.spec.clone());
if let Some(rt) = self.singletons.get_mut(label) {
rt.pending_handoff = None;
rt.ownership = Some(SingletonOwnership {
label: label.to_string(),
owner_node_id: Some(dest.new_owner.clone()),
generation: dest.new_generation,
phase: SingletonPhase::Starting,
});
}
let request_id = self.next_request_id();
self.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: label.to_string(),
target_node_id: dest.new_owner.clone(),
owner: Some(OwnerKey::Singleton(label.to_string())),
},
);
if let Some(spec) = spec {
self.send_singleton_spawn(
&dest.new_owner,
request_id,
&spec,
dest.new_generation.packed(),
);
}
tracing::info!(
"Singleton {label} drained — placing new owner {} (gen={:?})",
dest.new_owner,
dest.new_generation
);
}
}
}
fn handle_node_departure(
&mut self,
node_id: &str,
graceful: bool,
) -> Vec<(String, std::time::Duration)> {
let mut timers_needed = Vec::new();
let affected_labels: Vec<String> = self
.specs
.keys()
.filter(|label| {
let running_on = self
.cluster_view
.find_actor(label)
.is_some_and(|n| n == node_id);
let pending_on = self
.pending_spawns
.values()
.any(|ps| ps.spec_label == **label && ps.target_node_id == node_id);
running_on || pending_on
})
.cloned()
.collect();
self.pending_spawns
.retain(|_, ps| ps.target_node_id != node_id);
for label in affected_labels {
let spec = match self.specs.get(&label) {
Some(s) => s.clone(),
None => continue,
};
self.cluster_view.remove_actor(node_id, &label);
match &spec.crash_strategy {
CrashStrategy::Abandon => {
tracing::info!("Abandoning {label} (node {node_id} departed)");
self.specs.remove(&label);
}
CrashStrategy::WaitForReturn(duration) if !graceful => {
tracing::warn!(
"Waiting {:?} for node {} to return (spec: {label})",
duration,
node_id,
);
self.waiting_for_return.insert(
label.clone(),
WaitingSpec {
spec: spec.clone(),
failed_node_id: node_id.to_string(),
},
);
timers_needed.push((label, *duration));
}
_ => {
let reason = if graceful {
"graceful departure"
} else {
"failure"
};
tracing::info!("Redistributing {label} (node {node_id} {reason})");
if let Some(decision) = placement::select_node(
self.placement_strategy.as_ref(),
&spec,
&self.cluster_view,
) {
let request_id = self.next_request_id();
self.pending_spawns.insert(
request_id,
PendingSpawn {
spec_label: label.clone(),
target_node_id: decision.node_id.clone(),
owner: None,
},
);
self.send_spawn(&decision.node_id, request_id, &spec);
tracing::info!(
"Re-placed {label} on {} ({})",
decision.node_id,
decision.reason
);
} else {
tracing::warn!("No eligible node for {label} — spec remains unplaced");
}
}
}
}
timers_needed
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::app::election::OldestNode;
use crate::app::placement::LeastLoaded;
use crate::app::singleton::{SingletonAnchor, SingletonGeneration, SingletonSpec};
use crate::cluster::config::{NodeClass, NodeIdentity};
fn make_coordinator_state() -> CoordinatorState {
let alpha_identity = NodeIdentity::for_test("alpha", 1);
let local_node_id = alpha_identity.node_id_string();
let mut state = CoordinatorState::new(
&local_node_id,
Box::new(LeastLoaded),
Box::new(OldestNode::any()),
);
state.cluster_view.upsert_node(NodeInfo::new(
alpha_identity,
NodeClass::Worker,
HashMap::new(),
));
state.cluster_view.upsert_node(NodeInfo::new(
NodeIdentity {
name: "beta".into(),
endpoint_id: NodeIdentity::test_endpoint_id("beta"),
host: "127.0.0.1".into(),
port: 7200,
incarnation: 2,
},
NodeClass::Worker,
HashMap::new(),
));
state
}
fn make_system_and_coordinator() -> (crate::System, Endpoint<Coordinator>) {
let system = crate::System::local();
let ep = system.start("coordinator", Coordinator, make_coordinator_state());
(system, ep)
}
#[derive(Default)]
struct FailingPutSpec {
inner: CoordinatorGenerationSource,
}
#[async_trait::async_trait]
impl GenerationSource for FailingPutSpec {
async fn claim_term(
&self,
label: &str,
owner_node_id: &str,
) -> Result<SingletonGeneration, String> {
self.inner.claim_term(label, owner_node_id).await
}
async fn claim_seq(&self, label: &str) -> Result<SingletonGeneration, String> {
self.inner.claim_seq(label).await
}
async fn current(&self, label: &str) -> Result<Option<SingletonOwnership>, String> {
self.inner.current(label).await
}
async fn put_spec(&self, _label: &str, _spec: &SingletonSpec) -> Result<(), String> {
Err("simulated backend write failure".into())
}
async fn list(&self) -> Result<Vec<crate::app::singleton::SingletonRecord>, String> {
self.inner.list().await
}
}
#[tokio::test]
async fn test_start_singleton_fails_grant_when_put_spec_fails() {
let system = crate::System::local();
let state =
make_coordinator_state().with_generation_source(Arc::new(FailingPutSpec::default()));
let coordinator = system.start("coordinator", Coordinator, state);
let result = coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new("catalog", "app::Catalog", SingletonAnchor::Leader),
})
.await
.unwrap();
assert!(
result.is_err(),
"put_spec failure must fail the grant, got {result:?}"
);
let queried = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap();
assert!(
queried.is_none(),
"a failed grant must not leave a placed singleton, got {queried:?}"
);
}
#[tokio::test]
async fn test_submit_spec() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker"),
})
.await
.unwrap();
assert!(result.is_ok());
let decision = result.unwrap();
assert!(!decision.node_id.is_empty());
}
#[tokio::test]
async fn test_submit_duplicate_spec() {
let (_system, coordinator) = make_system_and_coordinator();
let _ = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker"),
})
.await
.unwrap();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker"),
})
.await
.unwrap();
assert!(result.is_err());
}
#[tokio::test]
async fn test_remove_spec() {
let (_system, coordinator) = make_system_and_coordinator();
let _ = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker"),
})
.await
.unwrap();
let result = coordinator
.send(RemoveSpec {
label: "worker/0".into(),
})
.await
.unwrap();
assert!(result.is_ok());
}
#[tokio::test]
async fn test_get_specs() {
let (_system, coordinator) = make_system_and_coordinator();
let _ = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker"),
})
.await
.unwrap();
let _ = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/1", "app::Worker"),
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 2);
}
#[tokio::test]
async fn test_node_failure_redistributes() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker")
.with_crash_strategy(CrashStrategy::Redistribute),
})
.await
.unwrap()
.unwrap();
let original_node = result.node_id.clone();
coordinator
.send_async(NotifyNodeFailed {
node_id: original_node.clone(),
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 1);
assert!(matches!(specs[0].state, SpecState::Pending));
coordinator
.send(NotifySpawnAck {
request_id: 1,
success: true,
error: None,
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 1);
assert!(matches!(specs[0].state, SpecState::Running));
assert_ne!(specs[0].placed_on.as_deref(), Some(original_node.as_str()));
}
#[tokio::test]
async fn test_node_failure_abandon() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker")
.with_crash_strategy(CrashStrategy::Abandon),
})
.await
.unwrap()
.unwrap();
coordinator
.send_async(NotifyNodeFailed {
node_id: result.node_id,
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 0);
}
#[tokio::test]
async fn test_get_cluster_view() {
let (_system, coordinator) = make_system_and_coordinator();
let view = coordinator.send(GetClusterView).await.unwrap();
assert_eq!(view.alive_count, 2);
assert_eq!(view.total_count, 2);
assert_eq!(view.nodes.len(), 2);
}
#[tokio::test]
async fn test_wait_for_return_timeout_redistributes() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker").with_crash_strategy(
CrashStrategy::WaitForReturn(std::time::Duration::from_millis(50)),
),
})
.await
.unwrap()
.unwrap();
let original_node = result.node_id.clone();
coordinator
.send_async(NotifyNodeFailed {
node_id: original_node.clone(),
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert!(matches!(specs[0].state, SpecState::WaitingForReturn));
coordinator
.send(WaitForReturnTimeout {
label: "worker/0".into(),
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 1);
assert!(!matches!(specs[0].state, SpecState::WaitingForReturn));
}
#[tokio::test]
async fn test_wait_for_return_node_rejoins() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send(SubmitSpec {
spec: ActorSpec::new("worker/0", "app::Worker").with_crash_strategy(
CrashStrategy::WaitForReturn(std::time::Duration::from_secs(60)),
),
})
.await
.unwrap()
.unwrap();
let original_node = result.node_id.clone();
coordinator
.send_async(NotifyNodeFailed {
node_id: original_node.clone(),
})
.await
.unwrap();
coordinator
.send(NotifyNodeJoined {
node_id: original_node.clone(),
info: SerializableNodeInfo {
name: "alpha".into(),
endpoint_id: crate::cluster::config::NodeIdentity::test_endpoint_id("alpha"),
host: "127.0.0.1".into(),
port: 7100,
incarnation: 1,
class: NodeClass::Worker,
metadata: HashMap::new(),
},
})
.await
.unwrap();
let specs = coordinator.send(GetSpecs).await.unwrap();
assert_eq!(specs.len(), 1);
assert!(matches!(specs[0].state, SpecState::Pending));
}
#[tokio::test]
async fn test_start_singleton_grants_ownership_on_leader() {
let (_system, coordinator) = make_system_and_coordinator();
let ownership = coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new("catalog", "app::Catalog", SingletonAnchor::Leader),
})
.await
.unwrap()
.unwrap();
assert!(
ownership.owner_node_id.as_deref().unwrap().contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("alpha").to_string()
),
"leader anchor should resolve to alpha, got {:?}",
ownership.owner_node_id
);
assert_eq!(
ownership.generation,
SingletonGeneration { term: 1, seq: 0 }
);
assert_eq!(ownership.phase, SingletonPhase::Starting);
let queried = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert_eq!(queried.generation, ownership.generation);
}
#[tokio::test]
async fn test_singleton_ack_promotes_to_active() {
let (_system, coordinator) = make_system_and_coordinator();
coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new("catalog", "app::Catalog", SingletonAnchor::Leader),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
let queried = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert_eq!(queried.phase, SingletonPhase::Active);
assert_eq!(queried.generation, SingletonGeneration { term: 1, seq: 0 });
}
#[tokio::test]
async fn test_start_singleton_idempotent_reasserts_with_seq_bump() {
let (_system, coordinator) = make_system_and_coordinator();
let spec = SingletonSpec::new("catalog", "app::Catalog", SingletonAnchor::Leader);
coordinator
.send_async(StartSingleton { spec: spec.clone() })
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
let reasserted = coordinator
.send_async(StartSingleton { spec })
.await
.unwrap()
.unwrap();
assert_eq!(
reasserted.generation,
SingletonGeneration { term: 1, seq: 1 }
);
assert_eq!(reasserted.phase, SingletonPhase::Active);
}
#[tokio::test]
async fn test_start_singleton_unsatisfiable_anchor_errors() {
let (_system, coordinator) = make_system_and_coordinator();
let result = coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new(
"catalog",
"app::Catalog",
SingletonAnchor::Node("ghost@127.0.0.1:9999#0".into()),
),
})
.await
.unwrap();
assert!(result.is_err(), "pinning to a missing node should fail");
let queried = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap();
assert!(queried.is_none());
}
#[tokio::test]
async fn test_singleton_redrives_on_owner_failure_with_higher_term() {
let (_system, coordinator) = make_system_and_coordinator();
coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new("catalog", "app::Catalog", SingletonAnchor::Leader),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
let before = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
let original_owner = before.owner_node_id.clone().unwrap();
assert!(original_owner.contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("alpha").to_string()
));
assert_eq!(before.generation, SingletonGeneration { term: 1, seq: 0 });
coordinator
.send_async(NotifyNodeFailed {
node_id: original_owner.clone(),
})
.await
.unwrap();
let after = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert!(
after.owner_node_id.as_deref().unwrap().contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("beta").to_string()
),
"must move off the failed node to beta, got {:?}",
after.owner_node_id
);
assert_eq!(after.generation, SingletonGeneration { term: 2, seq: 0 });
assert!(
after.generation > before.generation,
"term must strictly increase across failover (the fence)"
);
assert_eq!(after.phase, SingletonPhase::Starting);
coordinator
.send(NotifySpawnAck {
request_id: 1,
success: true,
error: None,
})
.await
.unwrap();
let active = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert_eq!(active.phase, SingletonPhase::Active);
}
#[tokio::test]
async fn test_singleton_node_anchor_owner_loss_clears_owner() {
let (_system, coordinator) = make_system_and_coordinator();
let alpha_id = NodeIdentity::for_test("alpha", 1).node_id_string();
coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new(
"catalog",
"app::Catalog",
SingletonAnchor::Node(alpha_id.clone()),
),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
coordinator
.send_async(NotifyNodeFailed { node_id: alpha_id })
.await
.unwrap();
let after = coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert!(
after.owner_node_id.is_none(),
"a Node-pinned singleton whose anchor died has no owner, got {:?}",
after.owner_node_id
);
}
fn make_coord_with_workers() -> (crate::System, Endpoint<Coordinator>) {
let system = crate::System::local();
let coord = NodeIdentity {
name: "coord".into(),
endpoint_id: NodeIdentity::test_endpoint_id("coord"),
host: "127.0.0.1".into(),
port: 7000,
incarnation: 0, };
let local = coord.node_id_string();
let mut state =
CoordinatorState::new(&local, Box::new(LeastLoaded), Box::new(OldestNode::any()));
state.cluster_view.upsert_node(NodeInfo::new(
coord,
NodeClass::Coordinator,
HashMap::new(),
));
for (name, port, inc) in [("w1", 7201u16, 5u64), ("w2", 7202, 3)] {
state.cluster_view.upsert_node(NodeInfo::new(
NodeIdentity {
name: name.into(),
endpoint_id: NodeIdentity::test_endpoint_id(name),
host: "127.0.0.1".into(),
port,
incarnation: inc,
},
NodeClass::Worker,
HashMap::new(),
));
}
let ep = system.start("coordinator", Coordinator, state);
(system, ep)
}
async fn join_older_worker(coordinator: &Endpoint<Coordinator>) -> String {
let id = NodeIdentity {
name: "w0".into(),
endpoint_id: NodeIdentity::test_endpoint_id("w0"),
host: "127.0.0.1".into(),
port: 7200,
incarnation: 1,
}
.node_id_string();
coordinator
.send(NotifyNodeJoined {
node_id: id.clone(),
info: SerializableNodeInfo {
name: "w0".into(),
endpoint_id: NodeIdentity::test_endpoint_id("w0"),
host: "127.0.0.1".into(),
port: 7200,
incarnation: 1,
class: NodeClass::Worker,
metadata: HashMap::new(),
},
})
.await
.unwrap();
id
}
async fn start_active_catalog(coordinator: &Endpoint<Coordinator>) {
coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new(
"catalog",
"app::Catalog",
SingletonAnchor::Class(NodeClass::Worker),
),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
}
async fn get_catalog(coordinator: &Endpoint<Coordinator>) -> Option<SingletonOwnership> {
coordinator
.send(GetSingleton {
label: "catalog".into(),
})
.await
.unwrap()
}
#[tokio::test]
async fn test_singleton_graceful_move_drains_then_places_with_higher_term() {
let (_system, coordinator) = make_coord_with_workers();
start_active_catalog(&coordinator).await;
let w0 = join_older_worker(&coordinator).await;
coordinator
.send_async(MoveSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
let draining = get_catalog(&coordinator).await.unwrap();
assert_eq!(draining.phase, SingletonPhase::Draining);
assert!(
draining.owner_node_id.as_deref().unwrap().contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("w2").to_string()
),
"ownership stays on the old owner while draining"
);
assert_eq!(draining.generation, SingletonGeneration { term: 1, seq: 0 });
coordinator
.send(NotifySingletonStopped {
label: "catalog".into(),
stopped_generation: SingletonGeneration { term: 1, seq: 0 },
})
.await
.unwrap();
let placing = get_catalog(&coordinator).await.unwrap();
assert_eq!(placing.owner_node_id.as_deref(), Some(w0.as_str()));
assert_eq!(placing.generation, SingletonGeneration { term: 2, seq: 0 });
assert_eq!(placing.phase, SingletonPhase::Starting);
coordinator
.send(NotifySpawnAck {
request_id: 1,
success: true,
error: None,
})
.await
.unwrap();
assert_eq!(
get_catalog(&coordinator).await.unwrap().phase,
SingletonPhase::Active
);
}
#[tokio::test]
async fn test_singleton_graceful_move_ignores_stale_stopped_ack() {
let (_system, coordinator) = make_coord_with_workers();
start_active_catalog(&coordinator).await;
let w0 = join_older_worker(&coordinator).await;
coordinator
.send_async(MoveSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySingletonStopped {
label: "catalog".into(),
stopped_generation: SingletonGeneration { term: 99, seq: 0 },
})
.await
.unwrap();
let still = get_catalog(&coordinator).await.unwrap();
assert_eq!(still.phase, SingletonPhase::Draining);
assert!(
still.owner_node_id.as_deref().unwrap().contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("w2").to_string()
)
);
coordinator
.send(NotifySingletonStopped {
label: "catalog".into(),
stopped_generation: SingletonGeneration { term: 1, seq: 0 },
})
.await
.unwrap();
let placing = get_catalog(&coordinator).await.unwrap();
assert_eq!(placing.owner_node_id.as_deref(), Some(w0.as_str()));
assert_eq!(placing.generation, SingletonGeneration { term: 2, seq: 0 });
}
#[tokio::test]
async fn test_singleton_graceful_teardown_removes_after_drain() {
let (_system, coordinator) = make_coord_with_workers();
start_active_catalog(&coordinator).await;
coordinator
.send(StopSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
assert_eq!(
get_catalog(&coordinator).await.unwrap().phase,
SingletonPhase::Draining
);
coordinator
.send(NotifySingletonStopped {
label: "catalog".into(),
stopped_generation: SingletonGeneration { term: 1, seq: 0 },
})
.await
.unwrap();
assert!(
get_catalog(&coordinator).await.is_none(),
"a torn-down singleton is forgotten"
);
}
#[tokio::test]
async fn test_singleton_move_is_noop_when_anchor_unchanged() {
let (_system, coordinator) = make_coord_with_workers();
start_active_catalog(&coordinator).await;
coordinator
.send_async(MoveSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
let after = get_catalog(&coordinator).await.unwrap();
assert_eq!(after.phase, SingletonPhase::Active);
assert_eq!(after.generation, SingletonGeneration { term: 1, seq: 0 });
assert!(
after.owner_node_id.as_deref().unwrap().contains(
&crate::cluster::config::NodeIdentity::test_endpoint_id("w2").to_string()
)
);
}
#[tokio::test]
async fn test_singleton_drain_timeout_force_proceeds() {
let (_system, coordinator) = make_coord_with_workers();
coordinator
.send_async(StartSingleton {
spec: SingletonSpec::new(
"catalog",
"app::Catalog",
SingletonAnchor::Class(NodeClass::Worker),
)
.with_drain_timeout(Duration::from_secs(30)),
})
.await
.unwrap()
.unwrap();
coordinator
.send(NotifySpawnAck {
request_id: 0,
success: true,
error: None,
})
.await
.unwrap();
let w0 = join_older_worker(&coordinator).await;
coordinator
.send_async(MoveSingleton {
label: "catalog".into(),
})
.await
.unwrap()
.unwrap();
coordinator
.send(SingletonDrainTimeout {
label: "catalog".into(),
drained_generation: SingletonGeneration { term: 1, seq: 0 },
})
.await
.unwrap();
let placing = get_catalog(&coordinator).await.unwrap();
assert_eq!(placing.owner_node_id.as_deref(), Some(w0.as_str()));
assert_eq!(placing.generation, SingletonGeneration { term: 2, seq: 0 });
assert_eq!(placing.phase, SingletonPhase::Starting);
}
}