// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs // version: 3 struct ActiveTaskGuard { active: std::sync::Arc, } impl ActiveTaskGuard { fn new(active: std::sync::Arc) -> Self { active.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return Self { active }; } } impl std::ops::Drop for ActiveTaskGuard { fn drop(&mut self) { self.active.fetch_sub(1, std::sync::atomic::Ordering::AcqRel); return; } } fn settings(network: &str) -> std::option::Option { let network = match ksp_store_lib::RawNetworkId::new(network) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let worker_id = match ksp_worker_api::WorkerId::new("raw-ingest-runtime-001") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; return std::option::Option::Some(crate::RawTransactionIngestSettings::with_defaults(network, worker_id)); } async fn wait_for_active_count(active: &std::sync::Arc, expected: usize) -> bool { for _ in 0..64 { if active.load(std::sync::atomic::Ordering::Acquire) == expected { return true; } tokio::task::yield_now().await; } return active.load(std::sync::atomic::Ordering::Acquire) == expected; } #[test] fn pre_004_current_runtime_is_required_before_spawn() { let result = super::current_runtime_handle(); assert!(result.is_err()); let error = match result { std::result::Result::Ok(_) => return, std::result::Result::Err(value) => value, }; assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID); assert_eq!(error.context().len(), 1); assert_eq!(error.context()[0].value(), "start.runtime_unavailable"); return; } #[test] fn pre_004_store_network_mismatch_is_rejected_without_echoing_network_values() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let other = match ksp_store_lib::RawNetworkId::new("devnet") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let result = super::validate_store_network(&settings, &other); assert!(result.is_err()); let error = match result { std::result::Result::Ok(_) => return, std::result::Result::Err(value) => value, }; assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID); let rendered = std::format!("{error:?}"); assert!(!rendered.contains("mainnet")); assert!(!rendered.contains("devnet")); return; } #[tokio::test(flavor = "current_thread")] async fn pre_004_immediate_stop_is_idempotent_and_reaches_stopped_terminal() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(*handle.terminal_receiver.borrow(), ksp_worker_api::WorkerState::Starting); assert!(handle.request_stop()); assert!(!handle.request_stop()); let terminal = handle.wait_terminal().await; assert!(terminal.is_ok()); let terminal = match terminal { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); assert!(terminal.is_terminal()); return; } #[tokio::test(flavor = "current_thread")] async fn pre_004_running_lifecycle_stops_through_cloned_handle() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; tokio::task::yield_now().await; assert_eq!(*handle.terminal_receiver.borrow(), ksp_worker_api::WorkerState::Running); let clone = handle.clone(); assert!(clone.request_stop()); assert!(!handle.request_stop()); let terminal = handle.wait_terminal().await; let terminal = match terminal { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); return; } #[tokio::test(flavor = "current_thread")] async fn pre_004_dropping_last_control_handle_causes_private_runtime_exit() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let handle = match super::start_foundation(settings, tokio::runtime::Handle::current(), std::option::Option::None) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let mut observer = handle.terminal_receiver.clone(); std::mem::drop(handle); loop { let current = *observer.borrow(); if current.is_terminal() { let closed = observer.changed().await; assert!(closed.is_err()); assert_eq!(current, ksp_worker_api::WorkerState::Stopped); return; } let changed = observer.changed().await; assert!(changed.is_ok(), "runtime closed before publishing terminal state"); if changed.is_err() { return; } } } #[tokio::test(flavor = "current_thread")] async fn pre_005_supervisor_joins_all_cooperative_source_tasks_before_terminal() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let source_active = std::sync::Arc::clone(&active); let handle = match super::start_foundation_with_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::None, move |children, stop_receiver, _admission_sender| { for _ in 0..3 { let active = std::sync::Arc::clone(&source_active); let mut child_stop = stop_receiver.clone(); let _abort_handle = children.spawn(async move { let _guard = ActiveTaskGuard::new(active); loop { let changed = child_stop.changed().await; if changed.is_err() || *child_stop.borrow() { return; } } }); } }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_active_count(&active, 3).await); assert!(handle.request_stop()); let terminal = handle.wait_terminal().await; let terminal = match terminal { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); assert_eq!(active.load(std::sync::atomic::Ordering::Acquire), 0); return; } #[tokio::test(flavor = "current_thread")] async fn pre_005_supervisor_reaps_completed_source_task_and_still_joins_live_child_on_stop() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let completed = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let source_active = std::sync::Arc::clone(&active); let source_completed = std::sync::Arc::clone(&completed); let handle = match super::start_foundation_with_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::None, move |children, stop_receiver, _admission_sender| { let completed = std::sync::Arc::clone(&source_completed); let _completed_abort_handle = children.spawn(async move { completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return; }); let active = std::sync::Arc::clone(&source_active); let mut child_stop = stop_receiver.clone(); let _active_abort_handle = children.spawn(async move { let _guard = ActiveTaskGuard::new(active); loop { let changed = child_stop.changed().await; if changed.is_err() || *child_stop.borrow() { return; } } }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_active_count(&active, 1).await); for _ in 0..64 { if completed.load(std::sync::atomic::Ordering::Acquire) == 1 { break; } tokio::task::yield_now().await; } assert_eq!(completed.load(std::sync::atomic::Ordering::Acquire), 1); assert!(handle.request_stop()); let terminal = handle.wait_terminal().await; let terminal = match terminal { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); assert_eq!(active.load(std::sync::atomic::Ordering::Acquire), 0); return; }