Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 83 additions & 39 deletions writer-service/src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
};
Expand Down Expand Up @@ -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,
Expand All @@ -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<StakeMetaCollection, reqwest::Error> =
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();
Comment thread
aoikurokawa marked this conversation as resolved.

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::<StakeMetaCollection>()
.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::<GeneratedMerkleTreeCollection>()
.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(
Expand Down
55 changes: 23 additions & 32 deletions writer-service/src/google_storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<GoogleStorageBucketFile> = loop {
Expand All @@ -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(),
Expand Down
3 changes: 3 additions & 0 deletions writer-service/src/result.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,9 @@ pub enum AppError {

#[error(transparent)]
KobeCore(#[from] KobeCoreError),

#[error("Writing MEV claims error: {0}")]
WriteMevInfoClaims(String),
}

impl From<BackoffError<ClientError>> for AppError {
Expand Down
Loading