use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use aion::Engine;
pub trait PeerLiveness: Send + Sync + 'static {
fn peer_connected(&self, peer_name: &str) -> bool;
fn read_shard_owner(&self, shard: usize) -> Option<String>;
}
#[cfg(feature = "haematite-backend")]
impl PeerLiveness for aion_store_haematite::HaematiteStore {
fn peer_connected(&self, peer_name: &str) -> bool {
Self::peer_connected(self, peer_name)
}
fn read_shard_owner(&self, shard: usize) -> Option<String> {
Self::read_shard_owner(self, shard).ok().flatten()
}
}
#[async_trait::async_trait]
pub trait ShardAdopter: Send + Sync + 'static {
async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String>;
}
#[async_trait::async_trait]
impl ShardAdopter for Engine {
async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
Engine::adopt_shards(self, shards)
.await
.map_err(|error| error.to_string())
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct WatchedPeer {
pub name: String,
pub owned_shards: Vec<usize>,
}
#[derive(Clone, Copy, Debug)]
pub struct SupervisorConfig {
pub poll_interval: Duration,
pub confirmations: u32,
}
#[derive(Default)]
struct PeerState {
consecutive_down: u32,
adopted: bool,
}
pub struct ClusterSupervisor<L: PeerLiveness, A: ShardAdopter> {
liveness: Arc<L>,
adopter: Arc<A>,
peers: Vec<WatchedPeer>,
config: SupervisorConfig,
state: BTreeMap<String, PeerState>,
}
impl<L: PeerLiveness, A: ShardAdopter> ClusterSupervisor<L, A> {
#[must_use]
pub fn new(
liveness: Arc<L>,
adopter: Arc<A>,
peers: Vec<WatchedPeer>,
config: SupervisorConfig,
) -> Self {
let peers: Vec<WatchedPeer> = peers
.into_iter()
.filter(|peer| !peer.owned_shards.is_empty())
.collect();
let state = peers
.iter()
.map(|peer| (peer.name.clone(), PeerState::default()))
.collect();
Self {
liveness,
adopter,
peers,
config,
state,
}
}
#[must_use]
pub fn watches_any(&self) -> bool {
!self.peers.is_empty()
}
#[must_use]
pub fn adopter(&self) -> &A {
&self.adopter
}
pub async fn tick(&mut self) -> Vec<String> {
let mut adopted_now = Vec::new();
for peer in &self.peers {
let connected = self.liveness.peer_connected(&peer.name);
let entry = self.state.entry(peer.name.clone()).or_default();
if connected {
entry.consecutive_down = 0;
entry.adopted = false;
continue;
}
entry.consecutive_down = entry.consecutive_down.saturating_add(1);
if entry.adopted || entry.consecutive_down < self.config.confirmations {
continue;
}
if Self::all_shards_handled_elsewhere(
self.liveness.as_ref(),
&peer.name,
&peer.owned_shards,
) {
entry.adopted = true;
tracing::info!(
peer = %peer.name,
shards = ?peer.owned_shards,
"downed peer's shards already adopted by another live owner; skipping"
);
continue;
}
match self.adopter.adopt_shards(&peer.owned_shards).await {
Ok(()) => {
entry.adopted = true;
adopted_now.push(peer.name.clone());
tracing::info!(
peer = %peer.name,
shards = ?peer.owned_shards,
"cluster supervisor adopted a downed peer's shards (SS-5b auto-failover)"
);
}
Err(error) => {
tracing::warn!(
peer = %peer.name,
shards = ?peer.owned_shards,
%error,
"cluster supervisor failed to adopt a downed peer's shards; will retry"
);
}
}
}
adopted_now
}
fn all_shards_handled_elsewhere(liveness: &L, peer_name: &str, shards: &[usize]) -> bool {
!shards.is_empty()
&& shards.iter().all(|&shard| {
liveness.read_shard_owner(shard).is_some_and(|owner| {
owner != peer_name && liveness.peer_connected(&owner)
})
})
}
pub async fn run(mut self, mut shutdown: tokio::sync::watch::Receiver<bool>) {
let mut interval = tokio::time::interval(self.config.poll_interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = interval.tick() => {
drop(self.tick().await);
}
changed = shutdown.changed() => {
if changed.is_err() || *shutdown.borrow() {
break;
}
}
}
}
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use super::*;
struct FakeLiveness {
connected: AtomicBool,
owners: Mutex<std::collections::BTreeMap<usize, String>>,
live_owners: Mutex<std::collections::BTreeSet<String>>,
}
impl FakeLiveness {
fn new(connected: bool) -> Self {
Self {
connected: AtomicBool::new(connected),
owners: Mutex::new(std::collections::BTreeMap::new()),
live_owners: Mutex::new(std::collections::BTreeSet::new()),
}
}
fn set(&self, connected: bool) {
self.connected.store(connected, Ordering::SeqCst);
}
fn publish(&self, shard: usize, owner: &str, live: bool) {
self.owners
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(shard, owner.to_owned());
if live {
self.live_owners
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(owner.to_owned());
}
}
}
impl PeerLiveness for FakeLiveness {
fn peer_connected(&self, peer_name: &str) -> bool {
if self
.live_owners
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(peer_name)
{
return true;
}
self.connected.load(Ordering::SeqCst)
}
fn read_shard_owner(&self, shard: usize) -> Option<String> {
self.owners
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&shard)
.cloned()
}
}
struct FakeAdopter {
calls: Mutex<Vec<Vec<usize>>>,
fail_first: AtomicBool,
}
impl FakeAdopter {
fn new(fail_first: bool) -> Self {
Self {
calls: Mutex::new(Vec::new()),
fail_first: AtomicBool::new(fail_first),
}
}
fn calls(&self) -> Vec<Vec<usize>> {
self.calls
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
#[async_trait::async_trait]
impl ShardAdopter for FakeAdopter {
async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
if self.fail_first.swap(false, Ordering::SeqCst) {
return Err("simulated election failure".to_owned());
}
self.calls
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(shards.to_vec());
Ok(())
}
}
fn supervisor(
liveness: Arc<FakeLiveness>,
adopter: Arc<FakeAdopter>,
confirmations: u32,
) -> ClusterSupervisor<FakeLiveness, FakeAdopter> {
ClusterSupervisor::new(
liveness,
adopter,
vec![WatchedPeer {
name: "node-1@127.0.0.1".to_owned(),
owned_shards: vec![1],
}],
SupervisorConfig {
poll_interval: Duration::from_millis(1),
confirmations,
},
)
}
#[tokio::test]
async fn does_not_adopt_while_peer_connected() {
let liveness = Arc::new(FakeLiveness::new(true));
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2);
for _ in 0..5 {
assert!(sup.tick().await.is_empty());
}
assert!(adopter.calls().is_empty(), "no adoption while peer is up");
}
#[tokio::test]
async fn debounce_requires_consecutive_down_before_adopting() {
let liveness = Arc::new(FakeLiveness::new(true));
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 3);
liveness.set(false);
assert!(sup.tick().await.is_empty(), "tick 1 down: below threshold");
liveness.set(true);
assert!(sup.tick().await.is_empty());
liveness.set(false);
assert!(
sup.tick().await.is_empty(),
"down again, counter reset to 1"
);
assert!(sup.tick().await.is_empty(), "2 consecutive: still below 3");
let fired = sup.tick().await;
assert_eq!(
fired,
vec!["node-1@127.0.0.1".to_owned()],
"3rd consecutive triggers"
);
assert_eq!(adopter.calls(), vec![vec![1]]);
}
#[tokio::test]
async fn adopts_once_then_stays_quiet_while_down() {
let liveness = Arc::new(FakeLiveness::new(false));
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
assert_eq!(sup.tick().await.len(), 1, "first down tick adopts");
for _ in 0..5 {
assert!(sup.tick().await.is_empty(), "no re-adopt while still down");
}
assert_eq!(adopter.calls(), vec![vec![1]], "adopted exactly once");
}
#[tokio::test]
async fn failed_adoption_is_retried_next_tick() {
let liveness = Arc::new(FakeLiveness::new(false));
let adopter = Arc::new(FakeAdopter::new(true)); let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
assert!(
sup.tick().await.is_empty(),
"first adopt fails, not recorded"
);
assert!(adopter.calls().is_empty());
assert_eq!(sup.tick().await.len(), 1, "retry succeeds next tick");
assert_eq!(adopter.calls(), vec![vec![1]]);
}
#[tokio::test]
async fn peer_with_no_shards_is_not_watched() {
let liveness = Arc::new(FakeLiveness::new(false));
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = ClusterSupervisor::new(
Arc::clone(&liveness),
Arc::clone(&adopter),
vec![WatchedPeer {
name: "node-2@127.0.0.1".to_owned(),
owned_shards: vec![],
}],
SupervisorConfig {
poll_interval: Duration::from_millis(1),
confirmations: 1,
},
);
assert!(!sup.watches_any());
assert!(sup.tick().await.is_empty());
assert!(adopter.calls().is_empty());
}
#[tokio::test]
async fn shard_already_published_to_live_owner_is_not_adopted() {
let liveness = Arc::new(FakeLiveness::new(false));
liveness.publish(1, "node-9@127.0.0.1", true);
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
assert!(
sup.tick().await.is_empty(),
"no adoption fires for a shard a live owner already holds"
);
assert!(
adopter.calls().is_empty(),
"the adopter is never invoked for an already-handled shard"
);
for _ in 0..3 {
assert!(sup.tick().await.is_empty());
}
assert!(adopter.calls().is_empty());
}
#[tokio::test]
async fn shard_published_to_a_down_owner_is_still_adopted() {
let liveness = Arc::new(FakeLiveness::new(false));
liveness.publish(1, "node-9@127.0.0.1", false);
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
assert_eq!(
sup.tick().await.len(),
1,
"a shard whose recorded owner is itself down is adoptable"
);
assert_eq!(adopter.calls(), vec![vec![1]]);
}
#[tokio::test]
async fn shard_published_to_the_dead_peer_itself_is_adopted() {
let liveness = Arc::new(FakeLiveness::new(false));
liveness.publish(1, "node-1@127.0.0.1", false);
let adopter = Arc::new(FakeAdopter::new(false));
let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
assert_eq!(
sup.tick().await.len(),
1,
"a record naming the dead peer itself is stale and still adoptable"
);
assert_eq!(adopter.calls(), vec![vec![1]]);
}
}