// file: crates/ksp-job-backfill-lib/unit_tests/persistence.rs // version: 1 #[derive(Clone, Copy)] enum FakeResponse { Outcome(ksp_store_lib::RawAcquisitionWriteOutcome), Conflict, Failure, } struct FakePersistencePort { network: ksp_store_lib::RawNetworkId, responses: std::sync::Mutex>, calls: std::sync::atomic::AtomicUsize, normal_mode_only: std::sync::atomic::AtomicBool, } impl FakePersistencePort { fn new(network: ksp_store_lib::RawNetworkId, responses: &[FakeResponse]) -> Self { return Self { network, responses: std::sync::Mutex::new(responses.iter().copied().collect()), calls: std::sync::atomic::AtomicUsize::new(0), normal_mode_only: std::sync::atomic::AtomicBool::new(true), }; } fn calls(&self) -> usize { return self.calls.load(std::sync::atomic::Ordering::SeqCst); } fn used_only_normal_mode(&self) -> bool { return self.normal_mode_only.load(std::sync::atomic::Ordering::SeqCst); } } impl super::RawTransactionPersistencePort for FakePersistencePort { 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.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst); if mode != ksp_store_lib::RawTransactionAcquisitionMode::Normal { self.normal_mode_only.store(false, std::sync::atomic::Ordering::SeqCst); } let result = if transaction.reference() != observation.transaction() { std::result::Result::Err(ksp_core_lib::Error::new( ksp_core_lib::ErrorCode::new("test", "reference_mismatch"), "fake persistence reference mismatch", )) } else { let response = match self.responses.lock() { std::result::Result::Ok(mut responses) => responses.pop_front(), std::result::Result::Err(_) => std::option::Option::None, }; match response { std::option::Option::Some(FakeResponse::Outcome(outcome)) => std::result::Result::Ok(outcome), std::option::Option::Some(FakeResponse::Conflict) => { std::result::Result::Err(ksp_core_lib::Error::new(ksp_store_lib::ERROR_CODE_RAW_CONFLICT, "fake canonical content conflict")) }, std::option::Option::Some(FakeResponse::Failure) | std::option::Option::None => { std::result::Result::Err(ksp_core_lib::Error::new(ksp_core_lib::ErrorCode::new("test", "store_failure"), "fake Store failure")) }, } }; return std::boxed::Box::pin(async move { return result; }); } } fn raw_network(value: &str) -> std::option::Option { return match ksp_store_lib::RawNetworkId::new(value) { std::result::Result::Ok(network) => std::option::Option::Some(network), std::result::Result::Err(_) => std::option::Option::None, }; } fn raw_reference(network: &str, signature_byte: u8) -> std::option::Option { let network = match raw_network(network) { std::option::Option::Some(network) => network, std::option::Option::None => return std::option::Option::None, }; return std::option::Option::Some(ksp_store_lib::RawTransactionReference::new(network, ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]))); } fn raw_acquisition_parts( network: &str, signature_byte: u8, observation_byte: u8, ) -> std::option::Option<(ksp_store_lib::RawTransactionReference, ksp_store_lib::RawTransaction, ksp_store_lib::RawTransactionObservation)> { let reference = match raw_reference(network, signature_byte) { std::option::Option::Some(reference) => reference, std::option::Option::None => return std::option::Option::None, }; let format_id = match ksp_store_lib::RawFormatId::new(crate::RAW_TRANSACTION_FORMAT_ID) { std::result::Result::Ok(format_id) => format_id, std::result::Result::Err(_) => return std::option::Option::None, }; let payload = ksp_store_lib::RawPayload::try_new( format_id, crate::RAW_TRANSACTION_FORMAT_VERSION, std::vec![signature_byte].into_boxed_slice(), ksp_store_lib::RawContentHash::new([signature_byte; 32]), ); let payload = match payload { std::result::Result::Ok(payload) => payload, 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(received_at) => received_at, std::result::Result::Err(_) => return std::option::Option::None, }; let provider = match ksp_store_lib::RawProvenanceCode::new("provider") { std::result::Result::Ok(provider) => provider, std::result::Result::Err(_) => return std::option::Option::None, }; let protocol = match ksp_store_lib::RawProvenanceCode::new("solana.http.json_rpc") { std::result::Result::Ok(protocol) => protocol, std::result::Result::Err(_) => return std::option::Option::None, }; let method = match ksp_store_lib::RawProvenanceCode::new("getTransaction") { std::result::Result::Ok(method) => method, std::result::Result::Err(_) => return std::option::Option::None, }; let provenance = ksp_store_lib::RawAcquisitionProvenance::new(provider, protocol, method, ksp_store_lib::RawAcquisitionOrigin::Backfill, received_at); let transaction = ksp_store_lib::RawTransaction::new(reference.clone(), 42, std::option::Option::None, payload); let observation = ksp_store_lib::RawTransactionObservation::new(ksp_store_lib::RawObservationKey::new([observation_byte; 32]), reference.clone(), provenance); return std::option::Option::Some((reference, transaction, observation)); } fn store_outcome( entity: ksp_store_lib::RawEntityWriteOutcome, observation: ksp_store_lib::RawObservationWriteOutcome, ) -> ksp_store_lib::RawAcquisitionWriteOutcome { return ksp_store_lib::RawAcquisitionWriteOutcome::new(entity, observation); } #[tokio::test] async fn pre_007_missing_skips_store_and_preserves_network_scoped_identity() { let reference = match raw_reference("devnet", 1) { std::option::Option::Some(reference) => reference, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let port = FakePersistencePort::new(network, &[]); let result = super::persist_hydration_with_port(&port, crate::BackfillHydrationOutcome::Missing(reference.clone())).await; assert!(result.is_ok()); if let std::result::Result::Ok(result) = result { assert_eq!(result.reference(), &reference); assert_eq!(result.entity(), crate::BackfillEntityPersistence::Missing); assert_eq!(result.observation(), crate::BackfillObservationPersistence::NotApplicable); } assert_eq!(port.calls(), 0); return; } #[tokio::test] async fn pre_007_store_network_mismatch_is_rejected_before_any_write() { let reference = match raw_reference("devnet", 2) { std::option::Option::Some(reference) => reference, std::option::Option::None => return, }; let network = match raw_network("mainnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let port = FakePersistencePort::new(network, &[]); let result = super::persist_hydration_with_port(&port, crate::BackfillHydrationOutcome::Missing(reference)).await; assert!(result.is_err()); if let std::result::Result::Err(error) = result { assert_eq!(error.code(), crate::ERROR_CODE_BACKFILL_PERSISTENCE_INVALID); } assert_eq!(port.calls(), 0); return; } #[tokio::test] async fn pre_007_atomic_insert_maps_entity_and_observation_without_second_write() { let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 3, 13) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let response = store_outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::Inserted); let port = FakePersistencePort::new(network, &[FakeResponse::Outcome(response)]); let result = super::persist_available_with_port(&port, reference.clone(), transaction, observation).await; assert!(result.is_ok()); if let std::result::Result::Ok(result) = result { assert_eq!(result.reference(), &reference); assert_eq!(result.entity(), crate::BackfillEntityPersistence::Inserted); assert_eq!(result.observation(), crate::BackfillObservationPersistence::Inserted); } assert_eq!(port.calls(), 1); assert!(port.used_only_normal_mode()); return; } #[tokio::test] async fn pre_007_existing_entity_distinguishes_new_from_idempotent_observation() { let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let first = store_outcome(ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent, ksp_store_lib::RawObservationWriteOutcome::Inserted); let second = store_outcome(ksp_store_lib::RawEntityWriteOutcome::AlreadyPresent, ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent); let port = FakePersistencePort::new(network, &[FakeResponse::Outcome(first), FakeResponse::Outcome(second)]); let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 4, 14) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let result = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(result.is_ok()); if let std::result::Result::Ok(result) = result { assert_eq!(result.entity(), crate::BackfillEntityPersistence::AlreadyPresent); assert_eq!(result.observation(), crate::BackfillObservationPersistence::Inserted); } let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 4, 14) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let rerun = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(rerun.is_ok()); if let std::result::Result::Ok(rerun) = rerun { assert_eq!(rerun.entity(), crate::BackfillEntityPersistence::AlreadyPresent); assert_eq!(rerun.observation(), crate::BackfillObservationPersistence::AlreadyPresent); } assert_eq!(port.calls(), 2); assert!(port.used_only_normal_mode()); return; } #[tokio::test] async fn pre_007_normal_backfill_respects_purged_tombstone_without_observation() { let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 5, 15) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let response = store_outcome(ksp_store_lib::RawEntityWriteOutcome::SkippedPurged, ksp_store_lib::RawObservationWriteOutcome::NotRecorded); let port = FakePersistencePort::new(network, &[FakeResponse::Outcome(response)]); let result = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(result.is_ok()); if let std::result::Result::Ok(result) = result { assert_eq!(result.entity(), crate::BackfillEntityPersistence::SkippedPurged); assert_eq!(result.observation(), crate::BackfillObservationPersistence::NotRecorded); } assert!(port.used_only_normal_mode()); return; } #[tokio::test] async fn pre_007_store_content_conflict_is_explicit_and_not_idempotent_success() { let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 6, 16) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let port = FakePersistencePort::new(network, &[FakeResponse::Conflict]); let result = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(result.is_ok()); if let std::result::Result::Ok(result) = result { assert_eq!(result.entity(), crate::BackfillEntityPersistence::Conflict); assert_eq!(result.observation(), crate::BackfillObservationPersistence::NotRecorded); assert_ne!(result.entity(), crate::BackfillEntityPersistence::AlreadyPresent); } return; } #[tokio::test] async fn pre_007_non_conflict_store_failure_propagates_unchanged() { let (reference, transaction, observation) = match raw_acquisition_parts("devnet", 7, 17) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let port = FakePersistencePort::new(network, &[FakeResponse::Failure]); let result = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(result.is_err()); if let std::result::Result::Err(error) = result { assert_eq!(error.code(), ksp_core_lib::ErrorCode::new("test", "store_failure")); } return; } #[test] fn pre_007_normal_mode_rejects_impossible_store_outcome_combinations() { let reference = match raw_reference("devnet", 8) { std::option::Option::Some(reference) => reference, std::option::Option::None => return, }; let impossible = [ store_outcome(ksp_store_lib::RawEntityWriteOutcome::Rehydrated, ksp_store_lib::RawObservationWriteOutcome::Inserted), store_outcome(ksp_store_lib::RawEntityWriteOutcome::Inserted, ksp_store_lib::RawObservationWriteOutcome::AlreadyPresent), store_outcome(ksp_store_lib::RawEntityWriteOutcome::SkippedPurged, ksp_store_lib::RawObservationWriteOutcome::Inserted), ]; for outcome in impossible { let result = super::map_store_outcome(reference.clone(), outcome); assert!(result.is_err()); if let std::result::Result::Err(error) = result { assert_eq!(error.code(), crate::ERROR_CODE_BACKFILL_PERSISTENCE_INVALID); } } return; } #[tokio::test] async fn pre_007_mismatched_transaction_observation_reference_is_rejected_before_store() { let (reference, transaction, _) = match raw_acquisition_parts("devnet", 9, 19) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let (_, _, observation) = match raw_acquisition_parts("devnet", 10, 20) { std::option::Option::Some(parts) => parts, std::option::Option::None => return, }; let network = match raw_network("devnet") { std::option::Option::Some(network) => network, std::option::Option::None => return, }; let port = FakePersistencePort::new(network, &[]); let result = super::persist_available_with_port(&port, reference, transaction, observation).await; assert!(result.is_err()); assert_eq!(port.calls(), 0); return; }