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 aion_store_haematite::LockAcquisitionRuling::from_millis(250, 5)?,
915 )
916 .await?,
917 );
918 let workflow_id = aion_core::WorkflowId::new_v4();
919 let envelope = |seq: u64| EventEnvelope {
920 seq,
921 recorded_at: chrono::Utc::now(),
922 workflow_id: workflow_id.clone(),
923 };
924 let events = vec![
925 Event::WorkflowStarted {
926 envelope: envelope(1),
927 workflow_type: String::from("dev_brief"),
928 input: aion_core::Payload::from_json(&serde_json::json!({}))?,
929 run_id: aion_core::RunId::new_v4(),
930 parent_run_id: None,
931 parent_workflow_id: None,
932 package_version: aion_core::PackageVersion::new("a".repeat(64)),
933 },
934 Event::WorkflowFailed {
935 envelope: envelope(2),
936 error: aion_core::WorkflowError {
937 message: String::from("boom"),
938 details: None,
939 },
940 },
941 ];
942 seeder
943 .append(WriteToken::recorder(), &workflow_id, &events, 0)
944 .await?;
945 let row = OutboxRow::pending(
946 workflow_id.clone(),
947 0,
948 String::from("norn_round"),
949 aion_core::Payload::from_json(&serde_json::json!({}))?,
950 chrono::Utc::now(),
951 );
952 let dispatch_key = row.dispatch_key.clone();
953 seeder
954 .append_outbox_batch(std::slice::from_ref(&row))
955 .await?;
956 assert_eq!(seeder.claim_outbox_rows(1).await?.len(), 1);
957
958 let event_store: Arc<dyn aion_store::EventStore> = Arc::clone(&seeder) as _;
959 let engine = Arc::new(
960 EngineBuilder::new()
961 .store_arc(event_store)
962 .in_memory_visibility()
963 .scheduler_threads(1)
964 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
965 Arc::new(aion::signal::ConcreteSignalRouter::new(runtime, handoff))
966 as Arc<dyn SignalRouter>
967 })
968 .build()
969 .await?,
970 );
971 let outbox_store: Arc<dyn OutboxStore> = Arc::clone(&seeder) as _;
972 let adopter =
973 OutboxSettlingAdopter::new(Arc::clone(&engine), Some(Arc::clone(&outbox_store)));
974
975 ShardAdopter::adopt_shards(&adopter, &[42])
976 .await
977 .map_err(|error| format!("adoption must succeed: {error}"))?;
978
979 let state = seeder
980 .outbox_row_state(&dispatch_key)
981 .await?
982 .ok_or("the stranded row must still exist")?;
983 assert_eq!(
984 state.status,
985 OutboxStatus::Cancelled,
986 "the adoption sweep must settle the terminal workflow's stranded row"
987 );
988 engine.shutdown()?;
989 Ok(())
990 }
991
992 #[tokio::test]
993 async fn does_not_adopt_while_peer_connected() {
994 let liveness = Arc::new(FakeLiveness::new(true));
995 let adopter = Arc::new(FakeAdopter::new(false));
996 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2);
997 for _ in 0..5 {
998 assert!(sup.tick().await.is_empty());
999 }
1000 assert!(adopter.calls().is_empty(), "no adoption while peer is up");
1001 }
1002
1003 #[tokio::test]
1004 async fn debounce_requires_consecutive_down_before_adopting() {
1005 let liveness = Arc::new(FakeLiveness::new(true));
1006 let adopter = Arc::new(FakeAdopter::new(false));
1007 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 3);
1008
1009 liveness.set(false);
1010 assert!(sup.tick().await.is_empty(), "tick 1 down: below threshold");
1011 liveness.set(true);
1013 assert!(sup.tick().await.is_empty());
1014 liveness.set(false);
1015 assert!(
1016 sup.tick().await.is_empty(),
1017 "down again, counter reset to 1"
1018 );
1019 assert!(sup.tick().await.is_empty(), "2 consecutive: still below 3");
1020 let fired = sup.tick().await;
1021 assert_eq!(
1022 fired,
1023 vec!["node-1@127.0.0.1".to_owned()],
1024 "3rd consecutive triggers"
1025 );
1026 assert_eq!(adopter.calls(), vec![vec![1]]);
1027 }
1028
1029 #[tokio::test]
1030 async fn adopts_once_then_stays_quiet_while_down() {
1031 let liveness = Arc::new(FakeLiveness::new(false));
1032 let adopter = Arc::new(FakeAdopter::new(false));
1033 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1034 assert_eq!(sup.tick().await.len(), 1, "first down tick adopts");
1035 for _ in 0..5 {
1036 assert!(sup.tick().await.is_empty(), "no re-adopt while still down");
1037 }
1038 assert_eq!(adopter.calls(), vec![vec![1]], "adopted exactly once");
1039 }
1040
1041 #[tokio::test]
1042 async fn failed_adoption_is_retried_next_tick() {
1043 let liveness = Arc::new(FakeLiveness::new(false));
1044 let adopter = Arc::new(FakeAdopter::new(true)); let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1046 assert!(
1047 sup.tick().await.is_empty(),
1048 "first adopt fails, not recorded"
1049 );
1050 assert!(adopter.calls().is_empty());
1051 assert_eq!(sup.tick().await.len(), 1, "retry succeeds next tick");
1052 assert_eq!(adopter.calls(), vec![vec![1]]);
1053 }
1054
1055 #[tokio::test]
1056 async fn peer_with_no_shards_is_not_watched() {
1057 let liveness = Arc::new(FakeLiveness::new(false));
1058 let adopter = Arc::new(FakeAdopter::new(false));
1059 let mut sup = ClusterSupervisor::new(
1060 Arc::clone(&liveness),
1061 Arc::clone(&adopter),
1062 vec![WatchedPeer {
1063 name: "node-2@127.0.0.1".to_owned(),
1064 owned_shards: vec![],
1065 }],
1066 SupervisorConfig {
1067 poll_interval: Duration::from_millis(1),
1068 confirmations: 1,
1069 },
1070 );
1071 assert!(!sup.watches_any());
1072 assert!(sup.tick().await.is_empty());
1073 assert!(adopter.calls().is_empty());
1074 }
1075
1076 #[tokio::test]
1080 async fn shard_already_published_to_live_owner_is_not_adopted() {
1081 let liveness = Arc::new(FakeLiveness::new(false));
1082 liveness.publish(1, "node-9@127.0.0.1", true);
1085 let adopter = Arc::new(FakeAdopter::new(false));
1086 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1087
1088 assert!(
1090 sup.tick().await.is_empty(),
1091 "no adoption fires for a shard a live owner already holds"
1092 );
1093 assert!(
1094 adopter.calls().is_empty(),
1095 "the adopter is never invoked for an already-handled shard"
1096 );
1097 for _ in 0..3 {
1099 assert!(sup.tick().await.is_empty());
1100 }
1101 assert!(adopter.calls().is_empty());
1102 }
1103
1104 #[tokio::test]
1108 async fn shard_published_to_a_down_owner_is_still_adopted() {
1109 let liveness = Arc::new(FakeLiveness::new(false));
1110 liveness.publish(1, "node-9@127.0.0.1", false);
1113 let adopter = Arc::new(FakeAdopter::new(false));
1114 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1115
1116 assert_eq!(
1117 sup.tick().await.len(),
1118 1,
1119 "a shard whose recorded owner is itself down is adoptable"
1120 );
1121 assert_eq!(adopter.calls(), vec![vec![1]]);
1122 }
1123
1124 #[tokio::test]
1130 async fn tick_emits_topology_deltas_through_the_publisher()
1131 -> Result<(), Box<dyn std::error::Error>> {
1132 use std::num::NonZeroUsize;
1133
1134 use aion_core::ClusterEvent;
1135 use futures::StreamExt;
1136
1137 use crate::cluster_publisher::ClusterEventPublisher;
1138
1139 let capacity = NonZeroUsize::new(64).ok_or("non-zero")?;
1140 let publisher = Arc::new(ClusterEventPublisher::new(capacity));
1141 let mut subscription = publisher.subscribe(0);
1142
1143 let liveness = Arc::new(FakeLiveness::new(true));
1144 let adopter = Arc::new(FakeAdopter::new(false));
1145 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 2)
1146 .with_publisher(Arc::clone(&publisher), "node-self@127.0.0.1");
1147
1148 liveness.set(false);
1150 drop(sup.tick().await);
1151 let fired = sup.tick().await;
1153 assert_eq!(fired, vec!["node-1@127.0.0.1".to_owned()]);
1154
1155 let first = next_event(&mut subscription).await?;
1157 assert!(
1158 matches!(
1159 &first,
1160 ClusterEvent::PeerDisconnected {
1161 confirmed: false,
1162 consecutive_down: 1,
1163 ..
1164 }
1165 ),
1166 "first delta must be an unconfirmed down: {first:?}"
1167 );
1168 let second = next_event(&mut subscription).await?;
1169 assert!(
1170 matches!(
1171 &second,
1172 ClusterEvent::PeerDisconnected {
1173 confirmed: true,
1174 consecutive_down: 2,
1175 ..
1176 }
1177 ),
1178 "second delta must be the confirmed down: {second:?}"
1179 );
1180 let third = next_event(&mut subscription).await?;
1181 let ClusterEvent::ShardAdopted {
1182 shards,
1183 adopted_by,
1184 from_peer,
1185 ..
1186 } = &third
1187 else {
1188 return Err(format!("third delta must be ShardAdopted: {third:?}").into());
1189 };
1190 assert_eq!(shards, &vec![1]);
1191 assert_eq!(adopted_by, "node-self@127.0.0.1");
1192 assert_eq!(from_peer, "node-1@127.0.0.1");
1193
1194 liveness.set(true);
1197 drop(sup.tick().await);
1198 let recovery = next_event(&mut subscription).await?;
1199 assert!(
1200 matches!(&recovery, ClusterEvent::PeerConnected { .. }),
1201 "recovery delta must be PeerConnected: {recovery:?}"
1202 );
1203
1204 let quiet = sup.tick().await;
1208 assert!(quiet.is_empty());
1209 assert!(
1211 tokio::time::timeout(std::time::Duration::from_millis(50), subscription.next())
1212 .await
1213 .is_err(),
1214 "a steady connected peer must not re-emit PeerConnected every tick"
1215 );
1216 Ok(())
1217 }
1218
1219 async fn next_event(
1220 subscription: &mut futures::stream::BoxStream<
1221 'static,
1222 Result<aion_core::ClusterEvent, crate::cluster_publisher::ClusterStreamLagged>,
1223 >,
1224 ) -> Result<aion_core::ClusterEvent, Box<dyn std::error::Error>> {
1225 use futures::StreamExt;
1226 tokio::time::timeout(std::time::Duration::from_secs(1), subscription.next())
1227 .await?
1228 .ok_or("cluster subscription ended")?
1229 .map_err(|lag| format!("unexpected lag: {lag:?}").into())
1230 }
1231
1232 #[tokio::test]
1235 async fn shard_published_to_the_dead_peer_itself_is_adopted() {
1236 let liveness = Arc::new(FakeLiveness::new(false));
1237 liveness.publish(1, "node-1@127.0.0.1", false);
1239 let adopter = Arc::new(FakeAdopter::new(false));
1240 let mut sup = supervisor(Arc::clone(&liveness), Arc::clone(&adopter), 1);
1241
1242 assert_eq!(
1243 sup.tick().await.len(),
1244 1,
1245 "a record naming the dead peer itself is stale and still adoptable"
1246 );
1247 assert_eq!(adopter.calls(), vec![vec![1]]);
1248 }
1249}