diff --git a/Cargo.lock b/Cargo.lock index 6560bda0..09ed9780 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5828,6 +5828,7 @@ dependencies = [ "tokio", "tokio-util", "tracing", + "tracing-subscriber 0.3.23", ] [[package]] diff --git a/bin/stateless-validator/Cargo.toml b/bin/stateless-validator/Cargo.toml index fabb866a..86f2a87d 100644 --- a/bin/stateless-validator/Cargo.toml +++ b/bin/stateless-validator/Cargo.toml @@ -48,6 +48,7 @@ tracing.workspace = true jsonrpsee.workspace = true jsonrpsee-types.workspace = true tempfile.workspace = true +tracing-subscriber.workspace = true # stateless stateless-test-utils = { path = "../../crates/stateless-test-utils", features = ["mock-rpc"] } diff --git a/bin/stateless-validator/src/runner.rs b/bin/stateless-validator/src/runner.rs index f3873593..cc8094f3 100644 --- a/bin/stateless-validator/src/runner.rs +++ b/bin/stateless-validator/src/runner.rs @@ -12,7 +12,8 @@ use alloy_primitives::B256; use eyre::Result; use stateless_common::RpcClient; use stateless_core::{ - BisectResolver, ChainStore, PipelineConfig, chain_spec::ChainSpec, pipeline::run_pipeline, + BisectResolver, ChainStore, PipelineConfig, chain_spec::ChainSpec, db::BlockMeta, + pipeline::run_pipeline, }; use stateless_db::ContractCache; use tokio::{signal, task}; @@ -34,9 +35,8 @@ const FINAL_REPORT_RETRY_DELAY: Duration = Duration::from_secs(1); /// /// Cleanly drains on SIGINT/SIGTERM and returns either the pipeline result or `Ok(())` /// on signal. Exception: on a fixed-range run (`--end-block`), a final validation report -/// that cannot land (or a detected validation gap) fails the run — a slice has no later -/// restart to re-report, so exiting 0 would let an orchestrator record the slice as -/// complete while upstream never saw the tip. +/// that cannot land fails the run — a slice has no later restart to re-report, so exiting 0 +/// would let an orchestrator record the slice as complete while upstream never saw the tip. pub async fn run_with_signals( client: Arc, r2_witness: Option>, @@ -136,29 +136,19 @@ pub async fn run_with_signals( if report_validation { let mut flush_failure: Option = None; for attempt in 1..=FINAL_REPORT_ATTEMPTS { - match report_range_once(&client, &validator_db, &last_reported).await { - Ok(true) => break, - Ok(false) if attempt < FINAL_REPORT_ATTEMPTS => { - warn!(attempt, "Final validation report failed, retrying"); - tokio::time::sleep(FINAL_REPORT_RETRY_DELAY).await; - } - Ok(false) => { - error!( - attempts = FINAL_REPORT_ATTEMPTS, - "Final validation report failed; the validated tail may be unreported \ - upstream" - ); - flush_failure = Some(eyre::eyre!( - "final validation report failed after {FINAL_REPORT_ATTEMPTS} attempts; \ - the validated tail may be unreported upstream" - )); - } - // A detected validation gap is deterministic — retrying cannot resolve it. - Err(e) => { - error!(error = %e, "Final validation report failed with a non-retryable error"); - flush_failure = Some(e); - break; - } + if report_range_once(&client, &validator_db, &last_reported).await { + break; + } else if attempt < FINAL_REPORT_ATTEMPTS { + warn!(attempt, "Final validation report failed, retrying"); + tokio::time::sleep(FINAL_REPORT_RETRY_DELAY).await; + } else { + error!( + attempts = FINAL_REPORT_ATTEMPTS, + "Final validation report failed; the validated tail may be unreported upstream" + ); + flush_failure = Some(eyre::eyre!( + "final validation report failed after {FINAL_REPORT_ATTEMPTS} attempts; the validated tail may be unreported upstream" + )); } } // A chain-following run re-reports anchor→tip on its next start, so an unreported tail @@ -213,7 +203,9 @@ async fn validation_reporter( } } - report_range_once(&client, &validator_db, &last_reported_block).await?; + if !report_range_once(&client, &validator_db, &last_reported_block).await { + warn!("Validation report was not accepted; will retry on next interval"); + } } } @@ -221,20 +213,19 @@ async fn validation_reporter( /// differs from `last_reported_block` (updated on an accepted report; a tip that regressed /// after a reorg rollback is deliberately re-reported). /// -/// Returns `Ok(true)` when the round settled (report accepted, or nothing to report) and -/// `Ok(false)` when the attempt failed in a way a retry could resolve (logged here). The only -/// `Err` is a detected validation gap, which is fatal to the reporter. +/// Returns `true` when the round settled (report accepted, or nothing to report) and `false` +/// when the attempt failed in a way a retry could resolve (logged here). async fn report_range_once( client: &RpcClient, validator_db: &ValidatorDB, last_reported_block: &AtomicU64, -) -> Result { +) -> bool { let (anchor, tip) = match (validator_db.get_anchor(), validator_db.get_canonical_tip()) { (Ok(Some(a)), Ok(Some(t))) => (a, t), - (Ok(None), _) | (_, Ok(None)) => return Ok(true), + (Ok(None), _) | (_, Ok(None)) => return true, (Err(e), _) | (_, Err(e)) => { warn!(error = %e, "Failed to read anchor/tip, retrying"); - return Ok(false); + return false; } }; @@ -242,7 +233,7 @@ async fn report_range_once( // the reporter is the only writer while it runs, and the final flush reads it only after // joining the reporter task. if tip.block_number == last_reported_block.load(Ordering::Relaxed) { - return Ok(true); + return true; } let result = client @@ -262,25 +253,375 @@ async fn report_range_once( "Reported blocks" ); last_reported_block.store(tip.block_number, Ordering::Relaxed); - Ok(true) + true } Ok(response) => { - if response.last_validated_block.0 < anchor.block_number { - return Err(eyre::eyre!( - "Validation gap detected: upstream at block {}, but local chain starts at {}", - response.last_validated_block.0, - anchor.block_number - )); + let upstream_number = response.last_validated_block.0.to::(); + let upstream_hash = response.last_validated_block.1; + warn!( + accepted = response.accepted, + local_anchor = anchor.block_number, + local_anchor_hash = %anchor.block_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + "Report rejected" + ); + + if upstream_number < anchor.block_number { + warn!( + local_anchor = anchor.block_number, + local_anchor_hash = %anchor.block_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + "validation gap detected; checking whether local chain can cover upstream last_validated" + ); + + return retry_report_from_upstream_pointer( + client, + validator_db, + last_reported_block, + &tip, + upstream_number, + upstream_hash, + ) + .await; } + + false + } + Err(e) => { + error!(error = %e, "Failed to report blocks"); + false + } + } +} + +/// Tries to heal a rejected anchor→tip report by covering the receiver's current pointer. +/// +/// `mega_setValidatedBlocks` accepts a range that covers the receiver's current validated block. +/// If the receiver is behind our current anchor but its pointer is still retained in the local +/// canonical-chain window, report that retained pointer as `first_block` so the receiver can +/// catch up without waiting for a process restart. If the pointer is absent or has a different +/// hash, keep the gap unresolved: this standalone validator must not claim to cover pre-anchor +/// heights it did not execute. +async fn retry_report_from_upstream_pointer( + client: &RpcClient, + validator_db: &ValidatorDB, + last_reported_block: &AtomicU64, + tip: &BlockMeta, + upstream_number: u64, + upstream_hash: B256, +) -> bool { + match validator_db.get_block_hash(upstream_number) { + Ok(Some(local_hash)) + if local_hash == alloy_primitives::BlockHash::from(upstream_hash.0) => + { + warn!( + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + "validation gap can be covered by local chain; retrying report from upstream last_validated" + ); + retry_covering_report(client, last_reported_block, tip, upstream_number, upstream_hash) + .await + } + Ok(Some(local_hash)) => { error!( - upstream_block = ?response.last_validated_block, - "Report rejected" + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + local_hash = %local_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + "validation gap cannot be covered: upstream last_validated hash mismatches local canonical chain" ); - Ok(false) + false + } + Ok(None) => { + error!( + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + local_hash = "missing", + "validation gap cannot be covered: upstream last_validated block is not retained locally" + ); + false + } + Err(e) => { + warn!( + error = %e, + upstream_last_validated = upstream_number, + upstream_last_validated_hash = %upstream_hash, + "Failed to read upstream last_validated hash from local chain, retrying" + ); + false + } + } +} + +async fn retry_covering_report( + client: &RpcClient, + last_reported_block: &AtomicU64, + tip: &BlockMeta, + upstream_number: u64, + upstream_hash: B256, +) -> bool { + match client + .set_validated_blocks( + (upstream_number, upstream_hash), + (tip.block_number, B256::from(tip.block_hash.0)), + ) + .await + { + Ok(response) if response.accepted => { + debug!( + start = upstream_number, + start_hash = %upstream_hash, + tip = tip.block_number, + tip_hash = %tip.block_hash, + "Reported blocks after validation gap resync" + ); + last_reported_block.store(tip.block_number, Ordering::Relaxed); + true + } + Ok(response) => { + warn!( + accepted = response.accepted, + upstream_last_validated = response.last_validated_block.0.to::(), + upstream_last_validated_hash = %response.last_validated_block.1, + attempted_start = upstream_number, + attempted_start_hash = %upstream_hash, + local_tip = tip.block_number, + local_tip_hash = %tip.block_hash, + "Validation gap resync report was rejected" + ); + false } Err(e) => { error!(error = %e, "Failed to report blocks"); - Ok(false) + false } } } + +#[cfg(test)] +mod tests { + use std::{ + collections::VecDeque, + future::Future, + io, + sync::{Arc, Mutex, atomic::AtomicU64}, + }; + + use jsonrpsee_types::ErrorObjectOwned; + use stateless_common::RpcClientConfig; + use stateless_core::ChainStore; + use stateless_test_utils::mock_rpc::serve; + use tracing::Level; + use tracing_subscriber::fmt::MakeWriter; + + use super::*; + use crate::test_support::make_block_meta; + + #[derive(Clone)] + struct CapturedLogs(Arc>>); + + impl CapturedLogs { + fn new() -> Self { + Self(Arc::default()) + } + + fn contents(&self) -> String { + String::from_utf8(self.0.lock().unwrap().clone()).unwrap() + } + } + + impl<'a> MakeWriter<'a> for CapturedLogs { + type Writer = CapturedLogWriter; + + fn make_writer(&'a self) -> Self::Writer { + CapturedLogWriter(Arc::clone(&self.0)) + } + } + + struct CapturedLogWriter(Arc>>); + + impl io::Write for CapturedLogWriter { + fn write(&mut self, buf: &[u8]) -> io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + async fn with_log_capture(future: impl Future) -> (T, String) { + let logs = CapturedLogs::new(); + let subscriber = tracing_subscriber::fmt() + .with_writer(logs.clone()) + .with_ansi(false) + .with_max_level(Level::WARN) + .finish(); + let _guard = tracing::subscriber::set_default(subscriber); + let output = future.await; + (output, logs.contents()) + } + + #[derive(Clone)] + struct ReportResponse { + accepted: bool, + last_validated_block: (u64, B256), + } + + type ReportCall = ((u64, B256), (u64, B256)); + + #[derive(Default)] + struct ReportServerState { + responses: Mutex>, + calls: Mutex>, + } + + async fn setup_report_client( + responses: Vec, + ) -> (Arc, Arc, jsonrpsee::server::ServerHandle) { + let state = Arc::new(ReportServerState { + responses: Mutex::new(responses.into()), + calls: Mutex::default(), + }); + let (handle, url) = serve(Arc::clone(&state), |module| { + module + .register_method("mega_setValidatedBlocks", |params, ctx, _| { + let (first, last): ((u64, B256), (u64, B256)) = params.parse().unwrap(); + ctx.calls.lock().unwrap().push((first, last)); + let response = ctx + .responses + .lock() + .unwrap() + .pop_front() + .expect("test must script every report response"); + Ok::(serde_json::json!({ + "accepted": response.accepted, + "lastValidatedBlock": [ + response.last_validated_block.0, + response.last_validated_block.1, + ], + })) + }) + .unwrap(); + }) + .await; + let client = Arc::new( + RpcClient::new_with_config( + &[url.as_str()], + &[url.as_str()], + RpcClientConfig::validator(), + Some(url.as_str()), + ) + .unwrap(), + ); + (client, state, handle) + } + + fn setup_report_db() -> (tempfile::TempDir, ValidatorDB) { + let dir = tempfile::tempdir().unwrap(); + let db = ValidatorDB::new(dir.path().join("validator.redb")).unwrap(); + db.reset_to_anchor(&make_block_meta(70)).unwrap(); + + let blocks: Vec<_> = (71..=80).map(make_block_meta).collect(); + db.advance_chain(&blocks).unwrap(); + (dir, db) + } + + async fn wait_for_report_calls(state: &ReportServerState, expected: usize) { + tokio::time::timeout(Duration::from_secs(2), async { + loop { + if state.calls.lock().unwrap().len() >= expected { + return; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("timed out waiting for report calls"); + } + + fn rejected_at(number: u64, hash: B256) -> ReportResponse { + ReportResponse { accepted: false, last_validated_block: (number, hash) } + } + + #[tokio::test(flavor = "current_thread")] + async fn rejected_gap_below_anchor_does_not_push_covering_range_and_reporter_stays_alive() { + let (_dir, db) = setup_report_db(); + let tip = make_block_meta(80); + let upstream_hash = B256::ZERO; + let responses = vec![rejected_at(0, upstream_hash); 4]; + let (client, state, handle) = setup_report_client(responses).await; + let shutdown = CancellationToken::new(); + let last_reported = Arc::new(AtomicU64::new(0)); + let reporter = tokio::spawn(validation_reporter( + client, + Arc::new(db), + Duration::from_millis(50), + shutdown.clone(), + Arc::clone(&last_reported), + )); + + wait_for_report_calls(&state, 1).await; + assert_eq!(last_reported.load(Ordering::Relaxed), 0); + assert!(!reporter.is_finished(), "reporter must keep looping after a validation gap"); + shutdown.cancel(); + reporter.await.unwrap().unwrap(); + + let calls = state.calls.lock().unwrap().clone(); + assert!(!calls.is_empty()); + assert!( + calls.iter().all(|call| call.0.0 == 70 && call.1.0 == 80), + "must not report a covering range below the local anchor: {calls:?}" + ); + assert_eq!(calls[0].1, (80, B256::from(tip.block_hash.0))); + handle.stop().unwrap(); + } + + #[tokio::test(flavor = "current_thread")] + async fn rejected_gap_missing_local_hash_does_not_push_covering_range_and_is_not_fatal() { + let upstream_hash = B256::from([99u8; 32]); + let (_dir, db) = setup_report_db(); + let (client, state, handle) = + setup_report_client(vec![rejected_at(15, upstream_hash)]).await; + let last_reported = AtomicU64::new(0); + + assert!(!report_range_once(&client, &db, &last_reported).await); + assert_eq!(last_reported.load(Ordering::Relaxed), 0); + + let calls = state.calls.lock().unwrap().clone(); + assert_eq!(calls.len(), 1, "must not report a range from a missing local hash"); + assert_eq!(calls[0].0.0, 70); + handle.stop().unwrap(); + } + + #[tokio::test(flavor = "current_thread")] + async fn rejected_gap_logs_uncovered_validation_gap() { + let upstream_hash = B256::from([99u8; 32]); + let (_dir, db) = setup_report_db(); + let (client, _state, handle) = + setup_report_client(vec![rejected_at(15, upstream_hash)]).await; + + let (settled, captured) = + with_log_capture(report_range_once(&client, &db, &AtomicU64::new(0))).await; + assert!(!settled); + assert!( + captured.contains( + "validation gap cannot be covered: upstream last_validated block is not retained locally" + ), + "gap path must emit the uncovered-gap log, got: {captured}" + ); + handle.stop().unwrap(); + } +}