diff --git a/apps/sequencer/src/feeds/feed_workers.rs b/apps/sequencer/src/feeds/feed_workers.rs index 07e51fc910..0678db5b4d 100644 --- a/apps/sequencer/src/feeds/feed_workers.rs +++ b/apps/sequencer/src/feeds/feed_workers.rs @@ -5,7 +5,7 @@ use crate::feeds::feeds_slots_manager::feeds_slots_manager_loop; use crate::feeds::votes_result_sender::votes_result_sender_loop; use crate::metrics_collector::metrics_collector_loop; use crate::providers::eth_send_utils::{ - create_and_collect_relayers_futures, BatchOfUpdatesToProcess, + create_and_collect_relayers_futures, create_per_network_reorg_trackers, BatchOfUpdatesToProcess, }; use crate::sequencer_state::SequencerState; use actix_web::web::Data; @@ -84,5 +84,7 @@ pub async fn prepare_app_workers( ) .await; + create_per_network_reorg_trackers(&collected_futures, sequencer_state.clone()).await; + collected_futures } diff --git a/apps/sequencer/src/http_handlers/admin.rs b/apps/sequencer/src/http_handlers/admin.rs index 0d1eed1285..bcc273895c 100644 --- a/apps/sequencer/src/http_handlers/admin.rs +++ b/apps/sequencer/src/http_handlers/admin.rs @@ -975,6 +975,7 @@ mod tests { publishing_criteria: vec![], should_load_rb_indices: true, contracts, + reorg: blocksense_config::ReorgConfig::default(), } }); diff --git a/apps/sequencer/src/providers/eth_send_utils.rs b/apps/sequencer/src/providers/eth_send_utils.rs index d8dd0a346c..43683f1758 100644 --- a/apps/sequencer/src/providers/eth_send_utils.rs +++ b/apps/sequencer/src/providers/eth_send_utils.rs @@ -6,12 +6,11 @@ use alloy::{ providers::{Provider, ProviderBuilder}, rpc::types::{eth::TransactionRequest, TransactionReceipt}, }; - -use alloy_primitives::{keccak256, FixedBytes, TxHash, U256}; -use blocksense_config::{FeedStrideAndDecimals, GNOSIS_SAFE_CONTRACT_NAME}; +use alloy_primitives::{keccak256, FixedBytes, TxHash, B256, U256}; +use blocksense_config::{FeedStrideAndDecimals, ReorgConfig, GNOSIS_SAFE_CONTRACT_NAME}; use blocksense_data_feeds::feeds_processing::{BatchedAggregatesToSend, VotedFeedUpdate}; use blocksense_registry::config::FeedConfig; -use blocksense_utils::{counter_unbounded_channel::CountedReceiver, EncodedFeedId}; +use blocksense_utils::{await_time, counter_unbounded_channel::CountedReceiver, EncodedFeedId}; use eyre::{bail, eyre, Result}; use std::{collections::HashMap, collections::HashSet, mem, sync::Arc}; use tokio::{ @@ -20,13 +19,15 @@ use tokio::{ }; use crate::{ - providers::provider::{ - parse_eth_address, HashValue, ProviderStatus, ProviderType, ProvidersMetrics, RpcProvider, - SharedRpcProviders, + providers::{ + provider::{ + parse_eth_address, HashValue, LatestRBIndex, ProviderStatus, ProviderType, + ProvidersMetrics, RpcProvider, SharedRpcProviders, + }, + reorg_tracking::ReorgTracker, }, sequencer_state::SequencerState, }; -use blocksense_feed_registry::registry::await_time; use blocksense_feeds_processing::adfs_gen_calldata::{ adfs_serialize_updates, get_neighbour_feed_ids, RoundBufferIndices, }; @@ -150,6 +151,7 @@ pub async fn get_serialized_updates_for_network( Ok(serialized_updates) } +#[derive(Clone)] pub struct BatchOfUpdatesToProcess { pub net: String, pub provider: Arc>, @@ -169,6 +171,7 @@ pub async fn create_and_collect_relayers_futures( ) { for (net, chan) in relayers_recv_channels.into_iter() { let feed_metrics_clone = feeds_metrics.clone(); + let net_clone = net.clone(); let provider_status_clone = provider_status.clone(); let relayer_name = format!("relayer_for_network {net}"); collected_futures.push( @@ -176,7 +179,7 @@ pub async fn create_and_collect_relayers_futures( .name(relayer_name.clone().as_str()) .spawn(async move { loop_processing_batch_of_updates( - net, + net_clone, relayer_name, feed_metrics_clone, provider_status_clone, @@ -185,7 +188,7 @@ pub async fn create_and_collect_relayers_futures( .await; Ok(()) }) - .expect("Failed to spawn metrics collector loop!"), + .expect("Failed to spawn {net} network relayer loop!"), ); } } @@ -200,7 +203,7 @@ pub async fn loop_processing_batch_of_updates( tracing::info!("Starting {relayer_name} loop..."); //TODO: Create a termination reason pattern in the future. At this point networks are not added/removed dynamically in the sequencer, - // therefore the loop in iterating over the lifetime of the sequencer. + // therefore the loop is iterating over the lifetime of the sequencer. loop { let cmd_opt = chan.recv().await; let msgs_in_queue = chan.len(); @@ -209,17 +212,7 @@ pub async fn loop_processing_batch_of_updates( let block_height = cmd.updates.block_height; tracing::info!("Processing updates for network {relayer_name}, block_height {block_height}, messages in queue = {msgs_in_queue}"); let provider = cmd.provider.clone(); - let result = eth_batch_send_to_contract( - cmd.net, - cmd.provider, - cmd.provider_settings, - cmd.updates, - cmd.feeds_config, - cmd.transaction_retry_timeout_secs, - cmd.transaction_retries_count_limit, - cmd.retry_fee_increment_fraction, - ) - .await; + let result = eth_batch_send_to_contract(cmd).await; let provider_metrics = provider.lock().await.provider_metrics.clone(); dec_metric!(provider_metrics, net, num_transactions_in_queue); @@ -326,17 +319,66 @@ pub async fn check_tx_hashes_for_inclusion( None } +pub async fn create_per_network_reorg_trackers( + collected_futures: &FuturesUnordered>>, + sequencer_state: Data, +) { + let providers_mutex = sequencer_state.providers.clone(); + let providers = providers_mutex.read().await; + + for (net, _p) in providers.iter() { + let reorg_trackers_name = format!("reorg_tracker for {net}"); + let net_clone = net.clone(); + let sequencer_state_providers_clone = sequencer_state.providers.clone(); + let sequencer_config = sequencer_state.sequencer_config.read().await; + let reorg_tracker_config = match sequencer_config.providers.get(net.as_str()) { + Some(c) => c.reorg.clone(), + None => { + error!("No config for provider for network {net} will set to default!"); + ReorgConfig::default() + } + }; + let relayer_send_channel = match sequencer_state + .relayers_send_channels + .read() + .await + .get(net.as_str()) + { + Some(chan) => chan.clone(), + None => { + panic!("Failed to spawn tracker for reorgs loop in network {net} because no updates sending relayer exists for it!"); + } + }; + let mut reorg_tracker = ReorgTracker::new( + net_clone, + reorg_tracker_config, + sequencer_state_providers_clone, + relayer_send_channel, + ); + collected_futures.push( + tokio::task::Builder::new() + .name(reorg_trackers_name.clone().as_str()) + .spawn(async move { + reorg_tracker.loop_tracking_for_reorg_in_network().await; + Ok(()) + }) + .expect("Failed to spawn tracker for reorgs loop in network {net}!"), + ); + } +} + #[allow(clippy::too_many_arguments)] pub async fn eth_batch_send_to_contract( - net: String, - provider_mutex: Arc>, - provider_settings: blocksense_config::Provider, - mut updates: BatchedAggregatesToSend, - feeds_config: Arc>>, - transaction_retry_timeout_secs: u64, - transaction_retries_count_limit: u64, - retry_fee_increment_fraction: f64, + cmd: BatchOfUpdatesToProcess, ) -> Result<(String, Vec)> { + let net = cmd.net.clone(); + let provider_mutex = cmd.provider.clone(); + let provider_settings = cmd.provider_settings.clone(); + let feeds_config = cmd.feeds_config.clone(); + let transaction_retry_timeout_secs = cmd.transaction_retry_timeout_secs; + let transaction_retries_count_limit = cmd.transaction_retries_count_limit; + let retry_fee_increment_fraction = cmd.retry_fee_increment_fraction; + let mut updates = cmd.updates.clone(); let mut feeds_rb_indices = HashMap::new(); let serialized_updates = get_serialized_updates_for_network( net.as_str(), @@ -487,6 +529,8 @@ pub async fn eth_batch_send_to_contract( break nonce; }; + let mut inclusion_block = None; + let mut inclusion_block_hash: Option = None; let mut generated_transaction_hashes = Vec::new(); loop { @@ -521,16 +565,17 @@ pub async fn eth_batch_send_to_contract( ); } + // If the nonce in the contract increased and the next state root hash is not as we expect, + // another sequencer was able to post updates for the current block height before this one. + // We need to take this into account and reread the round counters of the feeds. + info!("Updates to contract already posted, network {net}, block_height {block_height}, latest_nonce {latest_nonce}, previous_nonce {nonce}, merkle_root in contract {prev_calldata_merkle_tree_root:?}"); // TODO: maybe move into an else clause of the `if` above, // i.e. only do it when there were no included transactions found try_to_sync( net.as_str(), &mut provider, &contract_address, - block_height, - &next_calldata_merkle_tree_root, - latest_nonce, - nonce, + Some(&next_calldata_merkle_tree_root), ) .await; @@ -689,14 +734,12 @@ pub async fn eth_batch_send_to_contract( Err(err) => { warn!("Error while submitting transaction in network `{net}` block height {block_height} and address {sender_address} due to {err}"); if err.to_string().contains("execution revert") { + info!("Trying to sync due to tx revert, network {net}, block_height {block_height}, latest_nonce {latest_nonce}, previous_nonce {nonce}, merkle_root in contract {prev_calldata_merkle_tree_root:?}"); try_to_sync( net.as_str(), &mut provider, &contract_address, - block_height, - &next_calldata_merkle_tree_root, - latest_nonce, - nonce, + Some(&next_calldata_merkle_tree_root), ) .await; return Ok(("false".to_string(), feeds_to_update_ids)); @@ -765,6 +808,18 @@ pub async fn eth_batch_send_to_contract( tx_receipt }; + if let Some(eth_block_number) = tx_receipt.block_number { + info!("Transaction was included in block #{}", eth_block_number); + inclusion_block = Some(eth_block_number); + if let Some(h) = tx_receipt.block_hash { + inclusion_block_hash = Some(h); + } else { + error!("Receipt has no block hash!"); + } + } else { + error!("Receipt has no block number!"); + } + receipt = tx_receipt; break; @@ -777,6 +832,12 @@ pub async fn eth_batch_send_to_contract( log_gas_used(&net, &receipt, transaction_time, provider_metrics).await; + if let Some(b) = inclusion_block { + provider.insert_non_finalized_update(b, cmd); + if let Some(h) = inclusion_block_hash { + provider.insert_observed_block_hash(b, h); + } + } provider.update_history(&updates.updates); let result = receipt.status().to_string(); if result == "true" { @@ -888,14 +949,11 @@ pub async fn get_gas_limit( } } -async fn try_to_sync( +pub(crate) async fn try_to_sync( net: &str, provider: &mut RpcProvider, contract_address: &Address, - block_height: u64, - next_calldata_merkle_tree_root: &HashValue, - latest_nonce: u64, - previous_nonce: u64, + next_calldata_merkle_tree_root: Option<&HashValue>, ) { let rpc_handle = &provider.provider; match rpc_handle @@ -903,13 +961,30 @@ async fn try_to_sync( .await { Ok(root) => { - if root != next_calldata_merkle_tree_root.0.into() { - // If the nonce in the contract increased and the next state root hash is not as we expect, - // another sequencer was able to post updates for the current block height before this one. - // We need to take this into account and reread the round counters of the feeds. - info!("Updates to contract already posted, network {net}, block_height {block_height}, latest_nonce {latest_nonce}, previous_nonce {previous_nonce}, merkle_root in contract {root}"); + if next_calldata_merkle_tree_root.is_none_or(|val| root != val.0.into()) { provider.merkle_root_in_contract = Some(HashValue(root.into())); - // TODO: Read round counters from contract + let keys: Vec = if provider.rb_indices.is_empty() { + provider.feeds_variants.keys().cloned().collect() + } else { + provider.rb_indices.keys().cloned().collect() + }; + + let mut new_indices = provider.rb_indices.clone(); + for encoded_feed_id in keys { + match provider.get_latest_rb_index(&encoded_feed_id).await { + Ok(LatestRBIndex { + encoded_feed_id, + index, + }) => { + new_indices.insert(encoded_feed_id, index as u64); + } + Err(e) => { + warn!("Failed to refresh rb index for feed {encoded_feed_id} in network `{net}`: {e:?}"); + } + } + } + provider.rb_indices = new_indices; + debug!("Refreshed rb indices from chain for `{net}`"); } } Err(e) => { diff --git a/apps/sequencer/src/providers/inflight_observations.rs b/apps/sequencer/src/providers/inflight_observations.rs new file mode 100644 index 0000000000..961169fc28 --- /dev/null +++ b/apps/sequencer/src/providers/inflight_observations.rs @@ -0,0 +1,174 @@ +use alloy_primitives::B256; +use std::collections::HashMap; + +use crate::providers::eth_send_utils::BatchOfUpdatesToProcess; + +/// Holds information tied to blocks that may be reorged (non-finalized) +#[derive(Default)] +pub struct InflightObservations { + pub non_finalized_updates: HashMap, + pub observed_block_hashes: HashMap, +} + +impl InflightObservations { + pub fn new() -> Self { + Self { + non_finalized_updates: HashMap::new(), + observed_block_hashes: HashMap::new(), + } + } + + pub fn insert_non_finalized_update(&mut self, block_height: u64, cmd: BatchOfUpdatesToProcess) { + self.non_finalized_updates.insert(block_height, cmd); + } + + pub fn prune_observed_up_to(&mut self, finalized_block: u64) -> usize { + let keys_to_remove: Vec = self + .non_finalized_updates + .keys() + .copied() + .filter(|height| *height <= finalized_block) + .collect(); + let removed = keys_to_remove.len(); + for height in keys_to_remove { + tracing::info!("Pruning observed non_finalized_updates for block_height = {height}; finalized_block = {finalized_block}"); + self.non_finalized_updates.remove(&height); + } + + let hash_keys_to_remove: Vec = self + .observed_block_hashes + .keys() + .copied() + .filter(|height| *height <= finalized_block) + .collect(); + for height in hash_keys_to_remove { + tracing::info!("Pruning observed block hash for block_height = {height}; finalized_block = {finalized_block}"); + self.observed_block_hashes.remove(&height); + } + + removed + } + + pub fn insert_observed_block_hash(&mut self, block_height: u64, hash: B256) { + self.observed_block_hashes.insert(block_height, hash); + } +} + +#[cfg(test)] +mod tests { + use super::*; + use alloy::signers::local::PrivateKeySigner; + use alloy_primitives::B256; + use blocksense_config::{AllFeedsConfig, Provider as ProviderConfig}; + use blocksense_data_feeds::feeds_processing::BatchedAggregatesToSend; + use blocksense_metrics::metrics::ProviderMetrics; + use blocksense_registry::config::FeedConfig; + use blocksense_utils::EncodedFeedId; + use reqwest::Url; + use std::collections::HashMap; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::Arc; + use tokio::sync::{Mutex, RwLock}; + + use crate::providers::provider::RpcProvider; + + static METRICS_COUNTER: AtomicUsize = AtomicUsize::new(0); + + fn next_metrics_prefix() -> String { + let suffix = METRICS_COUNTER.fetch_add(1, Ordering::SeqCst); + format!("test_inflight_{}__", suffix) + } + + async fn dummy_batch() -> BatchOfUpdatesToProcess { + let provider_config = ProviderConfig { + private_key_path: "dummy".to_string(), + url: "http://localhost:8545".to_string(), + transaction_retries_count_limit: 3, + transaction_retry_timeout_secs: 10, + retry_fee_increment_fraction: 0.1, + transaction_gas_limit: 1_000_000, + impersonated_anvil_account: None, + is_enabled: true, + should_load_rb_indices: true, + allow_feeds: None, + publishing_criteria: Vec::new(), + contracts: Vec::new(), + receipt_polling_back_off_period_ms: 15000, + transaction_retry_back_off_ms: 15000, + reorg: blocksense_config::ReorgConfig::default(), + }; + let provider_config_clone = provider_config.clone(); + + let signer = PrivateKeySigner::from_bytes(&B256::from_slice(&[1u8; 32])).unwrap(); + let metrics_prefix = next_metrics_prefix(); + let provider_metrics = + Arc::new(RwLock::new(ProviderMetrics::new(&metrics_prefix).unwrap())); + let feeds_config = AllFeedsConfig { feeds: vec![] }; + + let rpc_provider = RpcProvider::new( + "test-network", + Url::parse("http://localhost:8545").unwrap(), + &signer, + &provider_config, + &provider_metrics, + &feeds_config, + ) + .await; + + BatchOfUpdatesToProcess { + net: "test-network".to_string(), + provider: Arc::new(Mutex::new(rpc_provider)), + provider_settings: provider_config_clone, + updates: BatchedAggregatesToSend { + block_height: 0, + updates: vec![], + }, + feeds_config: Arc::new(RwLock::new(HashMap::::new())), + transaction_retry_timeout_secs: 10, + transaction_retries_count_limit: 3, + retry_fee_increment_fraction: 0.1, + } + } + + #[tokio::test] + async fn insert_non_finalized_update_stores_entry() { + let mut inflight = InflightObservations::new(); + let batch = dummy_batch().await; + + inflight.insert_non_finalized_update(42, batch); + + assert_eq!(inflight.non_finalized_updates.len(), 1); + assert!(inflight.non_finalized_updates.contains_key(&42)); + } + + #[tokio::test] + async fn prune_observed_up_to_removes_data_and_returns_removed_count() { + let mut inflight = InflightObservations::new(); + inflight.insert_observed_block_hash(1, B256::from_slice(&[0x11; 32])); + inflight.insert_observed_block_hash(2, B256::from_slice(&[0x22; 32])); + inflight.insert_observed_block_hash(3, B256::from_slice(&[0x33; 32])); + let batch_one = dummy_batch().await; + let batch_two = dummy_batch().await; + inflight.insert_non_finalized_update(1, batch_one); + inflight.insert_non_finalized_update(3, batch_two); + + let removed = inflight.prune_observed_up_to(2); + + assert_eq!(removed, 1); + assert!(!inflight.non_finalized_updates.contains_key(&1)); + assert!(inflight.non_finalized_updates.contains_key(&3)); + assert!(!inflight.observed_block_hashes.contains_key(&1)); + assert!(!inflight.observed_block_hashes.contains_key(&2)); + assert!(inflight.observed_block_hashes.contains_key(&3)); + } + + #[tokio::test] + async fn insert_observed_block_hash_overwrites_previous_value() { + let mut inflight = InflightObservations::new(); + inflight.insert_observed_block_hash(7, B256::from_slice(&[0x44; 32])); + inflight.insert_observed_block_hash(7, B256::from_slice(&[0x55; 32])); + + let stored = inflight.observed_block_hashes.get(&7).unwrap(); + assert_eq!(stored, &B256::from_slice(&[0x55; 32])); + } +} diff --git a/apps/sequencer/src/providers/mod.rs b/apps/sequencer/src/providers/mod.rs index 3749fb7618..6a9983a8f2 100644 --- a/apps/sequencer/src/providers/mod.rs +++ b/apps/sequencer/src/providers/mod.rs @@ -1,2 +1,4 @@ pub mod eth_send_utils; +pub mod inflight_observations; pub mod provider; +pub mod reorg_tracking; diff --git a/apps/sequencer/src/providers/provider.rs b/apps/sequencer/src/providers/provider.rs index d9ef447031..5ba083dac9 100644 --- a/apps/sequencer/src/providers/provider.rs +++ b/apps/sequencer/src/providers/provider.rs @@ -44,7 +44,10 @@ use tokio::time::error::Elapsed; use tokio::time::Duration; use tracing::{debug, error, info, warn}; -use crate::providers::eth_send_utils::{get_gas_limit, get_tx_retry_params, GasFees}; +use crate::providers::eth_send_utils::{ + get_gas_limit, get_tx_retry_params, BatchOfUpdatesToProcess, GasFees, +}; +use crate::providers::inflight_observations::InflightObservations; use std::time::Instant; pub type ProviderType = @@ -122,6 +125,7 @@ pub struct RpcProvider { pub rpc_url: Url, pub rb_indices: RoundBufferIndices, num_tx_in_progress: u32, + pub inflight: InflightObservations, } #[derive(PartialEq, Debug, Serialize, Deserialize)] @@ -354,6 +358,7 @@ impl RpcProvider { rpc_url, rb_indices: RoundBufferIndices::new(), num_tx_in_progress: 0, + inflight: InflightObservations::new(), } } @@ -846,6 +851,18 @@ impl RpcProvider { pub fn get_num_tx_in_progress(&self) -> u32 { self.num_tx_in_progress } + + pub fn insert_non_finalized_update(&mut self, block_height: u64, cmd: BatchOfUpdatesToProcess) { + self.inflight.insert_non_finalized_update(block_height, cmd) + } + + pub fn prune_observed_up_to(&mut self, finalized_block: u64) -> usize { + self.inflight.prune_observed_up_to(finalized_block) + } + + pub fn insert_observed_block_hash(&mut self, block_height: u64, hash: B256) { + self.inflight.insert_observed_block_hash(block_height, hash) + } } async fn perform_post_deployment_transaction( diff --git a/apps/sequencer/src/providers/reorg_tracking.rs b/apps/sequencer/src/providers/reorg_tracking.rs new file mode 100644 index 0000000000..e72a0d67fc --- /dev/null +++ b/apps/sequencer/src/providers/reorg_tracking.rs @@ -0,0 +1,1237 @@ +use crate::providers::provider::{ProviderType, RpcProvider, SharedRpcProviders}; +use actix_web::rt::time::timeout; +use alloy::hex; +use alloy::rpc::types::Block; +use alloy::{eips::BlockNumberOrTag, providers::Provider}; +use alloy_primitives::B256; +use blocksense_config::ReorgConfig; +use blocksense_utils::counter_unbounded_channel::CountedSender; +use std::collections::HashMap; +use std::sync::Arc; +use tokio::sync::Mutex; +use tokio::time::Duration; +use tracing::{debug, error, info, trace, warn}; + +use blocksense_metrics::{inc_metric, inc_vec_metric, metrics::ProviderMetrics}; +use blocksense_utils::await_time; + +// Local helpers from eth_send_utils we need to call +use crate::providers::eth_send_utils::{try_to_sync, BatchOfUpdatesToProcess}; + +pub struct ReorgTracker { + observer_finalized_height: u64, + observed_latest_height: u64, + loop_count: u64, + rpc_timeout: Duration, + net: String, + providers_mutex: SharedRpcProviders, + updates_relayer_send_chan: CountedSender, +} + +impl ReorgTracker { + // Helper to handle reorg once already detected. Finds fork point and prints + // discarded observations, mirroring the existing log messages and structure. + async fn handle_reorg( + &mut self, + rpc_handle: &ProviderType, + provider_mutex: &Arc>, + observed_block_hashes: &HashMap, + observed_latest_height: u64, + ) -> Option { + let net = self.net.clone(); + let observer_finalized_height = self.observer_finalized_height; + let loop_count = self.loop_count; + let rpc_timeout = self.rpc_timeout; + let updates_relayer_send_chan = self.updates_relayer_send_chan.clone(); + + // Build list of observed heights up to observed_latest_height, highest to lowest + let mut observed_heights: Vec = observed_block_hashes + .keys() + .copied() + .filter(|h| *h <= observed_latest_height) + .collect(); + observed_heights.sort_unstable(); + observed_heights.reverse(); + + let fmt_hash = |hash: &B256| format!("0x{}", hex::encode(hash.as_slice())); + + let mut diverged_blocks = Vec::new(); + let mut first_common: Option<(u64, B256)> = None; + + for height in observed_heights { + let stored_hash = match observed_block_hashes.get(&height) { + Some(hash) => *hash, + None => continue, + }; + match timeout( + rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Number(height)), + ) + .await + { + Ok(Ok(Some(chain_block))) => { + let chain_hash = chain_block.header.hash; + if chain_hash == stored_hash { + first_common = Some((height, chain_hash)); + break; + } else { + diverged_blocks.push((height, chain_hash, stored_hash)); + } + } + Ok(Ok(None)) => warn!( + "Block {height} missing while inspecting reorg in network {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ), + Ok(Err(e)) => { + warn!( + "Failed to get block {height} in network {net} while inspecting reorg: {e:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ) + } + Err(_) => { + warn!( + "Timed out getting block {height} in network {net} while inspecting reorg (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ) + } + } + } + + if !diverged_blocks.is_empty() { + let diverged_description: Vec = diverged_blocks + .iter() + .map(|(height, chain_hash, stored_hash)| { + format!( + "height={height}, chain={}, stored={}", + fmt_hash(chain_hash), + fmt_hash(stored_hash), + ) + }) + .collect(); + warn!( + "Diverged blocks observed in network {net}: {} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + diverged_description.join("; ") + ); + } + + if let Some((common_height, common_hash)) = first_common { + info!( + "First common ancestor for reorg in network {net} at height {common_height} with hash {} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + fmt_hash(&common_hash) + ); + let fork_height = common_height + 1; + info!("Fork point for reorg in network {net} is at height {fork_height} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})"); + + // Print all non-finalized updates at or above the fork height + { + let mut provider = provider_mutex.lock().await; + let mut heights: Vec = provider + .inflight + .non_finalized_updates + .keys() + .copied() + .filter(|h| *h >= fork_height) + .collect(); + heights.sort_unstable(); + + if heights.is_empty() { + info!( + "No non_finalized_updates at or above fork height {fork_height} in network {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ); + } else { + for h in heights { + if let Some(batch) = provider.inflight.non_finalized_updates.remove(&h) { + let updates_count = batch.updates.updates.len(); + info!( + "non_finalized_update >= fork: height {h}, batch_block_height {}, updates_count {}, network {} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + batch.updates.block_height, + updates_count, + batch.net, + ); + let msgs_in_queue = updates_relayer_send_chan.len(); + let block_height = batch.updates.block_height; + let provider_metrics = &provider.provider_metrics; + match updates_relayer_send_chan.send(batch) { + Ok(()) => { + debug!("Resent updates to relayer for network {net} and block height {block_height}, messages in queue = {msgs_in_queue} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})"); + inc_metric!(provider_metrics, net, num_transactions_in_queue); + } + Err(e) => { + error!("Error while sending updates to relayer for network {net} and block height {block_height}, messages in queue = {msgs_in_queue}: {e} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})") + } + }; + } + } + } + } + + Some(fork_height) + } else { + warn!( + "Failed to find a common ancestor within stored block hashes for reorg in network {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ); + None + } + } + + pub async fn loop_tracking_for_reorg_in_network(&mut self) { + let net = self.net.clone(); + let providers_mutex = self.providers_mutex.clone(); + tracing::info!("Starting tracker for reorgs in network {net} loop..."); + + // Loop until block generation time is determined + let average_block_generation_time: u64 = loop { + let poll_period = 60 * 1000; + if let Some(t) = self + .calculate_block_generation_time_in_network(net.as_str(), &providers_mutex, 100) + .await + { + break t; + } else { + warn!( + "Could not determine block generation time for network: {net}. Will retry in {poll_period}ms (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ) + } + await_time(poll_period).await; + }; + + loop { + // Sleep between polls + await_time(average_block_generation_time).await; + + self.loop_count += 1; + debug!( + "BEGIN loop_tracking_for_reorg_in_network for {net} loop_count: {loop_count}: observer_finalized_height={observer_finalized_height} observed_latest_height={observed_latest_height}!", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + // Scope the lock on providers so we don't hold it across awaits/sleeps + { + let providers = providers_mutex.read().await; + if let Some(provider_mutex) = providers.get(net.as_str()) { + // Gather data and perform minimal work while holding the provider lock + let mut need_resync_indices = false; + { + let (rpc_handle, observed_block_hashes, provider_metrics) = { + let provider = provider_mutex.lock().await; + + for (k, v) in provider.inflight.non_finalized_updates.iter() { + trace!( + "We have stored the following observations for network {net} inflight[{k}] = {updates:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + updates = v.updates, + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + } + + ( + provider.provider.clone(), + provider.inflight.observed_block_hashes.clone(), + provider.provider_metrics.clone(), + ) + }; + + let latest_block_result = timeout( + self.rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Latest), + ) + .await; + if let Ok(Ok(Some(b))) = latest_block_result { + self.process_new_block( + net.as_str(), + &rpc_handle, + provider_mutex, + &observed_block_hashes, + provider_metrics, + b, + ) + .await; + } + + // 2) Check the on-chain ADFS root; if it differs from our local view, + // set it so subsequent txs use the correct prev-root and flag resync of indices. + { + let mut provider = provider_mutex.lock().await; + if let Some(contract) = provider.get_latest_contract() { + if let Some(contract_address) = contract.address { + match timeout( + self.rpc_timeout, + rpc_handle.get_storage_at( + contract_address, + alloy_primitives::U256::from(0), + ), + ) + .await + { + Ok(Ok(chain_root)) => { + let chain_root_h = crate::providers::provider::HashValue( + chain_root.into(), + ); + let local_frontier_root = + provider.calldata_merkle_tree_frontier.root(); + let tracked_contract_root = + provider.merkle_root_in_contract.clone(); + + let differs_from_local = + local_frontier_root.0 != chain_root_h.0; + + if differs_from_local { + info!( + "Detected state change on-chain for network {net} loop_count {loop_count}. Updating tracked contract root from {tracked_contract_root:?} / local {local_frontier_root:?} to {chain_root_h:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + provider.merkle_root_in_contract = Some(chain_root_h); + need_resync_indices = true; + } + } + Ok(Err(e)) => warn!( + "Failed to read ADFS root from network `{net}`: {e:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + , + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ), + Err(_) => { + warn!( + "Timed out reading ADFS root from network `{net}` (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ) + } + } + } else { + warn!( + "Could not get contract's address for network: {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + } + } else { + warn!( + "No ADFS contract set for network: {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + } + } + + // 1) Observe latest finalized block for logging/visibility + match timeout( + self.rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Finalized), + ) + .await + { + Ok(Ok(res)) => match res { + Some(eth_finalized_block) => { + let mut provider = provider_mutex.lock().await; + + // Prune non-finalized updates up to finalized height + let new_finalized_height = eth_finalized_block.header.inner.number; + if self.observer_finalized_height < new_finalized_height { + info!( + "Last finalized block in network {net} = {eth_finalized_block:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + + self.observer_finalized_height = new_finalized_height; + let removed = provider + .prune_observed_up_to(self.observer_finalized_height); + if removed > 0 { + info!( + "Pruned {removed} non-finalized updates up to finalized height {observer_finalized_height} in `{net}` (observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + } + } + if self.observed_latest_height < self.observer_finalized_height { + warn!( + "Lost track of chain in network {net} beyond a finalized checkpoint: {observer_finalized_height}, last observed block at height: {observed_latest_height} (loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + + self.observed_latest_height = self.observer_finalized_height; + // Insert current hash into provider cache + provider.insert_observed_block_hash( + self.observed_latest_height, + eth_finalized_block.header.hash, + ); + continue; + } + } + None => { + warn!( + "Could not get finalized block in network {net}, got None! (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ) + } + }, + Ok(Err(e)) => { + warn!( + "Could not get finalized block in network {net}: {e:?}! (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ) + } + Err(_) => warn!( + "Timed out getting finalized block in network {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ), + }; + }; + + // 3) If we detected divergence, resync round-buffer indices from chain + if need_resync_indices { + let mut provider = provider_mutex.lock().await; + if let Some(contract) = provider.get_latest_contract() { + if let Some(contract_address) = contract.address { + try_to_sync(net.as_str(), &mut provider, &contract_address, None) + .await; + } + } + } + } else { + info!( + "Terminating reorg tracker for network {net} since it no longer has an active provider! (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + break; + } + debug!( + "END loop_tracking_for_reorg_in_network for {net} loop_count: {loop_count}: observer_finalized_height={observer_finalized_height} observed_latest_height={observed_latest_height}!", + observer_finalized_height = self.observer_finalized_height, + observed_latest_height = self.observed_latest_height, + loop_count = self.loop_count, + ); + } + } + } + + async fn process_new_block( + &mut self, + net: &str, + rpc_handle: &ProviderType, + provider_mutex: &Arc>, + observed_block_hashes: &HashMap, + provider_metrics: Arc>, + b: Block, + ) { + let latest_height = b.header.inner.number; + let mut observed_latest_height = self.observed_latest_height; + let observer_finalized_height = self.observer_finalized_height; + let loop_count = self.loop_count; + let rpc_timeout = self.rpc_timeout; + + if latest_height > observed_latest_height { + if let Some(stored_hash) = observed_block_hashes.get(&observed_latest_height) { + match timeout( + rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Number( + observed_latest_height, + )), + ) + .await + { + Ok(Ok(Some(chain_block))) => { + if chain_block.header.hash != *stored_hash { + warn!( + "Reorg detected in network {net} at observed tip before processing new blocks (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ); + inc_vec_metric!(provider_metrics, observed_reorgs, net); + let _ = self + .handle_reorg( + rpc_handle, + provider_mutex, + observed_block_hashes, + observed_latest_height, + ) + .await; + } + } + Ok(Ok(None)) => warn!( + "Could not get block {observed_latest_height} in network {net} while pre-checking for reorg (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ), + Ok(Err(e)) => warn!( + "Failed to get block {observed_latest_height} in network {net} while pre-checking for reorg: {e:?} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ), + Err(_) => warn!( + "Timed out getting block {observed_latest_height} in network {net} while pre-checking for reorg (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ), + } + } + info!( + "Found new blocks in {net} loop_count = {loop_count} latest_height = {latest_height} old value {} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})", + observed_latest_height + ); + + // Check for reorg via parent mismatch + let first_new_block_height = observed_latest_height + 1; + if let Ok(Ok(Some(first_new_block))) = timeout( + rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Number(first_new_block_height)), + ) + .await + { + if let Some(observed_latest_block_hash) = + observed_block_hashes.get(&observed_latest_height) + { + if first_new_block.header.parent_hash != *observed_latest_block_hash { + warn!( + "Reorg detected in network {net} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height}, loop_count={loop_count})" + ); + inc_vec_metric!(provider_metrics, observed_reorgs, net); + let _ = self + .handle_reorg( + rpc_handle, + provider_mutex, + observed_block_hashes, + observed_latest_height, + ) + .await; + } else { + info!( + "Chain goes on in {net}, loop {loop_count} ... need to add {} new blocks (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})", + latest_height - observed_latest_height + ); + let mut provider = provider_mutex.lock().await; + info!( + "Adding block with height {first_new_block_height} in network {net} loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + provider.insert_observed_block_hash( + first_new_block_height, + first_new_block.header.hash, + ); + for block_height in first_new_block_height + 1..=latest_height { + if let Ok(Ok(Some(new_block))) = timeout( + rpc_timeout, + rpc_handle + .get_block_by_number(BlockNumberOrTag::Number(block_height)), + ) + .await + { + info!( + "Further adding block with height {block_height} in network {net} loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + provider.insert_observed_block_hash( + block_height, + new_block.header.hash, + ); + } else { + warn!( + "Could not get block {block_height} in network {net} loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + } + } + } + } else { + error!( + "No observed block hash for observed_latest_height = {observed_latest_height} in network {net} loop_count {loop_count}! (observer_finalized_height={observer_finalized_height})" + ); + } + } else { + warn!( + "Could not get block {first_new_block_height} in network {net} loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + } + observed_latest_height = latest_height; + } else if latest_height < observed_latest_height { + info!( + "Chain went back in network {net} loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + } else { + info!( + "No new blocks in {net} loop_count {loop_count} latest_height = {latest_height} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + // Even if there are no new blocks, a reorg could have occurred if + // the block at our observed_latest_height now has a different hash. + if let Some(stored_hash) = observed_block_hashes.get(&observed_latest_height) { + let chain_hash = b.header.hash; + if chain_hash != *stored_hash { + warn!( + "Reorg detected in network {net} without new tip advancement loop_count {loop_count} (observer_finalized_height={observer_finalized_height}, observed_latest_height={observed_latest_height})" + ); + inc_vec_metric!(provider_metrics, observed_reorgs, net); + let _ = self + .handle_reorg( + rpc_handle, + provider_mutex, + observed_block_hashes, + observed_latest_height, + ) + .await; + } + } + } + + self.observed_latest_height = observed_latest_height; + } + + pub fn new( + net: String, + config: ReorgConfig, + providers_mutex: SharedRpcProviders, + updates_relayer_send_chan: CountedSender, + ) -> ReorgTracker { + ReorgTracker { + observer_finalized_height: 0, + observed_latest_height: 0, + loop_count: 0, + rpc_timeout: Duration::from_secs(config.rpc_timeout_secs), + net, + providers_mutex, + updates_relayer_send_chan, + } + } + /// Calculates the average block generation time over the last `num_blocks` + /// as `(ts_latest - ts_{latest - num_blocks}) / num_blocks` (in milliseconds). + /// Returns `None` if the calculation cannot be performed. + pub async fn calculate_block_generation_time_in_network( + &self, + net: &str, + providers_mutex: &SharedRpcProviders, + num_blocks: u64, + ) -> Option { + if num_blocks == 0 { + warn!("num_blocks must be >= 1 for average block time in network {net}"); + return None; + } + let providers = providers_mutex.read().await; + let Some(provider_mutex) = providers.get(net) else { + warn!("No active provider found for network {net}"); + return None; + }; + // Clone the RPC handle and read configured timeout without holding the lock across awaits + let rpc_handle = { + let provider = provider_mutex.lock().await; + provider.provider.clone() + }; + + let latest_block = match timeout( + self.rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Latest), + ) + .await + { + Ok(Ok(Some(b))) => b, + Ok(Ok(None)) => { + warn!("Could not get latest block in network {net} (None)"); + return None; + } + Ok(Err(e)) => { + warn!("Error getting latest block in network {net}: {e:?}"); + return None; + } + Err(_) => { + warn!("Timed out getting latest block in network {net}"); + return None; + } + }; + + let latest_height = latest_block.header.inner.number; + if latest_height < num_blocks { + warn!( + "Latest block height {latest_height} < requested lookback {num_blocks} in network {net}" + ); + return None; + } + + let lookback_height = latest_height - num_blocks; + let prev_block = match timeout( + self.rpc_timeout, + rpc_handle.get_block_by_number(BlockNumberOrTag::Number(lookback_height)), + ) + .await + { + Ok(Ok(Some(b))) => b, + Ok(Ok(None)) => { + warn!("Could not get lookback block {lookback_height} in network {net} (None)"); + return None; + } + Ok(Err(e)) => { + warn!("Error getting lookback block {lookback_height} in network {net}: {e:?}"); + return None; + } + Err(_) => { + warn!("Timed out getting lookback block {lookback_height} in network {net}"); + return None; + } + }; + + let latest_ts = latest_block.header.inner.timestamp; + let prev_ts = prev_block.header.inner.timestamp; + if latest_ts < prev_ts { + warn!( + "Latest block timestamp < lookback block timestamp in {net}: {} < {}", + latest_ts, prev_ts + ); + return None; + } + let total_span = latest_ts.saturating_sub(prev_ts); + // Convert to milliseconds before averaging to preserve precision. + let total_span_ms = total_span.saturating_mul(1000); + Some(total_span_ms / num_blocks) + } +} +#[cfg(test)] +mod tests { + use super::*; + use alloy::hex::ToHexExt; + use alloy::node_bindings::Anvil; + use alloy::primitives::Address; + use blocksense_config::AllFeedsConfig; + use blocksense_config::{get_test_config_with_single_provider, test_feed_config}; + use blocksense_data_feeds::feeds_processing::BatchedAggregatesToSend; + use blocksense_registry::config::FeedConfig; + use blocksense_utils::counter_unbounded_channel::counted_unbounded_channel; + use blocksense_utils::EncodedFeedId; + use std::str::FromStr; + use tokio::sync::RwLock; + + use crate::providers::eth_send_utils::BatchOfUpdatesToProcess; + use crate::providers::provider::init_shared_rpc_providers; + + async fn mine_self_txs( + rpc: &ProviderType, + signer: &alloy::signers::local::PrivateKeySigner, + count: u64, + ) -> eyre::Result<()> { + use alloy::network::TransactionBuilder; + use alloy::primitives::U256 as U; + use alloy::rpc::types::TransactionRequest; + + // Fetch chain id and starting nonce + let chain_id = rpc.get_chain_id().await?; + let mut nonce = rpc + .get_transaction_count(signer.address()) + .pending() + .await?; + + for _ in 0..count { + let mut tx = TransactionRequest::default() + .to(signer.address()) + .from(signer.address()) + .with_chain_id(chain_id) + .with_nonce(nonce) + .value(U::from(0u8)); + // Provide minimal gas params since recommended fillers are disabled on the provider + tx = tx.with_gas_limit(21_000); + // Use legacy gas pricing for simplicity + tx = tx.with_gas_price(1_000_000_000u128); + + let pending = rpc.send_transaction(tx).await?; + let _ = pending.get_receipt().await?; + nonce += 1; + } + Ok(()) + } + + async fn mine_varied_txs( + rpc: &ProviderType, + signer: &alloy::signers::local::PrivateKeySigner, + count: u64, + gas_price: u128, + ) -> eyre::Result<()> { + use alloy::network::TransactionBuilder; + use alloy::primitives::{Address as Addr, U256 as U}; + use alloy::rpc::types::TransactionRequest; + + let chain_id = rpc.get_chain_id().await?; + let mut nonce = rpc + .get_transaction_count(signer.address()) + .pending() + .await?; + + for i in 0..count { + // Send to a pseudo-random address derived from i to ensure different block content + let mut bytes = [0u8; 20]; + bytes[0] = (i & 0xff) as u8; + bytes[1] = ((i >> 8) & 0xff) as u8; + let to_addr = Addr::from_slice(&bytes); + + let mut tx = TransactionRequest::default() + .to(to_addr) + .from(signer.address()) + .with_chain_id(chain_id) + .with_nonce(nonce) + .value(U::from(1u8)); + tx = tx.with_gas_limit(50_000); + tx = tx.with_gas_price(gas_price); + + let pending = rpc.send_transaction(tx).await?; + let _ = pending.get_receipt().await?; + nonce += 1; + } + Ok(()) + } + + async fn rpc_call( + url: &str, + method: &str, + params: serde_json::Value, + ) -> eyre::Result { + let client = reqwest::Client::new(); + let body = serde_json::json!({ + "jsonrpc": "2.0", + "id": 1, + "method": method, + "params": params, + }); + let resp = client.post(url).json(&body).send().await?; + let v: serde_json::Value = resp.json().await?; + if let Some(err) = v.get("error") { + eyre::bail!(format!("rpc error on {method}: {err}")); + } + Ok(v["result"].clone()) + } + + async fn anvil_snapshot(url: &str) -> eyre::Result { + let res = match rpc_call(url, "anvil_snapshot", serde_json::json!([])).await { + Ok(v) => v, + Err(_) => rpc_call(url, "evm_snapshot", serde_json::json!([])).await?, + }; + let snap_id = if let Some(s) = res.as_str() { + s.to_string() + } else if let Some(n) = res.as_u64() { + format!("0x{:x}", n) + } else { + eyre::bail!(format!("Unexpected snapshot id result: {res}")); + }; + Ok(snap_id) + } + + async fn anvil_revert(url: &str, snap: &str) -> eyre::Result<()> { + let params = serde_json::json!([snap]); + let res = match rpc_call(url, "anvil_revert", params.clone()).await { + Ok(v) => v, + Err(_) => rpc_call(url, "evm_revert", params).await?, + }; + let ok = res.as_bool().unwrap_or(false); + if !ok { + eyre::bail!(format!("Snapshot revert failed for id {snap}: {res}")); + } + Ok(()) + } + + // End-to-end style test that reproduces a fork and exercises the reorg tracker loop. + // It also verifies that a resync of indices is initiated when the on-chain root differs + // from the local calldata-merkle frontier (without deploying contracts). + #[tokio::test] + async fn test_loop_tracking_reorg_detect_and_resync_indices() { + let _ = tracing_subscriber::fmt().with_test_writer().try_init(); + + // 1) Spin up anvil and build a provider bound to its first funded key + let anvil = Anvil::new().try_spawn().unwrap(); + let net = "ETH1"; + + // Write the anvil key to a temp file that SequencerConfig will use + let signer = anvil.keys()[0].clone(); + let signer_hex = signer.to_bytes().encode_hex(); + let tmp_dir = tempfile::tempdir().unwrap(); + let key_path = tmp_dir.path().join("key"); + std::fs::write(&key_path, signer_hex).unwrap(); + + // Ensure block timestamps increase per mined block (1s interval) + let _ = rpc_call( + anvil.endpoint().as_str(), + "anvil_setBlockTimestampInterval", + serde_json::json!([1]), + ) + .await; + + // Minimal feeds config with a single feed so resync has keys to refresh + let feed = test_feed_config(13, 0); + let feeds_config = AllFeedsConfig { feeds: vec![feed] }; + + let mut cfg = get_test_config_with_single_provider( + net, + key_path.as_path(), + anvil.endpoint().as_str(), + ); + // Ensure we do not auto-load indices during init to keep control in the test + if let Some(p) = cfg.providers.get_mut(net) { + p.should_load_rb_indices = false; + } + + let providers = init_shared_rpc_providers(&cfg, Some("test_reorg_"), &feeds_config).await; + let provider_mutex = providers.read().await.get(net).unwrap().clone(); + + // Inject a dummy ADFS contract address so the reorg loop's root-check and resync paths are exercised + { + let mut provider = provider_mutex.lock().await; + let dummy_contract_address = + Address::from_str("0x1000000000000000000000000000000000000000").unwrap(); + provider.set_contract_address( + blocksense_config::ADFS_CONTRACT_NAME, + &dummy_contract_address, + ); + } + + // 2) Seed local state: non-zero local calldata frontier and initial observed hash for height 0 + { + let mut provider = provider_mutex.lock().await; + + // Make local frontier root non-zero so the loop triggers a resync vs on-chain (zero) root + let dummy_leaf = alloy_primitives::keccak256([0x42u8; 32]); + provider + .calldata_merkle_tree_frontier + .append(crate::providers::provider::HashValue(dummy_leaf)); + + // Pre-set an rb index value to detect refresh to on-chain (expected zero) + let encoded = EncodedFeedId::new(13u128, 0); + provider.rb_indices.insert(encoded, 7); + + // Insert observed hash for genesis (height 0) so the first loop can progress without finalized support + if let Ok(Some(genesis)) = provider + .provider + .get_block_by_number(BlockNumberOrTag::Number(0)) + .await + { + provider.insert_observed_block_hash(0, genesis.header.hash); + } + } + + // 3) Pre-mine 110 blocks by sending self txs to ensure avg block time > 0 + { + let provider = provider_mutex.lock().await; + mine_self_txs(&provider.provider, &provider.signer, 110) + .await + .expect("failed to pre-mine blocks"); + } + + // Query current tip height as T0, then snapshot the chain at T0 + let t0_height = { + let provider = provider_mutex.lock().await; + provider + .provider + .get_block_by_number(BlockNumberOrTag::Latest) + .await + .unwrap() + .unwrap() + .header + .inner + .number + }; + let snapshot_id = anvil_snapshot(anvil.endpoint().as_str()) + .await + .expect("snapshot should succeed"); + + let (feed_updates_send, mut feed_updates_recv) = counted_unbounded_channel(); + + let providers_clone = providers.clone(); + + // 4) Start the reorg tracking loop in background + let loop_handle = tokio::task::Builder::new() + .spawn(async move { + let mut reorg_tracker = ReorgTracker::new( + net.to_string(), + ReorgConfig { + rpc_timeout_secs: 5, + }, + providers_clone, + feed_updates_send, + ); + reorg_tracker.loop_tracking_for_reorg_in_network().await; + }) + .expect("Could not spawn reorg tracker"); + + // Give the loop time to run once and ingest the chain + tokio::time::sleep(std::time::Duration::from_millis(1500)).await; + + // Produce a few more blocks so observed_latest_height advances to T1 > T0 + { + let provider = provider_mutex.lock().await; + mine_self_txs(&provider.provider, &provider.signer, 5) + .await + .expect("failed to add blocks after snapshot"); + } + + // Give it time to incorporate the new tip + tokio::time::sleep(std::time::Duration::from_millis(1200)).await; + + // Capture current latest height (T1) and the observed hash we stored for it + let (t1_height, observed_t1_hash) = { + let provider = provider_mutex.lock().await; + let latest = provider + .provider + .get_block_by_number(BlockNumberOrTag::Latest) + .await + .unwrap() + .unwrap(); + let h = latest.header.inner.number; + let obs = provider.inflight.observed_block_hashes.get(&h).copied(); + (h, obs) + }; + + // Ensure the tracker has recorded observed hash at T1 before inducing the reorg + for _ in 0..20 { + let has_observed_t1 = { + let provider = provider_mutex.lock().await; + provider + .inflight + .observed_block_hashes + .contains_key(&t1_height) + }; + if has_observed_t1 { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + } + + // Inject non-finalized updates: one before the fork (<= T0), and two after the fork (T1-1, T1) + { + let mut provider = provider_mutex.lock().await; + let enc = EncodedFeedId::new(13u128, 0); + let provider_settings = cfg.providers.get(net).unwrap().clone(); + let feeds_map: Arc>> = + Arc::new(RwLock::new(HashMap::new())); + // Prepare three batches with block_height set to the intended heights + let mk_batch = |bh: u64| BatchOfUpdatesToProcess { + net: net.to_string(), + provider: provider_mutex.clone(), + provider_settings: provider_settings.clone(), + updates: BatchedAggregatesToSend { + block_height: 0, // cross-chain height; not used for asserting network key + updates: vec![blocksense_data_feeds::feeds_processing::VotedFeedUpdate { + encoded_feed_id: EncodedFeedId::new(bh as u128, 0), + value: blocksense_feed_registry::types::FeedType::Text(format!( + "marker_for_height_{bh}" + )), + end_slot_timestamp: 0u128, + }], + }, + feeds_config: feeds_map.clone(), + transaction_retry_timeout_secs: 3, + transaction_retries_count_limit: 1, + retry_fee_increment_fraction: 0.1, + }; + // Insert at T0 (pre-fork, should NOT be reintroduced) + provider + .inflight + .insert_non_finalized_update(t0_height, mk_batch(t0_height)); + // Insert at T1 and T1-1 (post-fork, should be reintroduced) + provider + .inflight + .insert_non_finalized_update(t1_height, mk_batch(t1_height)); + if t1_height > 0 { + provider + .inflight + .insert_non_finalized_update(t1_height - 1, mk_batch(t1_height - 1)); + } + // Also ensure the feed index map has our key + provider.rb_indices.insert(enc, 7); + } + + // 5) Revert to the snapshot (T0) and mine a different branch beyond T1 + anvil_revert(anvil.endpoint().as_str(), snapshot_id.as_str()) + .await + .expect("revert should succeed"); + + // Mine enough blocks to surpass T1 + 1 on the new branch + { + let provider = provider_mutex.lock().await; + // First, mine the first 10 blocks with a different gas price so that + // blocks [T0+1..T0+10] differ from the previously observed chain. + mine_varied_txs(&provider.provider, &provider.signer, 10, 2_000_000_000u128) + .await + .expect("failed to mine on new branch (differing blocks)"); + // Then mine additional blocks to move the tip clearly beyond T1, but keep + // total < 64 so finalized (latest-64) remains below ~110 + mine_self_txs(&provider.provider, &provider.signer, 40) + .await + .expect("failed to mine on new branch"); + } + + // Let the loop observe the new branch and trigger resync + tokio::time::sleep(std::time::Duration::from_millis(2000)).await; + + // 6) Assert reorg happened (chain hash at T1 changed vs previously observed) + let chain_t1_hash = { + let provider = provider_mutex.lock().await; + provider + .provider + .get_block_by_number(BlockNumberOrTag::Number(t1_height)) + .await + .unwrap() + .unwrap() + .header + .hash + }; + if let Some(prev_obs) = observed_t1_hash { + assert_ne!( + chain_t1_hash, prev_obs, + "Chain did not diverge at height {} as expected", + t1_height + ); + } + + // 7) Assert indices resync was initiated: merkle_root_in_contract is set from chain (zero) and rb_indices refreshed + { + let provider = provider_mutex.lock().await; + assert!( + provider.get_latest_contract().is_some(), + "Expected ADFS contract to be present in provider config" + ); + } + // Give the loop a bit more time to run the resync path if needed + for _ in 0..10 { + { + let provider = provider_mutex.lock().await; + if provider.merkle_root_in_contract.is_some() { + break; + } + } + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + } + { + let provider = provider_mutex.lock().await; + let root = provider.merkle_root_in_contract.clone(); + assert!( + root.is_some(), + "Expected merkle_root_in_contract to be set by resync loop" + ); + // On empty address storage, root is zero + assert_eq!( + root.unwrap().0, + B256::ZERO, + "Expected on-chain root to be zero" + ); + + let idx = provider + .rb_indices + .get(&EncodedFeedId::new(13u128, 0)) + .copied(); + assert_eq!( + idx, + Some(0), + "Expected rb index for the feed to be refreshed from chain (default 0)" + ); + } + + // Wait until the reorg metric increments to confirm detection + for _ in 0..100 { + let count = { + let provider_metrics = { + let provider = provider_mutex.lock().await; + provider.provider_metrics.clone() + }; + let metrics = provider_metrics.read().await; + metrics.observed_reorgs.with_label_values(&[net]).get() + }; + if count > 0 { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + } + { + let provider_metrics = { + let provider = provider_mutex.lock().await; + provider.provider_metrics.clone() + }; + let metrics = provider_metrics.read().await; + let count = metrics.observed_reorgs.with_label_values(&[net]).get(); + assert!(count > 0, "Expected at least one reorg to be detected"); + } + + // 8) Drain relayer channel and assert only post-fork updates were reintroduced + // Expect exactly the two batches for heights T1-1 and T1, none for T0. + let mut received_heights = Vec::new(); + // Allow generous time for messages to arrive to avoid flakiness + for _ in 0..20 { + if let Ok(Some(batch)) = + tokio::time::timeout(std::time::Duration::from_secs(1), feed_updates_recv.recv()) + .await + { + // Decode the marker we embedded via encoded_feed_id = height + if let Some(vu) = batch.updates.updates.first() { + received_heights.push(vu.encoded_feed_id.get_id()); + } + if received_heights.len() >= 2 { + break; + } + } else { + break; + } + } + received_heights.sort_unstable(); + let expected: Vec = if t1_height > 0 { + vec![(t1_height - 1) as u128, t1_height as u128] + } else { + vec![t1_height as u128] + }; + assert_eq!( + received_heights, expected, + "Expected only post-fork updates to be reintroduced" + ); + assert!( + !received_heights.contains(&(t0_height as u128)), + "Did not expect pre-fork update to be reintroduced" + ); + + // 9) Assert pre-fork update remains in inflight map + { + let provider = provider_mutex.lock().await; + assert!( + provider + .inflight + .non_finalized_updates + .contains_key(&t0_height), + "Expected pre-fork non-finalized update at T0 to remain in inflight.NON-finalized" + ); + } + + // 10) Mine 30 more blocks so finalized surpasses 110 and ensure the pre-fork update is pruned + { + let provider = provider_mutex.lock().await; + mine_self_txs(&provider.provider, &provider.signer, 30) + .await + .expect("failed to mine additional blocks to advance finalization"); + } + // Allow the loop to observe new finalized height and perform pruning + for _ in 0..20 { + { + let provider = provider_mutex.lock().await; + if !provider + .inflight + .non_finalized_updates + .contains_key(&t0_height) + { + break; + } + } + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + } + { + let provider = provider_mutex.lock().await; + assert!( + !provider + .inflight + .non_finalized_updates + .contains_key(&t0_height), + "Expected pre-fork non-finalized update at T0 to be pruned once finalized surpassed it" + ); + } + + // Abort the loop to finish the test + loop_handle.abort(); + } +} diff --git a/apps/sequencer_tests/bin/sequencer_tests.rs b/apps/sequencer_tests/bin/sequencer_tests.rs index 5f9b80019d..4e7808528d 100644 --- a/apps/sequencer_tests/bin/sequencer_tests.rs +++ b/apps/sequencer_tests/bin/sequencer_tests.rs @@ -14,8 +14,8 @@ use blocksense_config::get_sequencer_and_feed_configs; use blocksense_config::SequencerConfig; use blocksense_crypto::JsonSerializableSignature; use blocksense_data_feeds::generate_signature::generate_signature; -use blocksense_feed_registry::registry::await_time; use blocksense_feed_registry::types::{DataFeedPayload, FeedType, PayloadMetaData}; +use blocksense_utils::await_time; use curl::easy::Handler; use curl::easy::WriteError; use curl::easy::{Easy, Easy2}; diff --git a/libs/config/src/lib.rs b/libs/config/src/lib.rs index 7ea9660f60..23702b93a2 100644 --- a/libs/config/src/lib.rs +++ b/libs/config/src/lib.rs @@ -231,6 +231,10 @@ pub struct Provider { #[serde(default)] pub contracts: Vec, + + // Reorg tracking related configuration + #[serde(default)] + pub reorg: ReorgConfig, } fn default_is_enabled() -> bool { @@ -337,6 +341,21 @@ impl Provider { } } +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] +#[serde(default)] +pub struct ReorgConfig { + // Timeout for JSON-RPC requests used during reorg tracking (in seconds) + pub rpc_timeout_secs: u64, +} + +impl Default for ReorgConfig { + fn default() -> Self { + ReorgConfig { + rpc_timeout_secs: 4, + } + } +} + #[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] pub struct Reporter { pub id: u32, @@ -622,6 +641,7 @@ pub fn get_test_config_with_multiple_providers( min_quorum: None, } ], + reorg: ReorgConfig::default(), }, ); } diff --git a/libs/feed_registry/src/registry.rs b/libs/feed_registry/src/registry.rs index efdddd2c1d..8325813b92 100644 --- a/libs/feed_registry/src/registry.rs +++ b/libs/feed_registry/src/registry.rs @@ -6,7 +6,7 @@ use std::{ use crate::types::{DataFeedPayload, FeedMetaData, FeedType, Repeatability, Timestamp}; use blocksense_config::AllFeedsConfig; -use blocksense_utils::{time::current_unix_time, EncodedFeedId}; +use blocksense_utils::{await_time, time::current_unix_time, EncodedFeedId}; use chrono::{DateTime, TimeZone, Utc}; use ringbuf::{ storage::Heap, @@ -15,7 +15,7 @@ use ringbuf::{ }; use serde::{ser::SerializeMap, Deserialize, Serialize, Serializer}; use std::time::UNIX_EPOCH; -use tokio::{sync::RwLock, time}; +use tokio::sync::RwLock; use tracing::{debug, info}; /// Map representing feed_id -> FeedMetaData @@ -396,14 +396,6 @@ impl SlotTimeTracker { } } -pub async fn await_time(time_to_await_ms: u64) { - let time_to_await: Duration = Duration::from_millis(time_to_await_ms); - let mut interval = time::interval(time_to_await); - interval.tick().await; - // The first tick completes immediately. - interval.tick().await; -} - #[cfg(test)] mod tests { use blocksense_utils::time::current_unix_time; diff --git a/libs/metrics/src/metrics.rs b/libs/metrics/src/metrics.rs index 95dba15b17..5f38da421b 100644 --- a/libs/metrics/src/metrics.rs +++ b/libs/metrics/src/metrics.rs @@ -172,6 +172,7 @@ pub struct ProviderMetrics { pub total_mismatched_gnosis_safe_nonce: IntCounterVec, pub num_transactions_in_queue: IntGaugeVec, pub is_enabled: IntGaugeVec, + pub observed_reorgs: IntCounterVec, } impl ProviderMetrics { @@ -278,6 +279,11 @@ impl ProviderMetrics { "Whether the network is currently enabled or not", &["Network"] )?, + observed_reorgs: register_int_counter_vec!( + format!("{}observed_reorgs", prefix), + "Total number of observed chain reorganizations for the network", + &["Network"] + )?, }) } } diff --git a/libs/utils/src/lib.rs b/libs/utils/src/lib.rs index c5f5a7a568..9a7cbcd295 100644 --- a/libs/utils/src/lib.rs +++ b/libs/utils/src/lib.rs @@ -268,6 +268,8 @@ use std::{ use anyhow::{anyhow, Context, Result}; +use std::time::Duration; + pub fn get_env_var(key: &str) -> Result where T: FromStr, @@ -331,6 +333,14 @@ pub fn get_config_file_path(base_path_from_env: &str, config_file_name: &str) -> config_file_path.join(config_file_name) } +pub async fn await_time(time_to_await_ms: u64) { + let time_to_await: Duration = Duration::from_millis(time_to_await_ms); + let mut interval = tokio::time::interval(time_to_await); + interval.tick().await; + // The first tick completes immediately. + interval.tick().await; +} + #[cfg(test)] mod tests { use super::*;