// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/snapshot.rs // version: 2 fn snapshot_foundation() -> std::option::Option<(crate::RawTransactionIngestSnapshotPublisher, crate::RawTransactionIngestSnapshotSource)> { 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-snapshot-001") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let settings = crate::RawTransactionIngestSettings::with_defaults(network, worker_id.clone()); let kind = match ksp_worker_api::WorkerKindCode::new(crate::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let mut lifecycle = ksp_worker_api::WorkerLifecycle::new(worker_id, kind); if lifecycle.start().is_err() { return std::option::Option::None; } return std::option::Option::Some(crate::RawTransactionIngestSnapshotPublisher::new(&settings, &lifecycle)); } #[test] fn pre_008_initial_snapshot_and_common_projection_are_exact() { let (publisher, source) = match snapshot_foundation() { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let _publisher = publisher; let concrete = source.current(); assert_eq!(concrete.worker_snapshot().sequence().value(), 0); assert_eq!(concrete.worker_snapshot().state(), ksp_worker_api::WorkerState::Starting); assert_eq!(concrete.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Unknown); assert_eq!(concrete.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Unknown); assert_eq!(concrete.admission_queue_capacity(), crate::DEFAULT_RAW_TRANSACTION_INGEST_ADMISSION_QUEUE_CAPACITY); assert_eq!(concrete.admission_queue_depth(), 0); assert_eq!(concrete.persistence_concurrency(), crate::DEFAULT_RAW_TRANSACTION_INGEST_PERSISTENCE_CONCURRENCY); assert_eq!(concrete.in_flight_persistence(), 0); assert_eq!(concrete.admitted_total(), 0); assert_eq!(concrete.canonicalized_total(), 0); assert_eq!(concrete.persisted_total(), 0); assert_eq!(concrete.entity_inserted_total(), 0); assert_eq!(concrete.entity_already_present_total(), 0); assert_eq!(concrete.entity_skipped_purged_total(), 0); assert_eq!(concrete.observation_inserted_total(), 0); assert_eq!(concrete.observation_already_present_total(), 0); assert_eq!(concrete.content_conflict_total(), 0); assert_eq!(concrete.store_failure_total(), 0); assert_eq!(concrete.source_failure_total(), 0); assert_eq!(concrete.backpressure_wait_total(), 0); let common = ksp_worker_api::WorkerSnapshotSource::current(&source); assert_eq!(&common, concrete.worker_snapshot()); return; } #[test] fn pre_008_checked_counter_exhaustion_is_stable_and_never_wraps() { let result = super::checked_counter(u64::MAX, "admitted_total"); 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_COUNTER_EXHAUSTED); assert_eq!(error.context().len(), 1); assert_eq!(error.context()[0].key(), "field"); assert_eq!(error.context()[0].value(), "admitted_total"); return; } #[tokio::test(flavor = "current_thread")] async fn pre_008_slow_concrete_and_common_listeners_coalesce_to_same_latest_sequence() { let (mut publisher, source) = match snapshot_foundation() { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let observed = source.current().worker_snapshot().sequence(); assert!(publisher.publish_state(ksp_worker_api::WorkerState::Running, 0, 0).is_ok()); assert!(publisher.record_admission_success(ksp_worker_api::WorkerState::Running, 1, 1, false).is_ok()); let concrete = source.wait_for_change(observed).await; let common = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, observed).await; assert_eq!(concrete.worker_snapshot().sequence(), common.sequence()); assert_eq!(concrete.worker_snapshot().sequence().value(), 2); assert_eq!(concrete.worker_snapshot().state(), ksp_worker_api::WorkerState::Running); assert_eq!(concrete.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Healthy); assert_eq!(concrete.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Active); assert_eq!(concrete.admission_queue_depth(), 1); assert_eq!(concrete.in_flight_persistence(), 1); assert_eq!(concrete.admitted_total(), 1); assert_eq!(concrete.canonicalized_total(), 1); return; } #[tokio::test(flavor = "current_thread")] async fn pre_008_terminal_latest_value_is_retained_after_publisher_drop() { let (mut publisher, source) = match snapshot_foundation() { std::option::Option::Some(value) => value, std::option::Option::None => return, }; assert!(publisher.publish_state(ksp_worker_api::WorkerState::Running, 0, 0).is_ok()); assert!(publisher.publish_state(ksp_worker_api::WorkerState::Stopping, 0, 0).is_ok()); assert!(publisher.publish_state(ksp_worker_api::WorkerState::Stopped, 0, 0).is_ok()); std::mem::drop(publisher); let retained = source.current(); assert_eq!(retained.worker_snapshot().state(), ksp_worker_api::WorkerState::Stopped); assert_eq!(retained.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Healthy); assert_eq!(retained.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Idle); let same_terminal = source.wait_for_change(retained.worker_snapshot().sequence()).await; assert_eq!(same_terminal, retained); let common = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, retained.worker_snapshot().sequence()).await; assert_eq!(common, retained.worker_snapshot().clone()); return; } #[test] fn pre_009_source_failure_and_backpressure_counters_advance_without_changing_projection_contract() { let (mut publisher, source) = match snapshot_foundation() { std::option::Option::Some(value) => value, std::option::Option::None => return, }; assert!(publisher.publish_state(ksp_worker_api::WorkerState::Running, 0, 0).is_ok()); assert!(publisher.record_admission_success(ksp_worker_api::WorkerState::Running, 0, 1, true).is_ok()); assert!(publisher.record_source_failure(ksp_worker_api::WorkerState::Running, 0, 1).is_ok()); let concrete = source.current(); assert_eq!(concrete.admitted_total(), 1); assert_eq!(concrete.canonicalized_total(), 1); assert_eq!(concrete.backpressure_wait_total(), 1); assert_eq!(concrete.source_failure_total(), 1); assert_eq!(concrete.worker_snapshot().health(), ksp_worker_api::WorkerHealth::Healthy); assert_eq!(concrete.worker_snapshot().activity(), ksp_worker_api::WorkerActivity::Active); let common = ksp_worker_api::WorkerSnapshotSource::current(&source); assert_eq!(&common, concrete.worker_snapshot()); return; }