1use std::collections::BTreeMap;
35use std::sync::Arc;
36use std::time::Duration;
37
38use aion::Engine;
39use aion_core::ClusterEvent;
40
41use crate::cluster_publisher::ClusterEventPublisher;
42
43pub trait PeerLiveness: Send + Sync + 'static {
47 fn peer_connected(&self, peer_name: &str) -> bool;
49
50 fn read_shard_owner(&self, shard: usize) -> Option<String>;
61}
62
63impl PeerLiveness for aion_store_haematite::HaematiteStore {
64 fn peer_connected(&self, peer_name: &str) -> bool {
65 Self::peer_connected(self, peer_name)
66 }
67
68 fn read_shard_owner(&self, shard: usize) -> Option<String> {
69 Self::read_shard_owner(self, shard).ok().flatten()
71 }
72}
73
74#[async_trait::async_trait]
77pub trait ShardAdopter: Send + Sync + 'static {
78 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String>;
80}
81
82#[async_trait::async_trait]
83impl ShardAdopter for Engine {
84 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
85 Engine::adopt_shards(self, shards)
86 .await
87 .map_err(|error| error.to_string())
88 }
89}
90
91pub struct OutboxSettlingAdopter {
104 engine: Arc<Engine>,
105 outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
106}
107
108impl OutboxSettlingAdopter {
109 #[must_use]
113 pub fn new(
114 engine: Arc<Engine>,
115 outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
116 ) -> Self {
117 Self {
118 engine,
119 outbox_store,
120 }
121 }
122}
123
124#[async_trait::async_trait]
125impl ShardAdopter for OutboxSettlingAdopter {
126 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
127 ShardAdopter::adopt_shards(self.engine.as_ref(), shards).await?;
128 let Some(outbox_store) = &self.outbox_store else {
129 return Ok(());
130 };
131 match crate::worker::settle_terminal_outbox_rows(
132 self.engine.store().as_ref(),
133 outbox_store.as_ref(),
134 )
135 .await
136 {
137 Ok(settled) if settled.is_empty() => {}
138 Ok(settled) => {
139 tracing::info!(
140 ?shards,
141 settled = settled.len(),
142 "adoption sweep settled stranded outbox rows for terminal workflows"
143 );
144 }
145 Err(error) => {
146 tracing::error!(
147 ?shards,
148 %error,
149 "adoption sweep failed to settle terminal workflows' outbox rows; \
150 the reconciler liveness gate remains the backstop"
151 );
152 }
153 }
154 Ok(())
155 }
156}
157
158#[derive(Clone, Debug, PartialEq, Eq)]
161pub struct WatchedPeer {
162 pub name: String,
164 pub owned_shards: Vec<usize>,
166}
167
168#[derive(Clone, Copy, Debug)]
170pub struct SupervisorConfig {
171 pub poll_interval: Duration,
173 pub confirmations: u32,
176}
177
178#[derive(Default)]
180struct PeerState {
181 consecutive_down: u32,
183 adopted: bool,
185}
186
187pub struct ClusterSupervisor<L: PeerLiveness, A: ShardAdopter> {
189 liveness: Arc<L>,
190 adopter: Arc<A>,
191 peers: Vec<WatchedPeer>,
192 config: SupervisorConfig,
193 state: BTreeMap<String, PeerState>,
194 publisher: Option<Arc<ClusterEventPublisher>>,
199 self_node: String,
203}
204
205impl<L: PeerLiveness, A: ShardAdopter> ClusterSupervisor<L, A> {
206 #[must_use]
210 pub fn new(
211 liveness: Arc<L>,
212 adopter: Arc<A>,
213 peers: Vec<WatchedPeer>,
214 config: SupervisorConfig,
215 ) -> Self {
216 let peers: Vec<WatchedPeer> = peers
217 .into_iter()
218 .filter(|peer| !peer.owned_shards.is_empty())
219 .collect();
220 let state = peers
221 .iter()
222 .map(|peer| (peer.name.clone(), PeerState::default()))
223 .collect();
224 Self {
225 liveness,
226 adopter,
227 peers,
228 config,
229 state,
230 publisher: None,
231 self_node: String::new(),
232 }
233 }
234
235 #[must_use]
239 pub fn with_publisher(
240 mut self,
241 publisher: Arc<ClusterEventPublisher>,
242 self_node: impl Into<String>,
243 ) -> Self {
244 self.publisher = Some(publisher);
245 self.self_node = self_node.into();
246 self
247 }
248
249 fn emit<F>(&self, build: F)
253 where
254 F: FnOnce(aion_core::ClusterEventMeta) -> ClusterEvent,
255 {
256 if let Some(publisher) = &self.publisher {
257 drop(publisher.emit(build));
258 }
259 }
260
261 #[must_use]
264 pub fn watches_any(&self) -> bool {
265 !self.peers.is_empty()
266 }
267
268 #[must_use]
271 pub fn adopter(&self) -> &A {
272 &self.adopter
273 }
274
275 pub async fn tick(&mut self) -> Vec<String> {
283 let mut adopted_now = Vec::new();
284 let mut pending: Vec<ClusterEvent> = Vec::new();
288 let confirmations = self.config.confirmations;
289 for peer in &self.peers {
290 let connected = self.liveness.peer_connected(&peer.name);
291 let entry = self.state.entry(peer.name.clone()).or_default();
292 if connected {
293 let was_down = entry.consecutive_down > 0 || entry.adopted;
297 entry.consecutive_down = 0;
298 entry.adopted = false;
299 if was_down {
300 pending.push(ClusterEvent::PeerConnected {
301 meta: placeholder_meta(),
302 peer_name: peer.name.clone(),
303 forward_addr: None,
304 });
305 }
306 continue;
307 }
308 entry.consecutive_down = entry.consecutive_down.saturating_add(1);
309 let consecutive_down = entry.consecutive_down;
310 let confirmed = consecutive_down >= confirmations;
311 pending.push(ClusterEvent::PeerDisconnected {
314 meta: placeholder_meta(),
315 peer_name: peer.name.clone(),
316 consecutive_down,
317 confirmed,
318 });
319 if entry.adopted || consecutive_down < confirmations {
320 continue;
321 }
322 if Self::all_shards_handled_elsewhere(
329 self.liveness.as_ref(),
330 &peer.name,
331 &peer.owned_shards,
332 ) {
333 entry.adopted = true;
336 let held_by = Self::live_owner_of(self.liveness.as_ref(), &peer.owned_shards)
337 .unwrap_or_default();
338 pending.push(ClusterEvent::ShardAdoptionSkipped {
339 meta: placeholder_meta(),
340 shards: peer.owned_shards.clone(),
341 from_peer: peer.name.clone(),
342 held_by,
343 });
344 tracing::info!(
345 peer = %peer.name,
346 shards = ?peer.owned_shards,
347 "downed peer's shards already adopted by another live owner; skipping"
348 );
349 continue;
350 }
351 match self.adopter.adopt_shards(&peer.owned_shards).await {
352 Ok(()) => {
353 entry.adopted = true;
354 adopted_now.push(peer.name.clone());
355 pending.push(ClusterEvent::ShardAdopted {
356 meta: placeholder_meta(),
357 shards: peer.owned_shards.clone(),
358 from_peer: peer.name.clone(),
359 adopted_by: self.self_node.clone(),
360 });
361 tracing::info!(
362 peer = %peer.name,
363 shards = ?peer.owned_shards,
364 "cluster supervisor adopted a downed peer's shards (SS-5b auto-failover)"
365 );
366 }
367 Err(error) => {
368 pending.push(ClusterEvent::ShardAdoptionFailed {
369 meta: placeholder_meta(),
370 shards: peer.owned_shards.clone(),
371 from_peer: peer.name.clone(),
372 error: error.clone(),
373 });
374 tracing::warn!(
383 peer = %peer.name,
384 shards = ?peer.owned_shards,
385 %error,
386 "cluster supervisor failed to adopt a downed peer's shards; will retry"
387 );
388 }
389 }
390 }
391 for event in pending {
395 self.emit(|meta| with_meta(event, meta));
396 }
397 adopted_now
398 }
399
400 fn live_owner_of(liveness: &L, shards: &[usize]) -> Option<String> {
404 shards.iter().find_map(|&shard| {
405 liveness
406 .read_shard_owner(shard)
407 .filter(|owner| liveness.peer_connected(owner))
408 })
409 }
410
411 fn all_shards_handled_elsewhere(liveness: &L, peer_name: &str, shards: &[usize]) -> bool {
419 !shards.is_empty()
420 && shards.iter().all(|&shard| {
421 liveness.read_shard_owner(shard).is_some_and(|owner| {
422 owner != peer_name && liveness.peer_connected(&owner)
426 })
427 })
428 }
429
430 pub async fn run(mut self, mut shutdown: tokio::sync::watch::Receiver<bool>) {
433 let self_node = self.self_node.clone();
436 self.emit(|meta| ClusterEvent::SupervisorStarted {
437 meta,
438 node: self_node.clone(),
439 });
440 let mut interval = tokio::time::interval(self.config.poll_interval);
441 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
442 loop {
443 tokio::select! {
444 _ = interval.tick() => {
445 drop(self.tick().await);
446 }
447 changed = shutdown.changed() => {
448 if changed.is_err() || *shutdown.borrow() {
449 break;
450 }
451 }
452 }
453 }
454 let self_node = self.self_node.clone();
457 self.emit(|meta| ClusterEvent::SupervisorStopped {
458 meta,
459 node: self_node.clone(),
460 });
461 }
462}
463
464fn placeholder_meta() -> aion_core::ClusterEventMeta {
468 aion_core::ClusterEventMeta {
469 cluster_seq: 0,
470 observed_at: chrono::Utc::now(),
471 }
472}
473
474fn with_meta(event: ClusterEvent, meta: aion_core::ClusterEventMeta) -> ClusterEvent {
481 match event {
482 ClusterEvent::PeerAdded {
483 peer_name,
484 forward_addr,
485 ..
486 } => ClusterEvent::PeerAdded {
487 meta,
488 peer_name,
489 forward_addr,
490 },
491 ClusterEvent::PeerConnected {
492 peer_name,
493 forward_addr,
494 ..
495 } => ClusterEvent::PeerConnected {
496 meta,
497 peer_name,
498 forward_addr,
499 },
500 ClusterEvent::PeerDisconnected {
501 peer_name,
502 consecutive_down,
503 confirmed,
504 ..
505 } => ClusterEvent::PeerDisconnected {
506 meta,
507 peer_name,
508 consecutive_down,
509 confirmed,
510 },
511 ClusterEvent::ShardAdopted {
512 shards,
513 from_peer,
514 adopted_by,
515 ..
516 } => ClusterEvent::ShardAdopted {
517 meta,
518 shards,
519 from_peer,
520 adopted_by,
521 },
522 ClusterEvent::ShardAdoptionFailed {
523 shards,
524 from_peer,
525 error,
526 ..
527 } => ClusterEvent::ShardAdoptionFailed {
528 meta,
529 shards,
530 from_peer,
531 error,
532 },
533 ClusterEvent::ShardAdoptionSkipped {
534 shards,
535 from_peer,
536 held_by,
537 ..
538 } => ClusterEvent::ShardAdoptionSkipped {
539 meta,
540 shards,
541 from_peer,
542 held_by,
543 },
544 other => with_meta_worker_lifecycle(other, meta),
545 }
546}
547
548fn with_meta_worker_lifecycle(
558 event: ClusterEvent,
559 meta: aion_core::ClusterEventMeta,
560) -> ClusterEvent {
561 match event {
562 ClusterEvent::WorkerConnected {
563 worker_id,
564 namespaces,
565 task_queue,
566 transport,
567 node,
568 deployment,
569 deployment_association,
570 ..
571 } => ClusterEvent::WorkerConnected {
572 meta,
573 worker_id,
574 namespaces,
575 task_queue,
576 transport,
577 node,
578 deployment,
579 deployment_association,
580 },
581 ClusterEvent::WorkerDisconnected {
582 worker_id,
583 namespaces,
584 reason,
585 ..
586 } => ClusterEvent::WorkerDisconnected {
587 meta,
588 worker_id,
589 namespaces,
590 reason,
591 },
592 ClusterEvent::SupervisorStarted { node, .. } => {
593 ClusterEvent::SupervisorStarted { meta, node }
594 }
595 ClusterEvent::SupervisorStopped { node, .. } => {
596 ClusterEvent::SupervisorStopped { meta, node }
597 }
598 ClusterEvent::NamespaceCreated {
599 name,
600 created_at,
601 origin,
602 ..
603 } => ClusterEvent::NamespaceCreated {
604 meta,
605 name,
606 created_at,
607 origin,
608 },
609 ClusterEvent::NamespacePlacementChanged {
613 name, placement, ..
614 } => ClusterEvent::NamespacePlacementChanged {
615 meta,
616 name,
617 placement,
618 },
619 ClusterEvent::NamespaceQuotaState {
623 namespace,
624 in_flight,
625 ceiling,
626 ..
627 } => ClusterEvent::NamespaceQuotaState {
628 meta,
629 namespace,
630 in_flight,
631 ceiling,
632 },
633 parked_event @ ClusterEvent::DispatchParked { .. } => {
634 with_meta_dispatch_parked(parked_event, meta)
635 }
636 deployment_event @ (ClusterEvent::WorkerDeploymentPut { .. }
637 | ClusterEvent::WorkerDeploymentDesiredStateChanged { .. }
638 | ClusterEvent::WorkerDeploymentDeleted { .. }) => {
639 with_meta_worker_deployment(deployment_event, meta)
640 }
641 ClusterEvent::PeerAdded { .. }
642 | ClusterEvent::PeerConnected { .. }
643 | ClusterEvent::PeerDisconnected { .. }
644 | ClusterEvent::ShardAdopted { .. }
645 | ClusterEvent::ShardAdoptionFailed { .. }
646 | ClusterEvent::ShardAdoptionSkipped { .. } => {
647 unreachable!("peer/shard variants are re-stamped by with_meta, never delegated here")
648 }
649 }
650}
651
652fn with_meta_dispatch_parked(
659 event: ClusterEvent,
660 meta: aion_core::ClusterEventMeta,
661) -> ClusterEvent {
662 match event {
663 ClusterEvent::DispatchParked {
664 namespace,
665 task_queue,
666 activity_type,
667 node,
668 reason,
669 policy,
670 workflow_id,
671 activity_id,
672 waited_ms,
673 workers_in_pool,
674 workers_serving_activity,
675 compatible_workers,
676 last_compatible_poller_age_ms,
677 ..
678 } => ClusterEvent::DispatchParked {
679 meta,
680 namespace,
681 task_queue,
682 activity_type,
683 node,
684 reason,
685 policy,
686 workflow_id,
687 activity_id,
688 waited_ms,
689 workers_in_pool,
690 workers_serving_activity,
691 compatible_workers,
692 last_compatible_poller_age_ms,
693 },
694 ClusterEvent::WorkerConnected { .. }
695 | ClusterEvent::WorkerDisconnected { .. }
696 | ClusterEvent::SupervisorStarted { .. }
697 | ClusterEvent::SupervisorStopped { .. }
698 | ClusterEvent::NamespaceCreated { .. }
699 | ClusterEvent::NamespacePlacementChanged { .. }
700 | ClusterEvent::NamespaceQuotaState { .. }
701 | ClusterEvent::WorkerDeploymentPut { .. }
702 | ClusterEvent::WorkerDeploymentDesiredStateChanged { .. }
703 | ClusterEvent::WorkerDeploymentDeleted { .. }
704 | ClusterEvent::PeerAdded { .. }
705 | ClusterEvent::PeerConnected { .. }
706 | ClusterEvent::PeerDisconnected { .. }
707 | ClusterEvent::ShardAdopted { .. }
708 | ClusterEvent::ShardAdoptionFailed { .. }
709 | ClusterEvent::ShardAdoptionSkipped { .. } => {
710 unreachable!("only DispatchParked is delegated here")
711 }
712 }
713}
714
715fn with_meta_worker_deployment(
716 event: ClusterEvent,
717 meta: aion_core::ClusterEventMeta,
718) -> ClusterEvent {
719 match event {
720 ClusterEvent::WorkerDeploymentPut {
721 name,
722 outcome,
723 desired_state,
724 binary_version,
725 binary_content_hash,
726 ..
727 } => ClusterEvent::WorkerDeploymentPut {
728 meta,
729 name,
730 outcome,
731 desired_state,
732 binary_version,
733 binary_content_hash,
734 },
735 ClusterEvent::WorkerDeploymentDesiredStateChanged {
736 name,
737 desired_state,
738 ..
739 } => ClusterEvent::WorkerDeploymentDesiredStateChanged {
740 meta,
741 name,
742 desired_state,
743 },
744 ClusterEvent::WorkerDeploymentDeleted { name, .. } => {
745 ClusterEvent::WorkerDeploymentDeleted { meta, name }
746 }
747 ClusterEvent::WorkerConnected { .. }
748 | ClusterEvent::WorkerDisconnected { .. }
749 | ClusterEvent::SupervisorStarted { .. }
750 | ClusterEvent::SupervisorStopped { .. }
751 | ClusterEvent::NamespaceCreated { .. }
752 | ClusterEvent::NamespacePlacementChanged { .. }
753 | ClusterEvent::NamespaceQuotaState { .. }
754 | ClusterEvent::DispatchParked { .. }
755 | ClusterEvent::PeerAdded { .. }
756 | ClusterEvent::PeerConnected { .. }
757 | ClusterEvent::PeerDisconnected { .. }
758 | ClusterEvent::ShardAdopted { .. }
759 | ClusterEvent::ShardAdoptionFailed { .. }
760 | ClusterEvent::ShardAdoptionSkipped { .. } => {
761 unreachable!("only worker-deployment variants are delegated here")
762 }
763 }
764}
765
766#[cfg(test)]
767mod tests {
768 use std::sync::Mutex;
769 use std::sync::atomic::{AtomicBool, Ordering};
770
771 use super::*;
772
773 struct FakeLiveness {
778 connected: AtomicBool,
779 owners: Mutex<std::collections::BTreeMap<usize, String>>,
781 live_owners: Mutex<std::collections::BTreeSet<String>>,
783 }
784
785 impl FakeLiveness {
786 fn new(connected: bool) -> Self {
787 Self {
788 connected: AtomicBool::new(connected),
789 owners: Mutex::new(std::collections::BTreeMap::new()),
790 live_owners: Mutex::new(std::collections::BTreeSet::new()),
791 }
792 }
793 fn set(&self, connected: bool) {
794 self.connected.store(connected, Ordering::SeqCst);
795 }
796 fn publish(&self, shard: usize, owner: &str, live: bool) {
799 self.owners
800 .lock()
801 .unwrap_or_else(std::sync::PoisonError::into_inner)
802 .insert(shard, owner.to_owned());
803 if live {
804 self.live_owners
805 .lock()
806 .unwrap_or_else(std::sync::PoisonError::into_inner)
807 .insert(owner.to_owned());
808 }
809 }
810 }
811
812 impl PeerLiveness for FakeLiveness {
813 fn peer_connected(&self, peer_name: &str) -> bool {
814 if self
815 .live_owners
816 .lock()
817 .unwrap_or_else(std::sync::PoisonError::into_inner)
818 .contains(peer_name)
819 {
820 return true;
821 }
822 self.connected.load(Ordering::SeqCst)
823 }
824
825 fn read_shard_owner(&self, shard: usize) -> Option<String> {
826 self.owners
827 .lock()
828 .unwrap_or_else(std::sync::PoisonError::into_inner)
829 .get(&shard)
830 .cloned()
831 }
832 }
833
834 struct FakeAdopter {
836 calls: Mutex<Vec<Vec<usize>>>,
837 fail_first: AtomicBool,
838 }
839
840 impl FakeAdopter {
841 fn new(fail_first: bool) -> Self {
842 Self {
843 calls: Mutex::new(Vec::new()),
844 fail_first: AtomicBool::new(fail_first),
845 }
846 }
847 fn calls(&self) -> Vec<Vec<usize>> {
848 self.calls
849 .lock()
850 .unwrap_or_else(std::sync::PoisonError::into_inner)
851 .clone()
852 }
853 }
854
855 #[async_trait::async_trait]
856 impl ShardAdopter for FakeAdopter {
857 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
858 if self.fail_first.swap(false, Ordering::SeqCst) {
859 return Err("simulated election failure".to_owned());
860 }
861 self.calls
862 .lock()
863 .unwrap_or_else(std::sync::PoisonError::into_inner)
864 .push(shards.to_vec());
865 Ok(())
866 }
867 }
868
869 fn supervisor(
870 liveness: Arc<FakeLiveness>,
871 adopter: Arc<FakeAdopter>,
872 confirmations: u32,
873 ) -> ClusterSupervisor<FakeLiveness, FakeAdopter> {
874 ClusterSupervisor::new(
875 liveness,
876 adopter,
877 vec![WatchedPeer {
878 name: "node-1@127.0.0.1".to_owned(),
879 owned_shards: vec![1],
880 }],
881 SupervisorConfig {
882 poll_interval: Duration::from_millis(1),
883 confirmations,
884 },
885 )
886 }
887
888 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
894 async fn outbox_settling_adopter_settles_terminal_rows_after_adoption()
895 -> Result<(), Box<dyn std::error::Error>> {
896 use aion::{EngineBuilder, RuntimeHandle, SignalRouter};
897 use aion_core::{Event, EventEnvelope};
898 use aion_store::{OutboxRow, OutboxStatus, OutboxStore, WritableEventStore, WriteToken};
899 use aion_store_haematite::HaematiteStore;
900
901 let db_path = std::env::temp_dir().join(format!(
902 "aion-adopter-settle-{}-{}.db",
903 std::process::id(),
904 uuid::Uuid::new_v4()
905 ));
906
907 let seeder = Arc::new(
910 HaematiteStore::open_or_create(
911 db_path.clone(),
912 haematite::NodeCacheBudget::Unlimited,
913 )
915 .await?,
916 );
917 let workflow_id = aion_core::WorkflowId::new_v4();
918 let envelope = |seq: u64| EventEnvelope {
919 seq,
920 recorded_at: chrono::Utc::now(),
921 workflow_id: workflow_id.clone(),
922 };
923 let events = vec![
924 Event::WorkflowStarted {
925 envelope: envelope(1),
926 workflow_type: String::from("dev_brief"),
927 input: aion_core::Payload::from_json(&serde_json::json!({}))?,
928 run_id: aion_core::RunId::new_v4(),
929 parent_run_id: None,
930 parent_workflow_id: None,
931 package_version: aion_core::PackageVersion::new("a".repeat(64)),
932 },
933 Event::WorkflowFailed {
934 envelope: envelope(2),
935 error: aion_core::WorkflowError {
936 message: String::from("boom"),
937 details: None,
938 },
939 },
940 ];
941 seeder
942 .append(WriteToken::recorder(), &workflow_id, &events, 0)
943 .await?;
944 let row = OutboxRow::pending(
945 workflow_id.clone(),
946 0,
947 String::from("norn_round"),
948 aion_core::Payload::from_json(&serde_json::json!({}))?,
949 chrono::Utc::now(),
950 );
951 let dispatch_key = row.dispatch_key.clone();
952 seeder
953 .append_outbox_batch(std::slice::from_ref(&row))
954 .await?;
955 assert_eq!(seeder.claim_outbox_rows(1).await?.len(), 1);
956
957 let event_store: Arc<dyn aion_store::EventStore> = Arc::clone(&seeder) as _;
958 let engine = Arc::new(
959 EngineBuilder::new()
960 .store_arc(event_store)
961 .in_memory_visibility()
962 .scheduler_threads(1)
963 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
964 Arc::new(aion::signal::ConcreteSignalRouter::new(runtime, handoff))
965 as Arc<dyn SignalRouter>
966 })
967 .build()
968 .await?,
969 );
970 let outbox_store: Arc<dyn OutboxStore> = Arc::clone(&seeder) as _;
971 let adopter =
972 OutboxSettlingAdopter::new(Arc::clone(&engine), Some(Arc::clone(&outbox_store)));
973
974 ShardAdopter::adopt_shards(&adopter, &[42])
975 .await
976 .map_err(|error| format!("adoption must succeed: {error}"))?;
977
978 let state = seeder
979 .outbox_row_state(&dispatch_key)
980 .await?
981 .ok_or("the stranded row must still exist")?;
982 assert_eq!(
983 state.status,
984 OutboxStatus::Cancelled,
985 "the adoption sweep must settle the terminal workflow's stranded row"
986 );
987 engine.shutdown()?;
988 Ok(())
989 }
990
991 #[tokio::test]
992 async fn does_not_adopt_while_peer_connected() {
993 let liveness = Arc::new(FakeLiveness::new(true));
994 let adopter = Arc::new(FakeAdopter::new(false));
995 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2);
996 for _ in 0..5 {
997 assert!(sup.tick().await.is_empty());
998 }
999 assert!(adopter.calls().is_empty(), "no adoption while peer is up");
1000 }
1001
1002 #[tokio::test]
1003 async fn debounce_requires_consecutive_down_before_adopting() {
1004 let liveness = Arc::new(FakeLiveness::new(true));
1005 let adopter = Arc::new(FakeAdopter::new(false));
1006 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 3);
1007
1008 liveness.set(false);
1009 assert!(sup.tick().await.is_empty(), "tick 1 down: below threshold");
1010 liveness.set(true);
1012 assert!(sup.tick().await.is_empty());
1013 liveness.set(false);
1014 assert!(
1015 sup.tick().await.is_empty(),
1016 "down again, counter reset to 1"
1017 );
1018 assert!(sup.tick().await.is_empty(), "2 consecutive: still below 3");
1019 let fired = sup.tick().await;
1020 assert_eq!(
1021 fired,
1022 vec!["node-1@127.0.0.1".to_owned()],
1023 "3rd consecutive triggers"
1024 );
1025 assert_eq!(adopter.calls(), vec![vec![1]]);
1026 }
1027
1028 #[tokio::test]
1029 async fn adopts_once_then_stays_quiet_while_down() {
1030 let liveness = Arc::new(FakeLiveness::new(false));
1031 let adopter = Arc::new(FakeAdopter::new(false));
1032 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1033 assert_eq!(sup.tick().await.len(), 1, "first down tick adopts");
1034 for _ in 0..5 {
1035 assert!(sup.tick().await.is_empty(), "no re-adopt while still down");
1036 }
1037 assert_eq!(adopter.calls(), vec![vec![1]], "adopted exactly once");
1038 }
1039
1040 #[tokio::test]
1041 async fn failed_adoption_is_retried_next_tick() {
1042 let liveness = Arc::new(FakeLiveness::new(false));
1043 let adopter = Arc::new(FakeAdopter::new(true)); let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1045 assert!(
1046 sup.tick().await.is_empty(),
1047 "first adopt fails, not recorded"
1048 );
1049 assert!(adopter.calls().is_empty());
1050 assert_eq!(sup.tick().await.len(), 1, "retry succeeds next tick");
1051 assert_eq!(adopter.calls(), vec![vec![1]]);
1052 }
1053
1054 #[tokio::test]
1055 async fn peer_with_no_shards_is_not_watched() {
1056 let liveness = Arc::new(FakeLiveness::new(false));
1057 let adopter = Arc::new(FakeAdopter::new(false));
1058 let mut sup = ClusterSupervisor::new(
1059 Arc::clone(&liveness),
1060 Arc::clone(&adopter),
1061 vec![WatchedPeer {
1062 name: "node-2@127.0.0.1".to_owned(),
1063 owned_shards: vec![],
1064 }],
1065 SupervisorConfig {
1066 poll_interval: Duration::from_millis(1),
1067 confirmations: 1,
1068 },
1069 );
1070 assert!(!sup.watches_any());
1071 assert!(sup.tick().await.is_empty());
1072 assert!(adopter.calls().is_empty());
1073 }
1074
1075 #[tokio::test]
1079 async fn shard_already_published_to_live_owner_is_not_adopted() {
1080 let liveness = Arc::new(FakeLiveness::new(false));
1081 liveness.publish(1, "node-9@127.0.0.1", true);
1084 let adopter = Arc::new(FakeAdopter::new(false));
1085 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1086
1087 assert!(
1089 sup.tick().await.is_empty(),
1090 "no adoption fires for a shard a live owner already holds"
1091 );
1092 assert!(
1093 adopter.calls().is_empty(),
1094 "the adopter is never invoked for an already-handled shard"
1095 );
1096 for _ in 0..3 {
1098 assert!(sup.tick().await.is_empty());
1099 }
1100 assert!(adopter.calls().is_empty());
1101 }
1102
1103 #[tokio::test]
1107 async fn shard_published_to_a_down_owner_is_still_adopted() {
1108 let liveness = Arc::new(FakeLiveness::new(false));
1109 liveness.publish(1, "node-9@127.0.0.1", false);
1112 let adopter = Arc::new(FakeAdopter::new(false));
1113 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1114
1115 assert_eq!(
1116 sup.tick().await.len(),
1117 1,
1118 "a shard whose recorded owner is itself down is adoptable"
1119 );
1120 assert_eq!(adopter.calls(), vec![vec![1]]);
1121 }
1122
1123 #[tokio::test]
1129 async fn tick_emits_topology_deltas_through_the_publisher()
1130 -> Result<(), Box<dyn std::error::Error>> {
1131 use std::num::NonZeroUsize;
1132
1133 use aion_core::ClusterEvent;
1134 use futures::StreamExt;
1135
1136 use crate::cluster_publisher::ClusterEventPublisher;
1137
1138 let capacity = NonZeroUsize::new(64).ok_or("non-zero")?;
1139 let publisher = Arc::new(ClusterEventPublisher::new(capacity));
1140 let mut subscription = publisher.subscribe(0);
1141
1142 let liveness = Arc::new(FakeLiveness::new(true));
1143 let adopter = Arc::new(FakeAdopter::new(false));
1144 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2)
1145 .with_publisher(Arc::clone(&publisher), "node-self@127.0.0.1");
1146
1147 liveness.set(false);
1149 drop(sup.tick().await);
1150 let fired = sup.tick().await;
1152 assert_eq!(fired, vec!["node-1@127.0.0.1".to_owned()]);
1153
1154 let first = next_event(&mut subscription).await?;
1156 assert!(
1157 matches!(
1158 &first,
1159 ClusterEvent::PeerDisconnected {
1160 confirmed: false,
1161 consecutive_down: 1,
1162 ..
1163 }
1164 ),
1165 "first delta must be an unconfirmed down: {first:?}"
1166 );
1167 let second = next_event(&mut subscription).await?;
1168 assert!(
1169 matches!(
1170 &second,
1171 ClusterEvent::PeerDisconnected {
1172 confirmed: true,
1173 consecutive_down: 2,
1174 ..
1175 }
1176 ),
1177 "second delta must be the confirmed down: {second:?}"
1178 );
1179 let third = next_event(&mut subscription).await?;
1180 let ClusterEvent::ShardAdopted {
1181 shards,
1182 adopted_by,
1183 from_peer,
1184 ..
1185 } = &third
1186 else {
1187 return Err(format!("third delta must be ShardAdopted: {third:?}").into());
1188 };
1189 assert_eq!(shards, &vec![1]);
1190 assert_eq!(adopted_by, "node-self@127.0.0.1");
1191 assert_eq!(from_peer, "node-1@127.0.0.1");
1192
1193 liveness.set(true);
1196 drop(sup.tick().await);
1197 let recovery = next_event(&mut subscription).await?;
1198 assert!(
1199 matches!(&recovery, ClusterEvent::PeerConnected { .. }),
1200 "recovery delta must be PeerConnected: {recovery:?}"
1201 );
1202
1203 let quiet = sup.tick().await;
1207 assert!(quiet.is_empty());
1208 assert!(
1210 tokio::time::timeout(std::time::Duration::from_millis(50), subscription.next())
1211 .await
1212 .is_err(),
1213 "a steady connected peer must not re-emit PeerConnected every tick"
1214 );
1215 Ok(())
1216 }
1217
1218 async fn next_event(
1219 subscription: &mut futures::stream::BoxStream<
1220 'static,
1221 Result<aion_core::ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
1222 >,
1223 ) -> Result<aion_core::ClusterEvent, Box<dyn std::error::Error>> {
1224 use futures::StreamExt;
1225 tokio::time::timeout(std::time::Duration::from_secs(1), subscription.next())
1226 .await?
1227 .ok_or("cluster subscription ended")?
1228 .map_err(|lag| format!("unexpected lag: {lag:?}").into())
1229 }
1230
1231 #[tokio::test]
1234 async fn shard_published_to_the_dead_peer_itself_is_adopted() {
1235 let liveness = Arc::new(FakeLiveness::new(false));
1236 liveness.publish(1, "node-1@127.0.0.1", false);
1238 let adopter = Arc::new(FakeAdopter::new(false));
1239 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1240
1241 assert_eq!(
1242 sup.tick().await.len(),
1243 1,
1244 "a record naming the dead peer itself is stale and still adoptable"
1245 );
1246 assert_eq!(adopter.calls(), vec![vec![1]]);
1247 }
1248}