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 deployment,
571 deployment_association,
572 ..
573 } => ClusterEvent::WorkerConnected {
574 meta,
575 worker_id,
576 namespaces,
577 task_queue,
578 transport,
579 node,
580 deployment,
581 deployment_association,
582 },
583 ClusterEvent::WorkerDisconnected {
584 worker_id,
585 namespaces,
586 reason,
587 ..
588 } => ClusterEvent::WorkerDisconnected {
589 meta,
590 worker_id,
591 namespaces,
592 reason,
593 },
594 ClusterEvent::SupervisorStarted { node, .. } => {
595 ClusterEvent::SupervisorStarted { meta, node }
596 }
597 ClusterEvent::SupervisorStopped { node, .. } => {
598 ClusterEvent::SupervisorStopped { meta, node }
599 }
600 ClusterEvent::NamespaceCreated {
601 name,
602 created_at,
603 origin,
604 ..
605 } => ClusterEvent::NamespaceCreated {
606 meta,
607 name,
608 created_at,
609 origin,
610 },
611 ClusterEvent::NamespacePlacementChanged {
615 name, placement, ..
616 } => ClusterEvent::NamespacePlacementChanged {
617 meta,
618 name,
619 placement,
620 },
621 ClusterEvent::NamespaceQuotaState {
625 namespace,
626 in_flight,
627 ceiling,
628 ..
629 } => ClusterEvent::NamespaceQuotaState {
630 meta,
631 namespace,
632 in_flight,
633 ceiling,
634 },
635 deployment_event @ (ClusterEvent::WorkerDeploymentPut { .. }
636 | ClusterEvent::WorkerDeploymentDesiredStateChanged { .. }
637 | ClusterEvent::WorkerDeploymentDeleted { .. }) => {
638 with_meta_worker_deployment(deployment_event, meta)
639 }
640 ClusterEvent::PeerAdded { .. }
641 | ClusterEvent::PeerConnected { .. }
642 | ClusterEvent::PeerDisconnected { .. }
643 | ClusterEvent::ShardAdopted { .. }
644 | ClusterEvent::ShardAdoptionFailed { .. }
645 | ClusterEvent::ShardAdoptionSkipped { .. } => {
646 unreachable!("peer/shard variants are re-stamped by with_meta, never delegated here")
647 }
648 }
649}
650
651fn with_meta_worker_deployment(
652 event: ClusterEvent,
653 meta: aion_core::ClusterEventMeta,
654) -> ClusterEvent {
655 match event {
656 ClusterEvent::WorkerDeploymentPut {
657 name,
658 outcome,
659 desired_state,
660 binary_version,
661 binary_content_hash,
662 ..
663 } => ClusterEvent::WorkerDeploymentPut {
664 meta,
665 name,
666 outcome,
667 desired_state,
668 binary_version,
669 binary_content_hash,
670 },
671 ClusterEvent::WorkerDeploymentDesiredStateChanged {
672 name,
673 desired_state,
674 ..
675 } => ClusterEvent::WorkerDeploymentDesiredStateChanged {
676 meta,
677 name,
678 desired_state,
679 },
680 ClusterEvent::WorkerDeploymentDeleted { name, .. } => {
681 ClusterEvent::WorkerDeploymentDeleted { meta, name }
682 }
683 ClusterEvent::WorkerConnected { .. }
684 | ClusterEvent::WorkerDisconnected { .. }
685 | ClusterEvent::SupervisorStarted { .. }
686 | ClusterEvent::SupervisorStopped { .. }
687 | ClusterEvent::NamespaceCreated { .. }
688 | ClusterEvent::NamespacePlacementChanged { .. }
689 | ClusterEvent::NamespaceQuotaState { .. }
690 | ClusterEvent::PeerAdded { .. }
691 | ClusterEvent::PeerConnected { .. }
692 | ClusterEvent::PeerDisconnected { .. }
693 | ClusterEvent::ShardAdopted { .. }
694 | ClusterEvent::ShardAdoptionFailed { .. }
695 | ClusterEvent::ShardAdoptionSkipped { .. } => {
696 unreachable!("only worker-deployment variants are delegated here")
697 }
698 }
699}
700
701#[cfg(test)]
702mod tests {
703 use std::sync::Mutex;
704 use std::sync::atomic::{AtomicBool, Ordering};
705
706 use super::*;
707
708 struct FakeLiveness {
713 connected: AtomicBool,
714 owners: Mutex<std::collections::BTreeMap<usize, String>>,
716 live_owners: Mutex<std::collections::BTreeSet<String>>,
718 }
719
720 impl FakeLiveness {
721 fn new(connected: bool) -> Self {
722 Self {
723 connected: AtomicBool::new(connected),
724 owners: Mutex::new(std::collections::BTreeMap::new()),
725 live_owners: Mutex::new(std::collections::BTreeSet::new()),
726 }
727 }
728 fn set(&self, connected: bool) {
729 self.connected.store(connected, Ordering::SeqCst);
730 }
731 fn publish(&self, shard: usize, owner: &str, live: bool) {
734 self.owners
735 .lock()
736 .unwrap_or_else(std::sync::PoisonError::into_inner)
737 .insert(shard, owner.to_owned());
738 if live {
739 self.live_owners
740 .lock()
741 .unwrap_or_else(std::sync::PoisonError::into_inner)
742 .insert(owner.to_owned());
743 }
744 }
745 }
746
747 impl PeerLiveness for FakeLiveness {
748 fn peer_connected(&self, peer_name: &str) -> bool {
749 if self
750 .live_owners
751 .lock()
752 .unwrap_or_else(std::sync::PoisonError::into_inner)
753 .contains(peer_name)
754 {
755 return true;
756 }
757 self.connected.load(Ordering::SeqCst)
758 }
759
760 fn read_shard_owner(&self, shard: usize) -> Option<String> {
761 self.owners
762 .lock()
763 .unwrap_or_else(std::sync::PoisonError::into_inner)
764 .get(&shard)
765 .cloned()
766 }
767 }
768
769 struct FakeAdopter {
771 calls: Mutex<Vec<Vec<usize>>>,
772 fail_first: AtomicBool,
773 }
774
775 impl FakeAdopter {
776 fn new(fail_first: bool) -> Self {
777 Self {
778 calls: Mutex::new(Vec::new()),
779 fail_first: AtomicBool::new(fail_first),
780 }
781 }
782 fn calls(&self) -> Vec<Vec<usize>> {
783 self.calls
784 .lock()
785 .unwrap_or_else(std::sync::PoisonError::into_inner)
786 .clone()
787 }
788 }
789
790 #[async_trait::async_trait]
791 impl ShardAdopter for FakeAdopter {
792 async fn adopt_shards(&self, shards: &[usize]) -> Result<(), String> {
793 if self.fail_first.swap(false, Ordering::SeqCst) {
794 return Err("simulated election failure".to_owned());
795 }
796 self.calls
797 .lock()
798 .unwrap_or_else(std::sync::PoisonError::into_inner)
799 .push(shards.to_vec());
800 Ok(())
801 }
802 }
803
804 fn supervisor(
805 liveness: Arc<FakeLiveness>,
806 adopter: Arc<FakeAdopter>,
807 confirmations: u32,
808 ) -> ClusterSupervisor<FakeLiveness, FakeAdopter> {
809 ClusterSupervisor::new(
810 liveness,
811 adopter,
812 vec![WatchedPeer {
813 name: "node-1@127.0.0.1".to_owned(),
814 owned_shards: vec![1],
815 }],
816 SupervisorConfig {
817 poll_interval: Duration::from_millis(1),
818 confirmations,
819 },
820 )
821 }
822
823 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
829 async fn outbox_settling_adopter_settles_terminal_rows_after_adoption()
830 -> Result<(), Box<dyn std::error::Error>> {
831 use aion::{EngineBuilder, RuntimeHandle, SignalRouter};
832 use aion_core::{Event, EventEnvelope};
833 use aion_store::{OutboxRow, OutboxStatus, OutboxStore, WritableEventStore, WriteToken};
834 use aion_store_libsql::LibSqlStore;
835
836 let db_path = std::env::temp_dir().join(format!(
837 "aion-adopter-settle-{}-{}.db",
838 std::process::id(),
839 uuid::Uuid::new_v4()
840 ));
841
842 let seeder = LibSqlStore::open(db_path.clone()).await?;
845 let workflow_id = aion_core::WorkflowId::new_v4();
846 let envelope = |seq: u64| EventEnvelope {
847 seq,
848 recorded_at: chrono::Utc::now(),
849 workflow_id: workflow_id.clone(),
850 };
851 let events = vec![
852 Event::WorkflowStarted {
853 envelope: envelope(1),
854 workflow_type: String::from("dev_brief"),
855 input: aion_core::Payload::from_json(&serde_json::json!({}))?,
856 run_id: aion_core::RunId::new_v4(),
857 parent_run_id: None,
858 package_version: aion_core::PackageVersion::new("a".repeat(64)),
859 },
860 Event::WorkflowFailed {
861 envelope: envelope(2),
862 error: aion_core::WorkflowError {
863 message: String::from("boom"),
864 details: None,
865 },
866 },
867 ];
868 seeder
869 .append(WriteToken::recorder(), &workflow_id, &events, 0)
870 .await?;
871 let row = OutboxRow::pending(
872 workflow_id.clone(),
873 0,
874 String::from("norn_round"),
875 aion_core::Payload::from_json(&serde_json::json!({}))?,
876 chrono::Utc::now(),
877 );
878 let dispatch_key = row.dispatch_key.clone();
879 seeder
880 .append_outbox_batch(std::slice::from_ref(&row))
881 .await?;
882 assert_eq!(seeder.claim_outbox_rows(1).await?.len(), 1);
883
884 let engine = Arc::new(
885 EngineBuilder::new()
886 .store_arc(Arc::new(LibSqlStore::open(db_path.clone()).await?))
887 .in_memory_visibility()
888 .scheduler_threads(1)
889 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
890 Arc::new(aion::signal::ConcreteSignalRouter::new(runtime, handoff))
891 as Arc<dyn SignalRouter>
892 })
893 .build()
894 .await?,
895 );
896 let outbox_store: Arc<dyn OutboxStore> =
897 Arc::new(LibSqlStore::open(db_path.clone()).await?);
898 let adopter =
899 OutboxSettlingAdopter::new(Arc::clone(&engine), Some(Arc::clone(&outbox_store)));
900
901 ShardAdopter::adopt_shards(&adopter, &[42])
902 .await
903 .map_err(|error| format!("adoption must succeed: {error}"))?;
904
905 let state = seeder
906 .outbox_row_state(&dispatch_key)
907 .await?
908 .ok_or("the stranded row must still exist")?;
909 assert_eq!(
910 state.status,
911 OutboxStatus::Cancelled,
912 "the adoption sweep must settle the terminal workflow's stranded row"
913 );
914 engine.shutdown()?;
915 Ok(())
916 }
917
918 #[tokio::test]
919 async fn does_not_adopt_while_peer_connected() {
920 let liveness = Arc::new(FakeLiveness::new(true));
921 let adopter = Arc::new(FakeAdopter::new(false));
922 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2);
923 for _ in 0..5 {
924 assert!(sup.tick().await.is_empty());
925 }
926 assert!(adopter.calls().is_empty(), "no adoption while peer is up");
927 }
928
929 #[tokio::test]
930 async fn debounce_requires_consecutive_down_before_adopting() {
931 let liveness = Arc::new(FakeLiveness::new(true));
932 let adopter = Arc::new(FakeAdopter::new(false));
933 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 3);
934
935 liveness.set(false);
936 assert!(sup.tick().await.is_empty(), "tick 1 down: below threshold");
937 liveness.set(true);
939 assert!(sup.tick().await.is_empty());
940 liveness.set(false);
941 assert!(
942 sup.tick().await.is_empty(),
943 "down again, counter reset to 1"
944 );
945 assert!(sup.tick().await.is_empty(), "2 consecutive: still below 3");
946 let fired = sup.tick().await;
947 assert_eq!(
948 fired,
949 vec!["node-1@127.0.0.1".to_owned()],
950 "3rd consecutive triggers"
951 );
952 assert_eq!(adopter.calls(), vec![vec![1]]);
953 }
954
955 #[tokio::test]
956 async fn adopts_once_then_stays_quiet_while_down() {
957 let liveness = Arc::new(FakeLiveness::new(false));
958 let adopter = Arc::new(FakeAdopter::new(false));
959 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
960 assert_eq!(sup.tick().await.len(), 1, "first down tick adopts");
961 for _ in 0..5 {
962 assert!(sup.tick().await.is_empty(), "no re-adopt while still down");
963 }
964 assert_eq!(adopter.calls(), vec![vec![1]], "adopted exactly once");
965 }
966
967 #[tokio::test]
968 async fn failed_adoption_is_retried_next_tick() {
969 let liveness = Arc::new(FakeLiveness::new(false));
970 let adopter = Arc::new(FakeAdopter::new(true)); let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
972 assert!(
973 sup.tick().await.is_empty(),
974 "first adopt fails, not recorded"
975 );
976 assert!(adopter.calls().is_empty());
977 assert_eq!(sup.tick().await.len(), 1, "retry succeeds next tick");
978 assert_eq!(adopter.calls(), vec![vec![1]]);
979 }
980
981 #[tokio::test]
982 async fn peer_with_no_shards_is_not_watched() {
983 let liveness = Arc::new(FakeLiveness::new(false));
984 let adopter = Arc::new(FakeAdopter::new(false));
985 let mut sup = ClusterSupervisor::new(
986 Arc::clone(&liveness),
987 Arc::clone(&adopter),
988 vec![WatchedPeer {
989 name: "node-2@127.0.0.1".to_owned(),
990 owned_shards: vec![],
991 }],
992 SupervisorConfig {
993 poll_interval: Duration::from_millis(1),
994 confirmations: 1,
995 },
996 );
997 assert!(!sup.watches_any());
998 assert!(sup.tick().await.is_empty());
999 assert!(adopter.calls().is_empty());
1000 }
1001
1002 #[tokio::test]
1006 async fn shard_already_published_to_live_owner_is_not_adopted() {
1007 let liveness = Arc::new(FakeLiveness::new(false));
1008 liveness.publish(1, "node-9@127.0.0.1", true);
1011 let adopter = Arc::new(FakeAdopter::new(false));
1012 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1013
1014 assert!(
1016 sup.tick().await.is_empty(),
1017 "no adoption fires for a shard a live owner already holds"
1018 );
1019 assert!(
1020 adopter.calls().is_empty(),
1021 "the adopter is never invoked for an already-handled shard"
1022 );
1023 for _ in 0..3 {
1025 assert!(sup.tick().await.is_empty());
1026 }
1027 assert!(adopter.calls().is_empty());
1028 }
1029
1030 #[tokio::test]
1034 async fn shard_published_to_a_down_owner_is_still_adopted() {
1035 let liveness = Arc::new(FakeLiveness::new(false));
1036 liveness.publish(1, "node-9@127.0.0.1", false);
1039 let adopter = Arc::new(FakeAdopter::new(false));
1040 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1041
1042 assert_eq!(
1043 sup.tick().await.len(),
1044 1,
1045 "a shard whose recorded owner is itself down is adoptable"
1046 );
1047 assert_eq!(adopter.calls(), vec![vec![1]]);
1048 }
1049
1050 #[tokio::test]
1056 async fn tick_emits_topology_deltas_through_the_publisher()
1057 -> Result<(), Box<dyn std::error::Error>> {
1058 use std::num::NonZeroUsize;
1059
1060 use aion_core::ClusterEvent;
1061 use futures::StreamExt;
1062
1063 use crate::cluster_publisher::ClusterEventPublisher;
1064
1065 let capacity = NonZeroUsize::new(64).ok_or("non-zero")?;
1066 let publisher = Arc::new(ClusterEventPublisher::new(capacity));
1067 let mut subscription = publisher.subscribe(0);
1068
1069 let liveness = Arc::new(FakeLiveness::new(true));
1070 let adopter = Arc::new(FakeAdopter::new(false));
1071 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2)
1072 .with_publisher(Arc::clone(&publisher), "node-self@127.0.0.1");
1073
1074 liveness.set(false);
1076 drop(sup.tick().await);
1077 let fired = sup.tick().await;
1079 assert_eq!(fired, vec!["node-1@127.0.0.1".to_owned()]);
1080
1081 let first = next_event(&mut subscription).await?;
1083 assert!(
1084 matches!(
1085 &first,
1086 ClusterEvent::PeerDisconnected {
1087 confirmed: false,
1088 consecutive_down: 1,
1089 ..
1090 }
1091 ),
1092 "first delta must be an unconfirmed down: {first:?}"
1093 );
1094 let second = next_event(&mut subscription).await?;
1095 assert!(
1096 matches!(
1097 &second,
1098 ClusterEvent::PeerDisconnected {
1099 confirmed: true,
1100 consecutive_down: 2,
1101 ..
1102 }
1103 ),
1104 "second delta must be the confirmed down: {second:?}"
1105 );
1106 let third = next_event(&mut subscription).await?;
1107 let ClusterEvent::ShardAdopted {
1108 shards,
1109 adopted_by,
1110 from_peer,
1111 ..
1112 } = &third
1113 else {
1114 return Err(format!("third delta must be ShardAdopted: {third:?}").into());
1115 };
1116 assert_eq!(shards, &vec![1]);
1117 assert_eq!(adopted_by, "node-self@127.0.0.1");
1118 assert_eq!(from_peer, "node-1@127.0.0.1");
1119
1120 liveness.set(true);
1123 drop(sup.tick().await);
1124 let recovery = next_event(&mut subscription).await?;
1125 assert!(
1126 matches!(&recovery, ClusterEvent::PeerConnected { .. }),
1127 "recovery delta must be PeerConnected: {recovery:?}"
1128 );
1129
1130 let quiet = sup.tick().await;
1134 assert!(quiet.is_empty());
1135 assert!(
1137 tokio::time::timeout(std::time::Duration::from_millis(50), subscription.next())
1138 .await
1139 .is_err(),
1140 "a steady connected peer must not re-emit PeerConnected every tick"
1141 );
1142 Ok(())
1143 }
1144
1145 async fn next_event(
1146 subscription: &mut futures::stream::BoxStream<
1147 'static,
1148 Result<aion_core::ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
1149 >,
1150 ) -> Result<aion_core::ClusterEvent, Box<dyn std::error::Error>> {
1151 use futures::StreamExt;
1152 tokio::time::timeout(std::time::Duration::from_secs(1), subscription.next())
1153 .await?
1154 .ok_or("cluster subscription ended")?
1155 .map_err(|lag| format!("unexpected lag: {lag:?}").into())
1156 }
1157
1158 #[tokio::test]
1161 async fn shard_published_to_the_dead_peer_itself_is_adopted() {
1162 let liveness = Arc::new(FakeLiveness::new(false));
1163 liveness.publish(1, "node-1@127.0.0.1", false);
1165 let adopter = Arc::new(FakeAdopter::new(false));
1166 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1167
1168 assert_eq!(
1169 sup.tick().await.len(),
1170 1,
1171 "a record naming the dead peer itself is stale and still adoptable"
1172 );
1173 assert_eq!(adopter.calls(), vec![vec![1]]);
1174 }
1175}