diff --git a/writer-service/src/db.rs b/writer-service/src/db.rs index d45e13d..cd4e662 100644 --- a/writer-service/src/db.rs +++ b/writer-service/src/db.rs @@ -21,7 +21,7 @@ use solana_program::pubkey::Pubkey; use crate::{ google_storage, merkle_tree_parser, - result::Result, + result::{AppError, Result}, stake_pool_manager::StakePoolManager, tip_distributor_sdk::{GeneratedMerkleTreeCollection, StakeMetaCollection}, }; @@ -85,6 +85,12 @@ pub async fn upsert_to_db( Ok(()) } +/// Fetches and processes MEV claims data for a given epoch. +/// +/// This function attempts to download merkle tree and stake metadata files from multiple +/// GCP servers, trying each server in sequence until one successfully provides valid, +/// parsable files. This fallback mechanism handles cases where files may be corrupted +/// or unavailable on specific servers. pub async fn write_mev_claims_info( db: &Database, target_epoch: u64, @@ -102,50 +108,88 @@ pub async fn write_mev_claims_info( return Ok(()); } - let (merkle_tree_uri, stake_meta_uri) = - google_storage::get_file_uris(target_epoch, mainnet_gcp_server_names).await?; - let client = ReqwestClient::builder() .timeout(CoreDuration::from_secs(600)) .build()?; - let backoff = ExponentialBackoff::default(); - let stake_meta_collection_res: std::result::Result = - retry(backoff.clone(), || async { - let res = client - .get(stake_meta_uri.clone()) + for mainnet_gcp_server_name in mainnet_gcp_server_names { + info!("Trying to fetch files from {mainnet_gcp_server_name} for epoch {target_epoch}"); + + let (merkle_tree_uri, stake_meta_uri) = + match google_storage::get_file_uris(target_epoch, mainnet_gcp_server_name).await { + Ok(uris) => uris, + Err(e) => { + warn!("Files not found on {mainnet_gcp_server_name}: {e}"); + continue; + } + }; + + let backoff = ExponentialBackoff::default(); + + let stake_meta_collection = match retry(backoff.clone(), || async { + client + .get(&stake_meta_uri) .send() - .await? - .json() - .await; - Ok(res) + .await + .map_err(backoff::Error::transient)? + .json::() + .await + .map_err(backoff::Error::transient) }) - .await?; - let merkle_tree_collection_res: std::result::Result< - GeneratedMerkleTreeCollection, - reqwest::Error, - > = retry(backoff, || async { - let res = client - .get(merkle_tree_uri.clone()) - .send() - .await? - .json() - .await; - Ok(res) - }) - .await?; - info!("Successfully fetched merkle tree collection"); - - info!("Starting merkle tree parsing for epoch {target_epoch}"); - merkle_tree_parser::parse_merkle_tree( - db, - target_epoch, - &merkle_tree_collection_res?, - &stake_meta_collection_res?, - tip_distribution_program_id, - priority_fee_distribution_program_id, - ) - .await + .await + { + Ok(data) => data, + Err(e) => { + warn!("Failed to fetch stake_meta from {mainnet_gcp_server_name}: {e}"); + continue; + } + }; + + let merkle_tree_collection = match retry(backoff, || async { + client + .get(&merkle_tree_uri) + .send() + .await + .map_err(backoff::Error::transient)? + .json::() + .await + .map_err(backoff::Error::transient) + }) + .await + { + Ok(data) => data, + Err(e) => { + warn!("Failed to fetch merkle_tree from {mainnet_gcp_server_name}: {e}"); + continue; + } + }; + + info!("Successfully fetched merkle tree collection with {mainnet_gcp_server_name}"); + + match merkle_tree_parser::parse_merkle_tree( + db, + target_epoch, + &merkle_tree_collection, + &stake_meta_collection, + tip_distribution_program_id, + priority_fee_distribution_program_id, + ) + .await + { + Ok(()) => { + info!("Successfully processed files from {mainnet_gcp_server_name}"); + return Ok(()); + } + Err(e) => { + error!("Failed to parse files from {mainnet_gcp_server_name}: {e}"); + continue; + } + } + } + + Err(AppError::WriteMevInfoClaims(format!( + "Failed to process valid files for epoch {target_epoch} from any server" + ))) } pub async fn write_stake_pool_info( diff --git a/writer-service/src/google_storage.rs b/writer-service/src/google_storage.rs index c479d5b..5ec14b6 100644 --- a/writer-service/src/google_storage.rs +++ b/writer-service/src/google_storage.rs @@ -40,10 +40,7 @@ pub fn filter_file( .ok_or(AppError::FileNotFound(name)) } -pub async fn get_file_uris( - epoch: u64, - mainnet_gcp_server_names: &[String], -) -> Result<(String, String)> { +pub async fn get_file_uris(epoch: u64, mainnet_gcp_server_name: &str) -> Result<(String, String)> { let mut all_items = vec![]; let mut next_page_token = String::from(""); let items: Vec = loop { @@ -62,35 +59,29 @@ pub async fn get_file_uris( } }; - let merkle_tree_entry = mainnet_gcp_server_names - .iter() - .find_map(|gcp_name| { - filter_file( - &items, - String::from("merkle-tree"), - epoch, - gcp_name.to_owned(), - ) - .ok() - }) - .ok_or_else(|| { - AppError::FileNotFound(format!("Failed to find merkle-tree file of epoch {epoch}")) - })?; + let merkle_tree_entry = filter_file( + &items, + String::from("merkle-tree"), + epoch, + mainnet_gcp_server_name.to_string(), + ) + .map_err(|e| { + AppError::FileNotFound(format!( + "Failed to find merkle-tree file of epoch {epoch}: {e}" + )) + })?; - let stake_meta_entry = mainnet_gcp_server_names - .iter() - .find_map(|gcp_name| { - filter_file( - &items, - String::from("stake-meta"), - epoch, - gcp_name.to_owned(), - ) - .ok() - }) - .ok_or_else(|| { - AppError::FileNotFound(format!("Failed to find stake-meta file of epoch {epoch}")) - })?; + let stake_meta_entry = filter_file( + &items, + String::from("stake-meta"), + epoch, + mainnet_gcp_server_name.to_string(), + ) + .map_err(|e| { + AppError::FileNotFound(format!( + "Failed to find stake-meta file of epoch {epoch}: {e}" + )) + })?; Ok(( merkle_tree_entry.media_link.to_owned(), diff --git a/writer-service/src/result.rs b/writer-service/src/result.rs index 72dadf3..3f1048f 100644 --- a/writer-service/src/result.rs +++ b/writer-service/src/result.rs @@ -69,6 +69,9 @@ pub enum AppError { #[error(transparent)] KobeCore(#[from] KobeCoreError), + + #[error("Writing MEV claims error: {0}")] + WriteMevInfoClaims(String), } impl From> for AppError {