From c7f9ee50d99dcc315f7a6d8c8176d6c866dae159 Mon Sep 17 00:00:00 2001 From: aoikurokawa Date: Thu, 11 Dec 2025 05:49:07 +0900 Subject: [PATCH 1/2] fix: add parsing failure recovery --- writer-service/src/db.rs | 129 ++++++++++++++++++--------- writer-service/src/google_storage.rs | 55 +++++------- 2 files changed, 109 insertions(+), 75 deletions(-) diff --git a/writer-service/src/db.rs b/writer-service/src/db.rs index d45e13d..5fabd46 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, +/// parseable 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,87 @@ 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 client = ReqwestClient::builder() + .timeout(CoreDuration::from_secs(600)) + .build()?; + 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::FileNotFound(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(), From d9508570c09b221bbbf8d0a82a57683c3e5eac2b Mon Sep 17 00:00:00 2001 From: aoikurokawa Date: Thu, 11 Dec 2025 06:04:38 +0900 Subject: [PATCH 2/2] fix: addressed review --- writer-service/src/db.rs | 11 ++++++----- writer-service/src/result.rs | 3 +++ 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/writer-service/src/db.rs b/writer-service/src/db.rs index 5fabd46..cd4e662 100644 --- a/writer-service/src/db.rs +++ b/writer-service/src/db.rs @@ -89,7 +89,7 @@ pub async fn upsert_to_db( /// /// 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, -/// parseable files. This fallback mechanism handles cases where files may be corrupted +/// 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, @@ -108,6 +108,10 @@ pub async fn write_mev_claims_info( return Ok(()); } + let client = ReqwestClient::builder() + .timeout(CoreDuration::from_secs(600)) + .build()?; + for mainnet_gcp_server_name in mainnet_gcp_server_names { info!("Trying to fetch files from {mainnet_gcp_server_name} for epoch {target_epoch}"); @@ -120,9 +124,6 @@ pub async fn write_mev_claims_info( } }; - let client = ReqwestClient::builder() - .timeout(CoreDuration::from_secs(600)) - .build()?; let backoff = ExponentialBackoff::default(); let stake_meta_collection = match retry(backoff.clone(), || async { @@ -186,7 +187,7 @@ pub async fn write_mev_claims_info( } } - Err(AppError::FileNotFound(format!( + Err(AppError::WriteMevInfoClaims(format!( "Failed to process valid files for epoch {target_epoch} from any server" ))) } 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 {