413 lines
18 KiB
Rust
413 lines
18 KiB
Rust
// file: crates/ksp-job-backfill-lib/unit_tests/discovery.rs
|
|
// version: 6
|
|
|
|
#[derive(Clone, Debug, Eq, PartialEq)]
|
|
struct PageCall {
|
|
before: std::option::Option<std::string::String>,
|
|
until: std::option::Option<std::string::String>,
|
|
limit: std::option::Option<usize>,
|
|
commitment: std::option::Option<ksp_onchain_transport_lib::SolanaCommitment>,
|
|
min_context_slot: std::option::Option<u64>,
|
|
}
|
|
|
|
struct FakeSource {
|
|
pages: std::sync::Mutex<std::collections::VecDeque<std::vec::Vec<super::SignaturePageEntry>>>,
|
|
calls: std::sync::Mutex<std::vec::Vec<PageCall>>,
|
|
}
|
|
|
|
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::<ksp_core_lib::Result<std::vec::Vec<super::SignaturePageEntry>>>().await;
|
|
});
|
|
}
|
|
}
|
|
|
|
impl FakeSource {
|
|
fn new(pages: std::vec::Vec<std::vec::Vec<super::SignaturePageEntry>>) -> 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<PageCall> {
|
|
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<crate::BackfillSignature> {
|
|
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<crate::BackfillRequest> {
|
|
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<std::string::String> {
|
|
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;
|
|
}
|
|
|
|
#[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;
|
|
}
|