// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs // version: 10 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.snapshot_source().current().worker_snapshot().state(), 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.snapshot_source().current().worker_snapshot().state(), 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 observer = handle.snapshot_source(); std::mem::drop(handle); loop { let current = observer.current(); if current.worker_snapshot().state().is_terminal() { assert_eq!(current.worker_snapshot().state(), ksp_worker_api::WorkerState::Stopped); return; } let _changed = observer.wait_for_change(current.worker_snapshot().sequence()).await; } } #[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(()); } } }); } }, ) { 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 std::result::Result::Ok(()); }); 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(()); } } }); }, ) { 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; } #[derive(Clone, Copy)] enum RuntimePortResponse { Conflict, StoreFailure, StoreFailureBlocked, SuccessBlocked, } struct RuntimePortActiveGuard<'a> { active: &'a std::sync::atomic::AtomicUsize, } impl std::ops::Drop for RuntimePortActiveGuard<'_> { fn drop(&mut self) { self.active.fetch_sub(1, std::sync::atomic::Ordering::AcqRel); return; } } struct RuntimePersistencePort { active: std::sync::atomic::AtomicUsize, completed: std::sync::atomic::AtomicUsize, max_active: std::sync::atomic::AtomicUsize, network: ksp_store_lib::RawNetworkId, normal_mode_seen: std::sync::atomic::AtomicBool, notify: tokio::sync::Notify, released: std::sync::atomic::AtomicBool, response: RuntimePortResponse, } impl RuntimePersistencePort { fn new(network: ksp_store_lib::RawNetworkId, response: RuntimePortResponse, released: bool) -> Self { return Self { active: std::sync::atomic::AtomicUsize::new(0), completed: std::sync::atomic::AtomicUsize::new(0), max_active: std::sync::atomic::AtomicUsize::new(0), network, normal_mode_seen: std::sync::atomic::AtomicBool::new(false), notify: tokio::sync::Notify::new(), released: std::sync::atomic::AtomicBool::new(released), response, }; } fn release(&self) { self.released.store(true, std::sync::atomic::Ordering::Release); self.notify.notify_waiters(); return; } fn update_max_active(&self, current: usize) { let mut observed = self.max_active.load(std::sync::atomic::Ordering::Acquire); while current > observed { match self.max_active.compare_exchange_weak(observed, current, std::sync::atomic::Ordering::AcqRel, std::sync::atomic::Ordering::Acquire) { std::result::Result::Ok(_) => return, std::result::Result::Err(value) => observed = value, } } return; } } impl crate::RawTransactionIngestPersistencePort for RuntimePersistencePort { fn network_matches(&self, network: &ksp_store_lib::RawNetworkId) -> bool { return &self.network == network; } fn persist_acquisition<'a>( &'a self, _transaction: ksp_store_lib::RawTransaction, _observation: ksp_store_lib::RawTransactionObservation, mode: ksp_store_lib::RawTransactionAcquisitionMode, ) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result> { self.normal_mode_seen.store(mode == ksp_store_lib::RawTransactionAcquisitionMode::Normal, std::sync::atomic::Ordering::Release); let response = self.response; return std::boxed::Box::pin(async move { if matches!(response, RuntimePortResponse::StoreFailureBlocked | RuntimePortResponse::SuccessBlocked) { let current = self.active.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1; let _active_guard = RuntimePortActiveGuard { active: &self.active }; self.update_max_active(current); while !self.released.load(std::sync::atomic::Ordering::Acquire) { self.notify.notified().await; } if matches!(response, RuntimePortResponse::SuccessBlocked) { self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Ok(ksp_store_lib::RawAcquisitionWriteOutcome::new( ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted, )); } } if matches!(response, RuntimePortResponse::Conflict) { self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")); } self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Err(ksp_core_lib::Error::new( ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"), "synthetic store failure", )); }); } fn record_observation<'a>( &'a self, _observation: ksp_store_lib::RawTransactionObservation, ) -> ksp_store_lib::StoreApiFuture<'a, ksp_store_lib::Result> { let response = self.response; return std::boxed::Box::pin(async move { if matches!(response, RuntimePortResponse::StoreFailureBlocked | RuntimePortResponse::SuccessBlocked) { let current = self.active.fetch_add(1, std::sync::atomic::Ordering::AcqRel) + 1; let _active_guard = RuntimePortActiveGuard { active: &self.active }; self.update_max_active(current); while !self.released.load(std::sync::atomic::Ordering::Acquire) { self.notify.notified().await; } if matches!(response, RuntimePortResponse::SuccessBlocked) { self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Ok(ksp_store_lib::RawObservationWriteOutcome::Inserted); } } if matches!(response, RuntimePortResponse::Conflict) { self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "synthetic conflict")); } self.completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::result::Result::Err(ksp_core_lib::Error::new( ksp_core_lib::ErrorCode::new("synthetic_store", "write_failed"), "synthetic store failure", )); }); } } fn runtime_ingress(network: &ksp_store_lib::RawNetworkId, signature_byte: u8) -> std::option::Option { let signature = ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]); let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64( network.clone(), signature, 42, std::option::Option::Some(1_700_000_000), "AQID", ksp_raw_transaction_lib::RawTransactionWireField::Omitted, ksp_raw_transaction_lib::RawTransactionWireField::Omitted, ksp_raw_transaction_lib::RawTransactionWireField::Omitted, ); let provider = match ksp_store_lib::RawProvenanceCode::new("deterministic-harness") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let protocol = match ksp_store_lib::RawProvenanceCode::new("internal") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let method = match ksp_store_lib::RawProvenanceCode::new("runtime") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let received_at = match ksp_store_lib::RawTimestamp::from_unix_millis(1_700_000_000_000) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let provenance = ksp_store_lib::RawAcquisitionProvenance::new(provider, protocol, method, ksp_store_lib::RawAcquisitionOrigin::Live, received_at); return std::option::Option::Some(crate::RawTransactionIngress { material, network: network.clone(), provenance, source_key: [signature_byte; 32] }); } fn settings_with_persistence_concurrency(concurrency: usize) -> std::option::Option { return settings_with_runtime_limits(8, concurrency, std::time::Duration::from_secs(5)); } fn settings_with_runtime_limits( admission_queue_capacity: usize, persistence_concurrency: usize, shutdown_drain_timeout: std::time::Duration, ) -> std::option::Option { let network = match ksp_store_lib::RawNetworkId::new("mainnet") { 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-persistence-001") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; return match crate::RawTransactionIngestSettings::new(network, worker_id, admission_queue_capacity, persistence_concurrency, shutdown_drain_timeout) { std::result::Result::Ok(value) => std::option::Option::Some(value), std::result::Result::Err(_) => std::option::Option::None, }; } async fn wait_for_completed(port: &std::sync::Arc, expected: usize) -> bool { for _ in 0..128 { if port.completed.load(std::sync::atomic::Ordering::Acquire) == expected { return true; } tokio::task::yield_now().await; } return port.completed.load(std::sync::atomic::Ordering::Acquire) == expected; } async fn wait_for_max_active(port: &std::sync::Arc, expected: usize) -> bool { for _ in 0..128 { if port.max_active.load(std::sync::atomic::Ordering::Acquire) == expected { return true; } tokio::task::yield_now().await; } return port.max_active.load(std::sync::atomic::Ordering::Acquire) == expected; } #[tokio::test(flavor = "current_thread")] async fn pre_007_runtime_bounds_in_flight_store_persistence_to_configured_concurrency() { let settings = match settings_with_persistence_concurrency(2) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { for signature_byte in 1..=4 { let ingress = match runtime_ingress(&network, signature_byte) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_max_active(&port, 2).await); assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 2); assert_eq!(port.max_active.load(std::sync::atomic::Ordering::Acquire), 2); port.release(); assert!(wait_for_completed(&port, 4).await); assert!(port.normal_mode_seen.load(std::sync::atomic::Ordering::Acquire)); assert_eq!(port.max_active.load(std::sync::atomic::Ordering::Acquire), 2); 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!(port.active.load(std::sync::atomic::Ordering::Acquire), 0); return; } #[tokio::test(flavor = "current_thread")] async fn pre_007_content_conflict_becomes_terminal_after_private_drain() { let settings = match settings_with_persistence_concurrency(1) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::Conflict, true)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { let ingress = match runtime_ingress(&network, 11) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = handle.snapshot_source(); 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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_CONTENT_CONFLICT)); assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 1); let snapshot = source.current(); assert_eq!(snapshot.content_conflict_total(), 1); assert_eq!(snapshot.store_failure_total(), 0); assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy); return; } #[tokio::test(flavor = "current_thread")] async fn pre_007_store_failure_becomes_terminal_without_exposing_remote_error_text() { let settings = match settings_with_persistence_concurrency(1) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailure, true)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { let ingress = match runtime_ingress(&network, 12) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; 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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED)); assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 1); assert!(!std::format!("{handle:?}").contains("synthetic store failure")); return; } #[tokio::test(flavor = "current_thread")] async fn pre_008_runtime_snapshots_count_successful_pipeline_and_retain_terminal() { let settings = match settings_with_persistence_concurrency(2) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, true)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { for signature_byte in 21..=23 { let ingress = match runtime_ingress(&network, signature_byte) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let concrete_source = handle.snapshot_source(); let common_source = handle.worker_snapshot_source(); assert!(wait_for_completed(&port, 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); let snapshot = concrete_source.current(); assert_eq!(snapshot.worker_snapshot().state(), ksp_worker_api::WorkerState::Stopped); assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Healthy); assert_eq!(snapshot.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Idle); assert_eq!(snapshot.admitted_total(), 3); assert_eq!(snapshot.canonicalized_total(), 3); assert_eq!(snapshot.persisted_total(), 3); assert_eq!(snapshot.entity_inserted_total(), 3); assert_eq!(snapshot.entity_already_present_total(), 0); assert_eq!(snapshot.entity_skipped_purged_total(), 0); assert_eq!(snapshot.observation_inserted_total(), 3); assert_eq!(snapshot.observation_already_present_total(), 0); assert_eq!(snapshot.content_conflict_total(), 0); assert_eq!(snapshot.store_failure_total(), 0); assert_eq!(snapshot.source_failure_total(), 0); assert_eq!(snapshot.backpressure_wait_total(), 0); assert_eq!(snapshot.admission_queue_depth(), 0); assert_eq!(snapshot.in_flight_persistence(), 0); let common = ksp_worker_api::WorkerSnapshotSource::current(&common_source); assert_eq!(&common, snapshot.worker_snapshot()); let retained = concrete_source.wait_for_change(snapshot.worker_snapshot().sequence()).await; assert_eq!(retained, snapshot); return; } #[tokio::test(flavor = "current_thread")] async fn pre_008_fault_snapshots_count_classified_store_failures_and_project_unhealthy() { let settings = match settings_with_persistence_concurrency(1) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailure, true)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { let ingress = match runtime_ingress(&network, 24) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.source_failed")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = handle.snapshot_source(); 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::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED)); let snapshot = source.current(); assert_eq!(snapshot.worker_snapshot().state(), terminal); assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy); assert_eq!(snapshot.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Idle); assert_eq!(snapshot.admitted_total(), 1); assert_eq!(snapshot.canonicalized_total(), 1); assert_eq!(snapshot.persisted_total(), 0); assert_eq!(snapshot.store_failure_total(), 1); assert_eq!(snapshot.content_conflict_total(), 0); return; } #[tokio::test(flavor = "current_thread")] async fn pre_009_source_failure_is_counted_and_late_stop_cannot_replace_terminal_fault() { let settings = match settings("mainnet") { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let handle = match super::start_foundation_with_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::None, move |children: &mut tokio::task::JoinSet>, _stop_receiver, _admission_sender| { let _abort_handle = children.spawn(async move { return std::result::Result::Err(crate::runtime_error("test.source_failed")); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = handle.snapshot_source(); let terminal = match handle.wait_terminal().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_SOURCE_FAILED)); let snapshot = source.current(); assert_eq!(snapshot.source_failure_total(), 1); assert_eq!(snapshot.backpressure_wait_total(), 0); assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy); assert!(!handle.request_stop()); let retained = match handle.wait_terminal().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(retained, terminal); return; } #[tokio::test(flavor = "current_thread")] async fn pre_009_stop_does_not_hide_store_failure_observed_during_bounded_drain() { let settings = match settings_with_persistence_concurrency(1) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::StoreFailureBlocked, false)); let runtime_port: super::PersistencePort = port.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { let ingress = match runtime_ingress(&network, 41) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_max_active(&port, 1).await); assert!(handle.request_stop()); port.release(); let terminal = match handle.wait_terminal().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_STORE_FAILED)); let snapshot = handle.snapshot_source().current(); assert_eq!(snapshot.store_failure_total(), 1); assert_eq!(snapshot.in_flight_persistence(), 0); assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0); return; } #[tokio::test(flavor = "current_thread")] async fn pre_009_saturation_is_observable_without_drop_or_unbounded_admission() { let settings = match settings_with_runtime_limits(1, 1, std::time::Duration::from_secs(5)) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false)); let runtime_port: super::PersistencePort = port.clone(); let source_stage = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let source_stage_for_task = source_stage.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { for signature_byte in 51..=53 { let ingress = match runtime_ingress(&network, signature_byte) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } source_stage_for_task.fetch_add(1, std::sync::atomic::Ordering::AcqRel); } return std::result::Result::Ok(()); }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_max_active(&port, 1).await); for _ in 0..128 { if source_stage.load(std::sync::atomic::Ordering::Acquire) == 2 { break; } tokio::task::yield_now().await; } assert_eq!(source_stage.load(std::sync::atomic::Ordering::Acquire), 2); port.release(); assert!(wait_for_completed(&port, 3).await); for _ in 0..128 { if source_stage.load(std::sync::atomic::Ordering::Acquire) == 3 { break; } tokio::task::yield_now().await; } assert_eq!(source_stage.load(std::sync::atomic::Ordering::Acquire), 3); assert!(handle.request_stop()); let terminal = match handle.wait_terminal().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Stopped); let snapshot = handle.snapshot_source().current(); assert_eq!(snapshot.persisted_total(), 3); assert_eq!(snapshot.source_failure_total(), 0); assert!(snapshot.backpressure_wait_total() >= 1); assert_eq!(snapshot.admission_queue_depth(), 0); assert_eq!(snapshot.in_flight_persistence(), 0); return; } #[tokio::test(flavor = "current_thread")] async fn pre_009_drain_timeout_aborts_and_joins_all_owned_source_and_persistence_tasks() { let settings = match settings_with_runtime_limits(1, 1, crate::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let network = settings.network().clone(); let port = std::sync::Arc::new(RuntimePersistencePort::new(network.clone(), RuntimePortResponse::SuccessBlocked, false)); let runtime_port: super::PersistencePort = port.clone(); let source_active = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let source_active_for_task = source_active.clone(); let handle = match super::start_foundation_with_port_and_source_spawner( settings, tokio::runtime::Handle::current(), std::option::Option::Some(runtime_port), move |children, _stop_receiver, admission_sender| { let _abort_handle = children.spawn(async move { let ingress = match runtime_ingress(&network, 61) { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(crate::runtime_error("test.ingress_invalid")), }; if admission_sender.send(ingress).await.is_err() { return std::result::Result::Err(crate::runtime_error("test.source_failed")); } let _guard = ActiveTaskGuard::new(source_active_for_task); return std::future::pending::>().await; }); }, ) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert!(wait_for_max_active(&port, 1).await); assert!(wait_for_active_count(&source_active, 1).await); assert!(handle.request_stop()); let terminal = match handle.wait_terminal().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(terminal, ksp_worker_api::WorkerState::Faulted(crate::ERROR_CODE_RAW_TRANSACTION_INGEST_DRAIN_TIMEOUT)); let snapshot = handle.snapshot_source().current(); assert_eq!(snapshot.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unhealthy); assert_eq!(snapshot.admission_queue_depth(), 0); assert_eq!(snapshot.in_flight_persistence(), 0); assert_eq!(source_active.load(std::sync::atomic::Ordering::Acquire), 0); assert_eq!(port.active.load(std::sync::atomic::Ordering::Acquire), 0); assert_eq!(port.completed.load(std::sync::atomic::Ordering::Acquire), 0); return; }