1use std::collections::BTreeMap;
36use std::sync::Arc;
37use std::time::Duration;
38
39use aion::Engine;
40use aion_core::ClusterEvent;
41
42use crate::cluster_publisher::ClusterEventPublisher;
43
44pub trait PeerLiveness: Send + Sync + 'static {
48 fn peer_connected(&self, peer_name: &str) -> bool;
50
51 fn read_shard_owner(&self, shard: usize) -> Option<String>;
62}
63
64#[cfg(feature = "haematite-backend")]
65impl PeerLiveness for aion_store_haematite::HaematiteStore {
66 fn peer_connected(&self, peer_name: &str) -> bool {
67 Self::peer_connected(self, peer_name)
68 }
69
70 fn read_shard_owner(&self, shard: usize) -> Option<String> {
71 Self::read_shard_owner(self, shard).ok().flatten()
73 }
74}
75
76#[async_trait::async_trait]
79pub trait ShardAdopter: Send + Sync + 'static {
80 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String>;
82}
83
84#[async_trait::async_trait]
85impl ShardAdopter for Engine {
86 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
87 Engine::adopt_shards(self, shards)
88 .await
89 .map_err(|error| error.to_string())
90 }
91}
92
93pub struct OutboxSettlingAdopter {
106 engine: Arc<Engine>,
107 outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
108}
109
110impl OutboxSettlingAdopter {
111 #[must_use]
115 pub fn new(
116 engine: Arc<Engine>,
117 outbox_store: Option<Arc<dyn aion_store::OutboxStore>>,
118 ) -> Self {
119 Self {
120 engine,
121 outbox_store,
122 }
123 }
124}
125
126#[async_trait::async_trait]
127impl ShardAdopter for OutboxSettlingAdopter {
128 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
129 ShardAdopter::adopt_shards(self.engine.as_ref(), shards).await?;
130 let Some(outbox_store) = &self.outbox_store else {
131 return Ok(());
132 };
133 match crate::worker::settle_terminal_outbox_rows(
134 self.engine.store().as_ref(),
135 outbox_store.as_ref(),
136 )
137 .await
138 {
139 Ok(settled) if settled.is_empty() => {}
140 Ok(settled) => {
141 tracing::info!(
142 ?shards,
143 settled = settled.len(),
144 "adoption sweep settled stranded outbox rows for terminal workflows"
145 );
146 }
147 Err(error) => {
148 tracing::error!(
149 ?shards,
150 %error,
151 "adoption sweep failed to settle terminal workflows' outbox rows; \
152 the reconciler liveness gate remains the backstop"
153 );
154 }
155 }
156 Ok(())
157 }
158}
159
160#[derive(Clone, Debug, PartialEq, Eq)]
163pub struct WatchedPeer {
164 pub name: String,
166 pub owned_shards: Vec<usize>,
168}
169
170#[derive(Clone, Copy, Debug)]
172pub struct SupervisorConfig {
173 pub poll_interval: Duration,
175 pub confirmations: u32,
178}
179
180#[derive(Default)]
182struct PeerState {
183 consecutive_down: u32,
185 adopted: bool,
187}
188
189pub struct ClusterSupervisor<L: PeerLiveness, A: ShardAdopter> {
191 liveness: Arc<L>,
192 adopter: Arc<A>,
193 peers: Vec<WatchedPeer>,
194 config: SupervisorConfig,
195 state: BTreeMap<String, PeerState>,
196 publisher: Option<Arc<ClusterEventPublisher>>,
201 self_node: String,
205}
206
207impl<L: PeerLiveness, A: ShardAdopter> ClusterSupervisor<L, A> {
208 #[must_use]
212 pub fn new(
213 liveness: Arc<L>,
214 adopter: Arc<A>,
215 peers: Vec<WatchedPeer>,
216 config: SupervisorConfig,
217 ) -> Self {
218 let peers: Vec<WatchedPeer> = peers
219 .into_iter()
220 .filter(|peer| !peer.owned_shards.is_empty())
221 .collect();
222 let state = peers
223 .iter()
224 .map(|peer| (peer.name.clone(), PeerState::default()))
225 .collect();
226 Self {
227 liveness,
228 adopter,
229 peers,
230 config,
231 state,
232 publisher: None,
233 self_node: String::new(),
234 }
235 }
236
237 #[must_use]
241 pub fn with_publisher(
242 mut self,
243 publisher: Arc<ClusterEventPublisher>,
244 self_node: impl Into<String>,
245 ) -> Self {
246 self.publisher = Some(publisher);
247 self.self_node = self_node.into();
248 self
249 }
250
251 fn emit<F>(&self, build: F)
255 where
256 F: FnOnce(aion_core::ClusterEventMeta) -> ClusterEvent,
257 {
258 if let Some(publisher) = &self.publisher {
259 drop(publisher.emit(build));
260 }
261 }
262
263 #[must_use]
266 pub fn watches_any(&self) -> bool {
267 !self.peers.is_empty()
268 }
269
270 #[must_use]
273 pub fn adopter(&self) -> &A {
274 &self.adopter
275 }
276
277 pub async fn tick(&mut self) -> Vec<String> {
285 let mut adopted_now = Vec::new();
286 let mut pending: Vec<ClusterEvent> = Vec::new();
290 let confirmations = self.config.confirmations;
291 for peer in &self.peers {
292 let connected = self.liveness.peer_connected(&peer.name);
293 let entry = self.state.entry(peer.name.clone()).or_default();
294 if connected {
295 let was_down = entry.consecutive_down > 0 || entry.adopted;
299 entry.consecutive_down = 0;
300 entry.adopted = false;
301 if was_down {
302 pending.push(ClusterEvent::PeerConnected {
303 meta: placeholder_meta(),
304 peer_name: peer.name.clone(),
305 forward_addr: None,
306 });
307 }
308 continue;
309 }
310 entry.consecutive_down = entry.consecutive_down.saturating_add(1);
311 let consecutive_down = entry.consecutive_down;
312 let confirmed = consecutive_down >= confirmations;
313 pending.push(ClusterEvent::PeerDisconnected {
316 meta: placeholder_meta(),
317 peer_name: peer.name.clone(),
318 consecutive_down,
319 confirmed,
320 });
321 if entry.adopted || consecutive_down < confirmations {
322 continue;
323 }
324 if Self::all_shards_handled_elsewhere(
331 self.liveness.as_ref(),
332 &peer.name,
333 &peer.owned_shards,
334 ) {
335 entry.adopted = true;
338 let held_by = Self::live_owner_of(self.liveness.as_ref(), &peer.owned_shards)
339 .unwrap_or_default();
340 pending.push(ClusterEvent::ShardAdoptionSkipped {
341 meta: placeholder_meta(),
342 shards: peer.owned_shards.clone(),
343 from_peer: peer.name.clone(),
344 held_by,
345 });
346 tracing::info!(
347 peer = %peer.name,
348 shards = ?peer.owned_shards,
349 "downed peer's shards already adopted by another live owner; skipping"
350 );
351 continue;
352 }
353 match self.adopter.adopt_shards(&peer.owned_shards).await {
354 Ok(()) => {
355 entry.adopted = true;
356 adopted_now.push(peer.name.clone());
357 pending.push(ClusterEvent::ShardAdopted {
358 meta: placeholder_meta(),
359 shards: peer.owned_shards.clone(),
360 from_peer: peer.name.clone(),
361 adopted_by: self.self_node.clone(),
362 });
363 tracing::info!(
364 peer = %peer.name,
365 shards = ?peer.owned_shards,
366 "cluster supervisor adopted a downed peer's shards (SS-5b auto-failover)"
367 );
368 }
369 Err(error) => {
370 pending.push(ClusterEvent::ShardAdoptionFailed {
371 meta: placeholder_meta(),
372 shards: peer.owned_shards.clone(),
373 from_peer: peer.name.clone(),
374 error: error.clone(),
375 });
376 tracing::warn!(
385 peer = %peer.name,
386 shards = ?peer.owned_shards,
387 %error,
388 "cluster supervisor failed to adopt a downed peer's shards; will retry"
389 );
390 }
391 }
392 }
393 for event in pending {
397 self.emit(|meta| with_meta(event, meta));
398 }
399 adopted_now
400 }
401
402 fn live_owner_of(liveness: &L, shards: &[usize]) -> Option<String> {
406 shards.iter().find_map(|&shard| {
407 liveness
408 .read_shard_owner(shard)
409 .filter(|owner| liveness.peer_connected(owner))
410 })
411 }
412
413 fn all_shards_handled_elsewhere(liveness: &L, peer_name: &str, shards: &[usize]) -> bool {
421 !shards.is_empty()
422 && shards.iter().all(|&shard| {
423 liveness.read_shard_owner(shard).is_some_and(|owner| {
424 owner != peer_name && liveness.peer_connected(&owner)
428 })
429 })
430 }
431
432 pub async fn run(mut self, mut shutdown: tokio::sync::watch::Receiver<bool>) {
435 let self_node = self.self_node.clone();
438 self.emit(|meta| ClusterEvent::SupervisorStarted {
439 meta,
440 node: self_node.clone(),
441 });
442 let mut interval = tokio::time::interval(self.config.poll_interval);
443 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
444 loop {
445 tokio::select! {
446 _ = interval.tick() => {
447 drop(self.tick().await);
448 }
449 changed = shutdown.changed() => {
450 if changed.is_err() || *shutdown.borrow() {
451 break;
452 }
453 }
454 }
455 }
456 let self_node = self.self_node.clone();
459 self.emit(|meta| ClusterEvent::SupervisorStopped {
460 meta,
461 node: self_node.clone(),
462 });
463 }
464}
465
466fn placeholder_meta() -> aion_core::ClusterEventMeta {
470 aion_core::ClusterEventMeta {
471 cluster_seq: 0,
472 observed_at: chrono::Utc::now(),
473 }
474}
475
476fn with_meta(event: ClusterEvent, meta: aion_core::ClusterEventMeta) -> ClusterEvent {
483 match event {
484 ClusterEvent::PeerAdded {
485 peer_name,
486 forward_addr,
487 ..
488 } => ClusterEvent::PeerAdded {
489 meta,
490 peer_name,
491 forward_addr,
492 },
493 ClusterEvent::PeerConnected {
494 peer_name,
495 forward_addr,
496 ..
497 } => ClusterEvent::PeerConnected {
498 meta,
499 peer_name,
500 forward_addr,
501 },
502 ClusterEvent::PeerDisconnected {
503 peer_name,
504 consecutive_down,
505 confirmed,
506 ..
507 } => ClusterEvent::PeerDisconnected {
508 meta,
509 peer_name,
510 consecutive_down,
511 confirmed,
512 },
513 ClusterEvent::ShardAdopted {
514 shards,
515 from_peer,
516 adopted_by,
517 ..
518 } => ClusterEvent::ShardAdopted {
519 meta,
520 shards,
521 from_peer,
522 adopted_by,
523 },
524 ClusterEvent::ShardAdoptionFailed {
525 shards,
526 from_peer,
527 error,
528 ..
529 } => ClusterEvent::ShardAdoptionFailed {
530 meta,
531 shards,
532 from_peer,
533 error,
534 },
535 ClusterEvent::ShardAdoptionSkipped {
536 shards,
537 from_peer,
538 held_by,
539 ..
540 } => ClusterEvent::ShardAdoptionSkipped {
541 meta,
542 shards,
543 from_peer,
544 held_by,
545 },
546 other => with_meta_worker_lifecycle(other, meta),
547 }
548}
549
550fn with_meta_worker_lifecycle(
560 event: ClusterEvent,
561 meta: aion_core::ClusterEventMeta,
562) -> ClusterEvent {
563 match event {
564 ClusterEvent::WorkerConnected {
565 worker_id,
566 namespaces,
567 task_queue,
568 transport,
569 node,
570 ..
571 } => ClusterEvent::WorkerConnected {
572 meta,
573 worker_id,
574 namespaces,
575 task_queue,
576 transport,
577 node,
578 },
579 ClusterEvent::WorkerDisconnected {
580 worker_id,
581 namespaces,
582 reason,
583 ..
584 } => ClusterEvent::WorkerDisconnected {
585 meta,
586 worker_id,
587 namespaces,
588 reason,
589 },
590 ClusterEvent::SupervisorStarted { node, .. } => {
591 ClusterEvent::SupervisorStarted { meta, node }
592 }
593 ClusterEvent::SupervisorStopped { node, .. } => {
594 ClusterEvent::SupervisorStopped { meta, node }
595 }
596 ClusterEvent::NamespaceCreated {
597 name,
598 created_at,
599 origin,
600 ..
601 } => ClusterEvent::NamespaceCreated {
602 meta,
603 name,
604 created_at,
605 origin,
606 },
607 ClusterEvent::NamespacePlacementChanged {
611 name, placement, ..
612 } => ClusterEvent::NamespacePlacementChanged {
613 meta,
614 name,
615 placement,
616 },
617 ClusterEvent::NamespaceQuotaState {
621 namespace,
622 in_flight,
623 ceiling,
624 ..
625 } => ClusterEvent::NamespaceQuotaState {
626 meta,
627 namespace,
628 in_flight,
629 ceiling,
630 },
631 ClusterEvent::PeerAdded { .. }
632 | ClusterEvent::PeerConnected { .. }
633 | ClusterEvent::PeerDisconnected { .. }
634 | ClusterEvent::ShardAdopted { .. }
635 | ClusterEvent::ShardAdoptionFailed { .. }
636 | ClusterEvent::ShardAdoptionSkipped { .. } => {
637 unreachable!("peer/shard variants are re-stamped by with_meta, never delegated here")
638 }
639 }
640}
641
642#[cfg(test)]
643mod tests {
644 use std::sync::Mutex;
645 use std::sync::atomic::{AtomicBool, Ordering};
646
647 use super::*;
648
649 struct FakeLiveness {
654 connected: AtomicBool,
655 owners: Mutex<std::collections::BTreeMap<usize, String>>,
657 live_owners: Mutex<std::collections::BTreeSet<String>>,
659 }
660
661 impl FakeLiveness {
662 fn new(connected: bool) -> Self {
663 Self {
664 connected: AtomicBool::new(connected),
665 owners: Mutex::new(std::collections::BTreeMap::new()),
666 live_owners: Mutex::new(std::collections::BTreeSet::new()),
667 }
668 }
669 fn set(&self, connected: bool) {
670 self.connected.store(connected, Ordering::SeqCst);
671 }
672 fn publish(&self, shard: usize, owner: &str, live: bool) {
675 self.owners
676 .lock()
677 .unwrap_or_else(std::sync::PoisonError::into_inner)
678 .insert(shard, owner.to_owned());
679 if live {
680 self.live_owners
681 .lock()
682 .unwrap_or_else(std::sync::PoisonError::into_inner)
683 .insert(owner.to_owned());
684 }
685 }
686 }
687
688 impl PeerLiveness for FakeLiveness {
689 fn peer_connected(&self, peer_name: &str) -> bool {
690 if self
691 .live_owners
692 .lock()
693 .unwrap_or_else(std::sync::PoisonError::into_inner)
694 .contains(peer_name)
695 {
696 return true;
697 }
698 self.connected.load(Ordering::SeqCst)
699 }
700
701 fn read_shard_owner(&self, shard: usize) -> Option<String> {
702 self.owners
703 .lock()
704 .unwrap_or_else(std::sync::PoisonError::into_inner)
705 .get(&shard)
706 .cloned()
707 }
708 }
709
710 struct FakeAdopter {
712 calls: Mutex<Vec<Vec<usize>>>,
713 fail_first: AtomicBool,
714 }
715
716 impl FakeAdopter {
717 fn new(fail_first: bool) -> Self {
718 Self {
719 calls: Mutex::new(Vec::new()),
720 fail_first: AtomicBool::new(fail_first),
721 }
722 }
723 fn calls(&self) -> Vec<Vec<usize>> {
724 self.calls
725 .lock()
726 .unwrap_or_else(std::sync::PoisonError::into_inner)
727 .clone()
728 }
729 }
730
731 #[async_trait::async_trait]
732 impl ShardAdopter for FakeAdopter {
733 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
734 if self.fail_first.swap(false, Ordering::SeqCst) {
735 return Err("simulated election failure".to_owned());
736 }
737 self.calls
738 .lock()
739 .unwrap_or_else(std::sync::PoisonError::into_inner)
740 .push(shards.to_vec());
741 Ok(())
742 }
743 }
744
745 fn supervisor(
746 liveness: Arc<FakeLiveness>,
747 adopter: Arc<FakeAdopter>,
748 confirmations: u32,
749 ) -> ClusterSupervisor<FakeLiveness, FakeAdopter> {
750 ClusterSupervisor::new(
751 liveness,
752 adopter,
753 vec![WatchedPeer {
754 name: "node-1@127.0.0.1".to_owned(),
755 owned_shards: vec![1],
756 }],
757 SupervisorConfig {
758 poll_interval: Duration::from_millis(1),
759 confirmations,
760 },
761 )
762 }
763
764 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
770 async fn outbox_settling_adopter_settles_terminal_rows_after_adoption()
771 -> Result<(), Box<dyn std::error::Error>> {
772 use aion::{EngineBuilder, RuntimeHandle, SignalRouter};
773 use aion_core::{Event, EventEnvelope};
774 use aion_store::{OutboxRow, OutboxStatus, OutboxStore, WritableEventStore, WriteToken};
775 use aion_store_libsql::LibSqlStore;
776
777 let db_path = std::env::temp_dir().join(format!(
778 "aion-adopter-settle-{}-{}.db",
779 std::process::id(),
780 uuid::Uuid::new_v4()
781 ));
782
783 let seeder = LibSqlStore::open(db_path.clone()).await?;
786 let workflow_id = aion_core::WorkflowId::new_v4();
787 let envelope = |seq: u64| EventEnvelope {
788 seq,
789 recorded_at: chrono::Utc::now(),
790 workflow_id: workflow_id.clone(),
791 };
792 let events = vec![
793 Event::WorkflowStarted {
794 envelope: envelope(1),
795 workflow_type: String::from("dev_brief"),
796 input: aion_core::Payload::from_json(&serde_json::json!({}))?,
797 run_id: aion_core::RunId::new_v4(),
798 parent_run_id: None,
799 package_version: aion_core::PackageVersion::new("a".repeat(64)),
800 },
801 Event::WorkflowFailed {
802 envelope: envelope(2),
803 error: aion_core::WorkflowError {
804 message: String::from("boom"),
805 details: None,
806 },
807 },
808 ];
809 seeder
810 .append(WriteToken::recorder(), &workflow_id, &events, 0)
811 .await?;
812 let row = OutboxRow::pending(
813 workflow_id.clone(),
814 0,
815 String::from("norn_round"),
816 aion_core::Payload::from_json(&serde_json::json!({}))?,
817 chrono::Utc::now(),
818 );
819 let dispatch_key = row.dispatch_key.clone();
820 seeder
821 .append_outbox_batch(std::slice::from_ref(&row))
822 .await?;
823 assert_eq!(seeder.claim_outbox_rows(1).await?.len(), 1);
824
825 let engine = Arc::new(
826 EngineBuilder::new()
827 .store_arc(Arc::new(LibSqlStore::open(db_path.clone()).await?))
828 .in_memory_visibility()
829 .scheduler_threads(1)
830 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
831 Arc::new(aion::signal::ConcreteSignalRouter::new(runtime, handoff))
832 as Arc<dyn SignalRouter>
833 })
834 .build()
835 .await?,
836 );
837 let outbox_store: Arc<dyn OutboxStore> =
838 Arc::new(LibSqlStore::open(db_path.clone()).await?);
839 let adopter =
840 OutboxSettlingAdopter::new(Arc::clone(&engine), Some(Arc::clone(&outbox_store)));
841
842 ShardAdopter::adopt_shards(&adopter, &[42])
843 .await
844 .map_err(|error| format!("adoption must succeed: {error}"))?;
845
846 let state = seeder
847 .outbox_row_state(&dispatch_key)
848 .await?
849 .ok_or("the stranded row must still exist")?;
850 assert_eq!(
851 state.status,
852 OutboxStatus::Cancelled,
853 "the adoption sweep must settle the terminal workflow's stranded row"
854 );
855 engine.shutdown()?;
856 Ok(())
857 }
858
859 #[tokio::test]
860 async fn does_not_adopt_while_peer_connected() {
861 let liveness = Arc::new(FakeLiveness::new(true));
862 let adopter = Arc::new(FakeAdopter::new(false));
863 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2);
864 for _ in 0..5 {
865 assert!(sup.tick().await.is_empty());
866 }
867 assert!(adopter.calls().is_empty(), "no adoption while peer is up");
868 }
869
870 #[tokio::test]
871 async fn debounce_requires_consecutive_down_before_adopting() {
872 let liveness = Arc::new(FakeLiveness::new(true));
873 let adopter = Arc::new(FakeAdopter::new(false));
874 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 3);
875
876 liveness.set(false);
877 assert!(sup.tick().await.is_empty(), "tick 1 down: below threshold");
878 liveness.set(true);
880 assert!(sup.tick().await.is_empty());
881 liveness.set(false);
882 assert!(
883 sup.tick().await.is_empty(),
884 "down again, counter reset to 1"
885 );
886 assert!(sup.tick().await.is_empty(), "2 consecutive: still below 3");
887 let fired = sup.tick().await;
888 assert_eq!(
889 fired,
890 vec!["node-1@127.0.0.1".to_owned()],
891 "3rd consecutive triggers"
892 );
893 assert_eq!(adopter.calls(), vec![vec![1]]);
894 }
895
896 #[tokio::test]
897 async fn adopts_once_then_stays_quiet_while_down() {
898 let liveness = Arc::new(FakeLiveness::new(false));
899 let adopter = Arc::new(FakeAdopter::new(false));
900 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
901 assert_eq!(sup.tick().await.len(), 1, "first down tick adopts");
902 for _ in 0..5 {
903 assert!(sup.tick().await.is_empty(), "no re-adopt while still down");
904 }
905 assert_eq!(adopter.calls(), vec![vec![1]], "adopted exactly once");
906 }
907
908 #[tokio::test]
909 async fn failed_adoption_is_retried_next_tick() {
910 let liveness = Arc::new(FakeLiveness::new(false));
911 let adopter = Arc::new(FakeAdopter::new(true)); let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
913 assert!(
914 sup.tick().await.is_empty(),
915 "first adopt fails, not recorded"
916 );
917 assert!(adopter.calls().is_empty());
918 assert_eq!(sup.tick().await.len(), 1, "retry succeeds next tick");
919 assert_eq!(adopter.calls(), vec![vec![1]]);
920 }
921
922 #[tokio::test]
923 async fn peer_with_no_shards_is_not_watched() {
924 let liveness = Arc::new(FakeLiveness::new(false));
925 let adopter = Arc::new(FakeAdopter::new(false));
926 let mut sup = ClusterSupervisor::new(
927 Arc::clone(&liveness),
928 Arc::clone(&adopter),
929 vec![WatchedPeer {
930 name: "node-2@127.0.0.1".to_owned(),
931 owned_shards: vec![],
932 }],
933 SupervisorConfig {
934 poll_interval: Duration::from_millis(1),
935 confirmations: 1,
936 },
937 );
938 assert!(!sup.watches_any());
939 assert!(sup.tick().await.is_empty());
940 assert!(adopter.calls().is_empty());
941 }
942
943 #[tokio::test]
947 async fn shard_already_published_to_live_owner_is_not_adopted() {
948 let liveness = Arc::new(FakeLiveness::new(false));
949 liveness.publish(1, "node-9@127.0.0.1", true);
952 let adopter = Arc::new(FakeAdopter::new(false));
953 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
954
955 assert!(
957 sup.tick().await.is_empty(),
958 "no adoption fires for a shard a live owner already holds"
959 );
960 assert!(
961 adopter.calls().is_empty(),
962 "the adopter is never invoked for an already-handled shard"
963 );
964 for _ in 0..3 {
966 assert!(sup.tick().await.is_empty());
967 }
968 assert!(adopter.calls().is_empty());
969 }
970
971 #[tokio::test]
975 async fn shard_published_to_a_down_owner_is_still_adopted() {
976 let liveness = Arc::new(FakeLiveness::new(false));
977 liveness.publish(1, "node-9@127.0.0.1", false);
980 let adopter = Arc::new(FakeAdopter::new(false));
981 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
982
983 assert_eq!(
984 sup.tick().await.len(),
985 1,
986 "a shard whose recorded owner is itself down is adoptable"
987 );
988 assert_eq!(adopter.calls(), vec![vec![1]]);
989 }
990
991 #[tokio::test]
997 async fn tick_emits_topology_deltas_through_the_publisher()
998 -> Result<(), Box<dyn std::error::Error>> {
999 use std::num::NonZeroUsize;
1000
1001 use aion_core::ClusterEvent;
1002 use futures::StreamExt;
1003
1004 use crate::cluster_publisher::ClusterEventPublisher;
1005
1006 let capacity = NonZeroUsize::new(64).ok_or("non-zero")?;
1007 let publisher = Arc::new(ClusterEventPublisher::new(capacity));
1008 let mut subscription = publisher.subscribe(0);
1009
1010 let liveness = Arc::new(FakeLiveness::new(true));
1011 let adopter = Arc::new(FakeAdopter::new(false));
1012 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2)
1013 .with_publisher(Arc::clone(&publisher), "node-self@127.0.0.1");
1014
1015 liveness.set(false);
1017 drop(sup.tick().await);
1018 let fired = sup.tick().await;
1020 assert_eq!(fired, vec!["node-1@127.0.0.1".to_owned()]);
1021
1022 let first = next_event(&mut subscription).await?;
1024 assert!(
1025 matches!(
1026 &first,
1027 ClusterEvent::PeerDisconnected {
1028 confirmed: false,
1029 consecutive_down: 1,
1030 ..
1031 }
1032 ),
1033 "first delta must be an unconfirmed down: {first:?}"
1034 );
1035 let second = next_event(&mut subscription).await?;
1036 assert!(
1037 matches!(
1038 &second,
1039 ClusterEvent::PeerDisconnected {
1040 confirmed: true,
1041 consecutive_down: 2,
1042 ..
1043 }
1044 ),
1045 "second delta must be the confirmed down: {second:?}"
1046 );
1047 let third = next_event(&mut subscription).await?;
1048 let ClusterEvent::ShardAdopted {
1049 shards,
1050 adopted_by,
1051 from_peer,
1052 ..
1053 } = &third
1054 else {
1055 return Err(format!("third delta must be ShardAdopted: {third:?}").into());
1056 };
1057 assert_eq!(shards, &vec![1]);
1058 assert_eq!(adopted_by, "node-self@127.0.0.1");
1059 assert_eq!(from_peer, "node-1@127.0.0.1");
1060
1061 liveness.set(true);
1064 drop(sup.tick().await);
1065 let recovery = next_event(&mut subscription).await?;
1066 assert!(
1067 matches!(&recovery, ClusterEvent::PeerConnected { .. }),
1068 "recovery delta must be PeerConnected: {recovery:?}"
1069 );
1070
1071 let quiet = sup.tick().await;
1075 assert!(quiet.is_empty());
1076 assert!(
1078 tokio::time::timeout(std::time::Duration::from_millis(50), subscription.next())
1079 .await
1080 .is_err(),
1081 "a steady connected peer must not re-emit PeerConnected every tick"
1082 );
1083 Ok(())
1084 }
1085
1086 async fn next_event(
1087 subscription: &mut futures::stream::BoxStream<
1088 'static,
1089 Result<aion_core::ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
1090 >,
1091 ) -> Result<aion_core::ClusterEvent, Box<dyn std::error::Error>> {
1092 use futures::StreamExt;
1093 tokio::time::timeout(std::time::Duration::from_secs(1), subscription.next())
1094 .await?
1095 .ok_or("cluster subscription ended")?
1096 .map_err(|lag| format!("unexpected lag: {lag:?}").into())
1097 }
1098
1099 #[tokio::test]
1102 async fn shard_published_to_the_dead_peer_itself_is_adopted() {
1103 let liveness = Arc::new(FakeLiveness::new(false));
1104 liveness.publish(1, "node-1@127.0.0.1", false);
1106 let adopter = Arc::new(FakeAdopter::new(false));
1107 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1108
1109 assert_eq!(
1110 sup.tick().await.len(),
1111 1,
1112 "a record naming the dead peer itself is stale and still adoptable"
1113 );
1114 assert_eq!(adopter.calls(), vec![vec![1]]);
1115 }
1116}