137 lines
7.2 KiB
Rust
137 lines
7.2 KiB
Rust
// 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;
|
|
}
|