use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use aion::Engine;
use aion_core::ClusterEvent;
use crate::cluster_publisher::ClusterEventPublisher;
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())
}
}
pub struct OutboxSettlingAdopter {
engine: Arc<Engine>,
outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
}
impl OutboxSettlingAdopter {
#[must_use]
pub fn new(
engine: Arc<Engine>,
outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
) -> Self {
Self {
engine,
outbox_store,
}
}
}
#[async_trait::async_trait]
impl ShardAdopter for OutboxSettlingAdopter {
async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
ShardAdopter::adopt_shards(self.engine.as_ref(), shards).await?;
let Some(outbox_store) = &self.outbox_store else {
return Ok(());
};
match crate::worker::settle_terminal_outbox_rows(
self.engine.store().as_ref(),
outbox_store.as_ref(),
)
.await
{
Ok(settled) if settled.is_empty() => {}
Ok(settled) => {
tracing::info!(
?shards,
settled = settled.len(),
"adoption sweep settled stranded outbox rows for terminal workflows"
);
}
Err(error) => {
tracing::error!(
?shards,
%error,
"adoption sweep failed to settle terminal workflows' outbox rows; \
the reconciler liveness gate remains the backstop"
);
}
}
Ok(())
}
}
#[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>,
publisher: Option<Arc<ClusterEventPublisher>>,
self_node: String,
}
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,
publisher: None,
self_node: String::new(),
}
}
#[must_use]
pub fn with_publisher(
mut self,
publisher: Arc<ClusterEventPublisher>,
self_node: impl Into<String>,
) -> Self {
self.publisher = Some(publisher);
self.self_node = self_node.into();
self
}
fn emit<F>(&self, build: F)
where
F: FnOnce(aion_core::ClusterEventMeta) -> ClusterEvent,
{
if let Some(publisher) = &self.publisher {
drop(publisher.emit(build));
}
}
#[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();
let mut pending: Vec<ClusterEvent> = Vec::new();
let confirmations = self.config.confirmations;
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 {
let was_down = entry.consecutive_down > 0 || entry.adopted;
entry.consecutive_down = 0;
entry.adopted = false;
if was_down {
pending.push(ClusterEvent::PeerConnected {
meta: placeholder_meta(),
peer_name: peer.name.clone(),
forward_addr: None,
});
}
continue;
}
entry.consecutive_down = entry.consecutive_down.saturating_add(1);
let consecutive_down = entry.consecutive_down;
let confirmed = consecutive_down >= confirmations;
pending.push(ClusterEvent::PeerDisconnected {
meta: placeholder_meta(),
peer_name: peer.name.clone(),
consecutive_down,
confirmed,
});
if entry.adopted || consecutive_down < confirmations {
continue;
}
if Self::all_shards_handled_elsewhere(
self.liveness.as_ref(),
&peer.name,
&peer.owned_shards,
) {
entry.adopted = true;
let held_by = Self::live_owner_of(self.liveness.as_ref(), &peer.owned_shards)
.unwrap_or_default();
pending.push(ClusterEvent::ShardAdoptionSkipped {
meta: placeholder_meta(),
shards: peer.owned_shards.clone(),
from_peer: peer.name.clone(),
held_by,
});
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());
pending.push(ClusterEvent::ShardAdopted {
meta: placeholder_meta(),
shards: peer.owned_shards.clone(),
from_peer: peer.name.clone(),
adopted_by: self.self_node.clone(),
});
tracing::info!(
peer = %peer.name,
shards = ?peer.owned_shards,
"cluster supervisor adopted a downed peer's shards (SS-5b auto-failover)"
);
}
Err(error) => {
pending.push(ClusterEvent::ShardAdoptionFailed {
meta: placeholder_meta(),
shards: peer.owned_shards.clone(),
from_peer: peer.name.clone(),
error: error.clone(),
});
tracing::warn!(
peer = %peer.name,
shards = ?peer.owned_shards,
%error,
"cluster supervisor failed to adopt a downed peer's shards; will retry"
);
}
}
}
for event in pending {
self.emit(|meta| with_meta(event, meta));
}
adopted_now
}
fn live_owner_of(liveness: &L, shards: &[usize]) -> Option<String> {
shards.iter().find_map(|&shard| {
liveness
.read_shard_owner(shard)
.filter(|owner| liveness.peer_connected(owner))
})
}
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 self_node = self.self_node.clone();
self.emit(|meta| ClusterEvent::SupervisorStarted {
meta,
node: self_node.clone(),
});
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;
}
}
}
}
let self_node = self.self_node.clone();
self.emit(|meta| ClusterEvent::SupervisorStopped {
meta,
node: self_node.clone(),
});
}
}
fn placeholder_meta() -> aion_core::ClusterEventMeta {
aion_core::ClusterEventMeta {
cluster_seq: 0,
observed_at: chrono::Utc::now(),
}
}
fn with_meta(event: ClusterEvent, meta: aion_core::ClusterEventMeta) -> ClusterEvent {
match event {
ClusterEvent::PeerAdded {
peer_name,
forward_addr,
..
} => ClusterEvent::PeerAdded {
meta,
peer_name,
forward_addr,
},
ClusterEvent::PeerConnected {
peer_name,
forward_addr,
..
} => ClusterEvent::PeerConnected {
meta,
peer_name,
forward_addr,
},
ClusterEvent::PeerDisconnected {
peer_name,
consecutive_down,
confirmed,
..
} => ClusterEvent::PeerDisconnected {
meta,
peer_name,
consecutive_down,
confirmed,
},
ClusterEvent::ShardAdopted {
shards,
from_peer,
adopted_by,
..
} => ClusterEvent::ShardAdopted {
meta,
shards,
from_peer,
adopted_by,
},
ClusterEvent::ShardAdoptionFailed {
shards,
from_peer,
error,
..
} => ClusterEvent::ShardAdoptionFailed {
meta,
shards,
from_peer,
error,
},
ClusterEvent::ShardAdoptionSkipped {
shards,
from_peer,
held_by,
..
} => ClusterEvent::ShardAdoptionSkipped {
meta,
shards,
from_peer,
held_by,
},
other => with_meta_worker_lifecycle(other, meta),
}
}
fn with_meta_worker_lifecycle(
event: ClusterEvent,
meta: aion_core::ClusterEventMeta,
) -> ClusterEvent {
match event {
ClusterEvent::WorkerConnected {
worker_id,
namespaces,
task_queue,
transport,
node,
..
} => ClusterEvent::WorkerConnected {
meta,
worker_id,
namespaces,
task_queue,
transport,
node,
},
ClusterEvent::WorkerDisconnected {
worker_id,
namespaces,
reason,
..
} => ClusterEvent::WorkerDisconnected {
meta,
worker_id,
namespaces,
reason,
},
ClusterEvent::SupervisorStarted { node, .. } => {
ClusterEvent::SupervisorStarted { meta, node }
}
ClusterEvent::SupervisorStopped { node, .. } => {
ClusterEvent::SupervisorStopped { meta, node }
}
ClusterEvent::NamespaceCreated {
name,
created_at,
origin,
..
} => ClusterEvent::NamespaceCreated {
meta,
name,
created_at,
origin,
},
ClusterEvent::NamespacePlacementChanged {
name, placement, ..
} => ClusterEvent::NamespacePlacementChanged {
meta,
name,
placement,
},
ClusterEvent::NamespaceQuotaState {
namespace,
in_flight,
ceiling,
..
} => ClusterEvent::NamespaceQuotaState {
meta,
namespace,
in_flight,
ceiling,
},
ClusterEvent::PeerAdded { .. }
| ClusterEvent::PeerConnected { .. }
| ClusterEvent::PeerDisconnected { .. }
| ClusterEvent::ShardAdopted { .. }
| ClusterEvent::ShardAdoptionFailed { .. }
| ClusterEvent::ShardAdoptionSkipped { .. } => {
unreachable!("peer/shard variants are re-stamped by with_meta, never delegated here")
}
}
}
#[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(flavor = "multi_thread", worker_threads = 2)]
async fn outbox_settling_adopter_settles_terminal_rows_after_adoption()
-> Result<(), Box<dyn std::error::Error>> {
use aion::{EngineBuilder, RuntimeHandle, SignalRouter};
use aion_core::{Event, EventEnvelope};
use aion_store::{OutboxRow, OutboxStatus, OutboxStore, WritableEventStore, WriteToken};
use aion_store_libsql::LibSqlStore;
let db_path = std::env::temp_dir().join(format!(
"aion-adopter-settle-{}-{}.db",
std::process::id(),
uuid::Uuid::new_v4()
));
let seeder = LibSqlStore::open(db_path.clone()).await?;
let workflow_id = aion_core::WorkflowId::new_v4();
let envelope = |seq: u64| EventEnvelope {
seq,
recorded_at: chrono::Utc::now(),
workflow_id: workflow_id.clone(),
};
let events = vec![
Event::WorkflowStarted {
envelope: envelope(1),
workflow_type: String::from("dev_brief"),
input: aion_core::Payload::from_json(&serde_json::json!({}))?,
run_id: aion_core::RunId::new_v4(),
parent_run_id: None,
package_version: aion_core::PackageVersion::new("a".repeat(64)),
},
Event::WorkflowFailed {
envelope: envelope(2),
error: aion_core::WorkflowError {
message: String::from("boom"),
details: None,
},
},
];
seeder
.append(WriteToken::recorder(), &workflow_id, &events, 0)
.await?;
let row = OutboxRow::pending(
workflow_id.clone(),
0,
String::from("norn_round"),
aion_core::Payload::from_json(&serde_json::json!({}))?,
chrono::Utc::now(),
);
let dispatch_key = row.dispatch_key.clone();
seeder
.append_outbox_batch(std::slice::from_ref(&row))
.await?;
assert_eq!(seeder.claim_outbox_rows(1).await?.len(), 1);
let engine = Arc::new(
EngineBuilder::new()
.store_arc(Arc::new(LibSqlStore::open(db_path.clone()).await?))
.in_memory_visibility()
.scheduler_threads(1)
.signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
Arc::new(aion::signal::ConcreteSignalRouter::new(runtime, handoff))
as Arc<dyn SignalRouter>
})
.build()
.await?,
);
let outbox_store: Arc<dyn OutboxStore> =
Arc::new(LibSqlStore::open(db_path.clone()).await?);
let adopter =
OutboxSettlingAdopter::new(Arc::clone(&engine), Some(Arc::clone(&outbox_store)));
ShardAdopter::adopt_shards(&adopter, &[42])
.await
.map_err(|error| format!("adoption must succeed: {error}"))?;
let state = seeder
.outbox_row_state(&dispatch_key)
.await?
.ok_or("the stranded row must still exist")?;
assert_eq!(
state.status,
OutboxStatus::Cancelled,
"the adoption sweep must settle the terminal workflow's stranded row"
);
engine.shutdown()?;
Ok(())
}
#[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 tick_emits_topology_deltas_through_the_publisher()
-> Result<(), Box<dyn std::error::Error>> {
use std::num::NonZeroUsize;
use aion_core::ClusterEvent;
use futures::StreamExt;
use crate::cluster_publisher::ClusterEventPublisher;
let capacity = NonZeroUsize::new(64).ok_or("non-zero")?;
let publisher = Arc::new(ClusterEventPublisher::new(capacity));
let mut subscription = publisher.subscribe(0);
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)
.with_publisher(Arc::clone(&publisher), "node-self@127.0.0.1");
liveness.set(false);
drop(sup.tick().await);
let fired = sup.tick().await;
assert_eq!(fired, vec!["node-1@127.0.0.1".to_owned()]);
let first = next_event(&mut subscription).await?;
assert!(
matches!(
&first,
ClusterEvent::PeerDisconnected {
confirmed: false,
consecutive_down: 1,
..
}
),
"first delta must be an unconfirmed down: {first:?}"
);
let second = next_event(&mut subscription).await?;
assert!(
matches!(
&second,
ClusterEvent::PeerDisconnected {
confirmed: true,
consecutive_down: 2,
..
}
),
"second delta must be the confirmed down: {second:?}"
);
let third = next_event(&mut subscription).await?;
let ClusterEvent::ShardAdopted {
shards,
adopted_by,
from_peer,
..
} = &third
else {
return Err(format!("third delta must be ShardAdopted: {third:?}").into());
};
assert_eq!(shards, &vec![1]);
assert_eq!(adopted_by, "node-self@127.0.0.1");
assert_eq!(from_peer, "node-1@127.0.0.1");
liveness.set(true);
drop(sup.tick().await);
let recovery = next_event(&mut subscription).await?;
assert!(
matches!(&recovery, ClusterEvent::PeerConnected { .. }),
"recovery delta must be PeerConnected: {recovery:?}"
);
let quiet = sup.tick().await;
assert!(quiet.is_empty());
assert!(
tokio::time::timeout(std::time::Duration::from_millis(50), subscription.next())
.await
.is_err(),
"a steady connected peer must not re-emit PeerConnected every tick"
);
Ok(())
}
async fn next_event(
subscription: &mut futures::stream::BoxStream<
'static,
Result<aion_core::ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
>,
) -> Result<aion_core::ClusterEvent, Box<dyn std::error::Error>> {
use futures::StreamExt;
tokio::time::timeout(std::time::Duration::from_secs(1), subscription.next())
.await?
.ok_or("cluster subscription ended")?
.map_err(|lag| format!("unexpected lag: {lag:?}").into())
}
#[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]]);
}
}