// file: crates/ksp-job-backfill-lib/unit_tests/discovery.rs // version: 7 #[derive(Clone, Debug, Eq, PartialEq)] struct PageCall { before: std::option::Option, until: std::option::Option, limit: std::option::Option, commitment: std::option::Option, min_context_slot: std::option::Option, } struct FakeSource { pages: std::sync::Mutex>>, calls: std::sync::Mutex>, } struct PendingSource { calls: std::sync::atomic::AtomicUsize, } impl PendingSource { fn new() -> Self { return Self { calls: std::sync::atomic::AtomicUsize::new(0) }; } fn calls(&self) -> usize { return self.calls.load(std::sync::atomic::Ordering::Acquire); } } impl super::SignaturePageSource for PendingSource { fn fetch_signature_page<'a>( &'a self, _role: &'a ksp_onchain_transport_lib::HttpRoleName, _address: &'a ksp_core_lib::Pubkey, _config: ksp_onchain_transport_lib::SolanaSignaturesForAddressConfig, ) -> super::SignaturePageFuture<'a> { self.calls.fetch_add(1, std::sync::atomic::Ordering::AcqRel); return std::boxed::Box::pin(async { return std::future::pending::>>().await; }); } } impl FakeSource { fn new(pages: std::vec::Vec>) -> Self { return Self { pages: std::sync::Mutex::new(pages.into()), calls: std::sync::Mutex::new(std::vec::Vec::new()) }; } fn calls(&self) -> std::vec::Vec { let guard = self.calls.lock(); return match guard { std::result::Result::Ok(value) => value.clone(), std::result::Result::Err(_) => std::vec::Vec::new(), }; } } impl super::SignaturePageSource for FakeSource { fn fetch_signature_page<'a>( &'a self, _role: &'a ksp_onchain_transport_lib::HttpRoleName, _address: &'a ksp_core_lib::Pubkey, config: ksp_onchain_transport_lib::SolanaSignaturesForAddressConfig, ) -> super::SignaturePageFuture<'a> { let call = PageCall { before: config.before().map(str::to_owned), until: config.until().map(str::to_owned), limit: config.limit(), commitment: config.commitment(), min_context_slot: config.min_context_slot(), }; let calls_result = self.calls.lock(); match calls_result { std::result::Result::Ok(mut calls) => calls.push(call), std::result::Result::Err(_) => { return std::boxed::Box::pin(async { return std::result::Result::Err(ksp_core_lib::Error::new( crate::ERROR_CODE_BACKFILL_DISCOVERY_INVALID, "test call recorder lock poisoned", )); }); }, } let pages_result = self.pages.lock(); let page = match pages_result { std::result::Result::Ok(mut pages) => match pages.pop_front() { std::option::Option::Some(value) => value, std::option::Option::None => std::vec::Vec::new(), }, std::result::Result::Err(_) => { return std::boxed::Box::pin(async { return std::result::Result::Err(ksp_core_lib::Error::new(crate::ERROR_CODE_BACKFILL_DISCOVERY_INVALID, "test page queue lock poisoned")); }); }, }; return std::boxed::Box::pin(async move { return std::result::Result::Ok(page); }); } } fn page_entry(character: char, slot: u64) -> super::SignaturePageEntry { return super::SignaturePageEntry { signature: character.to_string().repeat(crate::MIN_BACKFILL_SIGNATURE_TEXT_BYTES), slot }; } fn signature(character: char) -> std::option::Option { return match crate::BackfillSignature::new(character.to_string().repeat(crate::MIN_BACKFILL_SIGNATURE_TEXT_BYTES)) { std::result::Result::Ok(value) => std::option::Option::Some(value), std::result::Result::Err(_) => std::option::Option::None, }; } fn request(scope: crate::BackfillScope, page_size: usize, max_pages: usize, max_candidates: usize) -> std::option::Option { let job_id = match ksp_job_api::JobId::new("backfill:discovery-test") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let network = match ksp_store_lib::RawNetworkId::new("devnet") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::option::Option::None, }; let result = crate::BackfillRequest::new( job_id, network, ksp_onchain_transport_lib::HttpRoleName::new("history"), crate::BackfillCommitment::Confirmed, scope, page_size, max_pages, max_candidates, 1, std::option::Option::Some(42), ); return match result { std::result::Result::Ok(value) => std::option::Option::Some(value), std::result::Result::Err(_) => std::option::Option::None, }; } fn candidate_signatures(discovery: &crate::BackfillDiscovery) -> std::vec::Vec { let mut signatures = std::vec::Vec::with_capacity(discovery.candidates().len()); for candidate in discovery.candidates() { signatures.push(candidate.identity().signature().as_str().to_owned()); } return signatures; } #[test] fn pre_008_fix_001_signature_page_future_contract_is_send() { fn require_send(_value: T) { return; } let future: super::SignaturePageFuture<'static> = std::boxed::Box::pin(async { return std::result::Result::Ok(std::vec::Vec::new()); }); require_send(future); return; } #[tokio::test] async fn pre_005_latest_paginates_newest_first_and_deduplicates_pages_stably() { let scope = crate::BackfillScope::latest_address(ksp_core_lib::Pubkey::new_from_array([1_u8; 32])); let request = match request(scope, 3, 4, 10) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = FakeSource::new(std::vec![ std::vec![page_entry('6', 60), page_entry('5', 50), page_entry('5', 50)], std::vec![page_entry('4', 40), page_entry('3', 30)], ]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; assert!(result.is_ok()); let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(candidate_signatures(&discovery), std::vec!["6".repeat(64), "5".repeat(64), "4".repeat(64), "3".repeat(64)]); assert_eq!(discovery.pages_fetched(), 2); assert_eq!(discovery.boundary(), crate::BackfillDiscoveryBoundary::RpcBoundary); assert!(!discovery.is_partial()); let calls = source.calls(); assert_eq!(calls.len(), 2); assert_eq!(calls[0].before, std::option::Option::None); let expected_cursor = "5".repeat(64); assert_eq!(calls[1].before.as_deref(), std::option::Option::Some(expected_cursor.as_str())); assert_eq!(calls[0].until, std::option::Option::None); assert_eq!(calls[0].limit, std::option::Option::Some(3)); assert_eq!(calls[0].commitment, std::option::Option::Some(ksp_onchain_transport_lib::SolanaCommitment::Confirmed)); assert_eq!(calls[0].min_context_slot, std::option::Option::Some(42)); return; } #[tokio::test] async fn pre_005_before_uses_exclusive_anchor_then_advances_rpc_cursor() { let anchor = match signature('7') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let scope = crate::BackfillScope::before_address(ksp_core_lib::Pubkey::new_from_array([2_u8; 32]), anchor.clone()); let request = match request(scope, 2, 3, 5) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = FakeSource::new(std::vec![std::vec![page_entry('6', 60), page_entry('5', 50)], std::vec![page_entry('4', 40)]]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(candidate_signatures(&discovery), std::vec!["6".repeat(64), "5".repeat(64), "4".repeat(64)]); let calls = source.calls(); assert_eq!(calls.len(), 2); assert_eq!(calls[0].before.as_deref(), std::option::Option::Some(anchor.as_str())); let expected_cursor = "5".repeat(64); assert_eq!(calls[1].before.as_deref(), std::option::Option::Some(expected_cursor.as_str())); assert_eq!(calls[0].until, std::option::Option::None); return; } #[tokio::test] async fn pre_005_after_keeps_only_nearest_newer_window_and_preserves_rpc_order() { let anchor = match signature('1') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let scope = crate::BackfillScope::after_address(ksp_core_lib::Pubkey::new_from_array([3_u8; 32]), anchor.clone()); let request = match request(scope, 3, 3, 3) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = FakeSource::new(std::vec![ std::vec![page_entry('7', 70), page_entry('6', 60), page_entry('5', 50)], std::vec![page_entry('4', 40), page_entry('3', 30)], ]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(candidate_signatures(&discovery), std::vec!["5".repeat(64), "4".repeat(64), "3".repeat(64)]); assert_eq!(discovery.boundary(), crate::BackfillDiscoveryBoundary::RpcBoundary); assert!(!discovery.is_partial()); let calls = source.calls(); assert_eq!(calls.len(), 2); assert_eq!(calls[0].until.as_deref(), std::option::Option::Some(anchor.as_str())); assert_eq!(calls[1].until.as_deref(), std::option::Option::Some(anchor.as_str())); let expected_cursor = "5".repeat(64); assert_eq!(calls[1].before.as_deref(), std::option::Option::Some(expected_cursor.as_str())); return; } #[tokio::test] async fn pre_005_after_page_bound_is_partial_and_does_not_claim_anchor_completion() { let anchor = match signature('1') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let scope = crate::BackfillScope::after_address(ksp_core_lib::Pubkey::new_from_array([4_u8; 32]), anchor); let request = match request(scope, 2, 2, 3) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = FakeSource::new(std::vec![std::vec![page_entry('7', 70), page_entry('6', 60)], std::vec![page_entry('5', 50), page_entry('4', 40)],]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(candidate_signatures(&discovery), std::vec!["6".repeat(64), "5".repeat(64), "4".repeat(64)]); assert_eq!(discovery.boundary(), crate::BackfillDiscoveryBoundary::AfterAnchorNotReached); assert!(discovery.is_partial()); assert_eq!(discovery.pages_fetched(), 2); return; } #[tokio::test] async fn pre_005_latest_page_bound_is_partial_when_full_pages_leave_more_history_possible() { let scope = crate::BackfillScope::latest_address(ksp_core_lib::Pubkey::new_from_array([5_u8; 32])); let request = match request(scope, 2, 1, 5) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = FakeSource::new(std::vec![std::vec![page_entry('7', 70), page_entry('6', 60)]]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(discovery.boundary(), crate::BackfillDiscoveryBoundary::PageLimit); assert!(discovery.is_partial()); return; } #[tokio::test] async fn pre_005_explicit_scope_never_calls_transport_and_preserves_network_scoped_identity() { let first = match signature('2') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let second = match signature('3') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let scope_result = crate::BackfillScope::explicit_signatures(std::vec![first.clone(), second.clone(), first]); let scope = match scope_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let job_id = match ksp_job_api::JobId::new("backfill:explicit") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let network = match ksp_store_lib::RawNetworkId::new("synthetic") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let request_result = crate::BackfillRequest::new( job_id, network.clone(), ksp_onchain_transport_lib::HttpRoleName::new("unused"), crate::BackfillCommitment::Finalized, scope, 100, 10, 10, 1, std::option::Option::None, ); let request = match request_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = FakeSource::new(std::vec::Vec::new()); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; let discovery = match result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; assert_eq!(discovery.boundary(), crate::BackfillDiscoveryBoundary::ExplicitInput); assert_eq!(discovery.pages_fetched(), 0); assert_eq!(discovery.candidates().len(), 2); assert_eq!(discovery.candidates()[0].identity().network(), &network); assert!(source.calls().is_empty()); return; } #[tokio::test] async fn pre_008_before_resume_uses_checkpoint_cursor_instead_of_original_anchor() { let anchor = match signature('8') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let resume = match signature('5') { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let scope = crate::BackfillScope::before_address(ksp_core_lib::Pubkey::new_from_array([8_u8; 32]), anchor); let request = match request(scope, 2, 2, 4) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let checkpoint = crate::BackfillCheckpoint::new(request.job_id().clone(), request.scope_fingerprint(), 2, std::option::Option::Some(resume.clone())); let request = match request.with_checkpoint(checkpoint) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = FakeSource::new(std::vec![std::vec![page_entry('4', 40)]]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; assert!(result.is_ok()); let calls = source.calls(); assert_eq!(calls.len(), 1); assert_eq!(calls[0].before.as_deref(), std::option::Option::Some(resume.as_str())); return; } #[tokio::test] async fn pre_008_latest_resume_restarts_from_current_latest_without_rpc_cursor() { let scope = crate::BackfillScope::latest_address(ksp_core_lib::Pubkey::new_from_array([9_u8; 32])); let request = match request(scope, 2, 2, 4) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let checkpoint = crate::BackfillCheckpoint::new(request.job_id().clone(), request.scope_fingerprint(), 2, std::option::Option::None); let request = match request.with_checkpoint(checkpoint) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return, }; let source = FakeSource::new(std::vec![std::vec![page_entry('7', 70)]]); let result = super::discover_with_source(&source, &request, std::option::Option::None).await; assert!(result.is_ok()); let calls = source.calls(); assert_eq!(calls.len(), 1); assert_eq!(calls[0].before, std::option::Option::None); return; } #[tokio::test] async fn pre_009_discovery_rpc_wait_is_cancelled_cooperatively() { let scope = crate::BackfillScope::latest_address(ksp_core_lib::Pubkey::new_from_array([9_u8; 32])); let request = match request(scope, 2, 2, 5) { std::option::Option::Some(value) => value, std::option::Option::None => return, }; let source = PendingSource::new(); let token = ksp_job_api::JobCancellationToken::new(); let (cancel_sender, cancel_receiver) = tokio::sync::watch::channel(false); let cancellation = crate::BackfillCancellationSignal::new(token.clone(), cancel_receiver); let discovery = super::discover_with_source(&source, &request, std::option::Option::Some(&cancellation)); let cancel = async { tokio::task::yield_now().await; assert!(token.cancel()); assert!(cancel_sender.send(true).is_ok()); }; let (result, ()) = tokio::join!(discovery, cancel); let error = match result { std::result::Result::Ok(_) => return, std::result::Result::Err(error) => error, }; assert_eq!(error.code(), crate::ERROR_CODE_BACKFILL_CANCELLED); assert_eq!(source.calls(), 1); return; }