From 7dc7c6fd8ea9489ee290fa114af72cc19b36bd1f Mon Sep 17 00:00:00 2001 From: Tim Denisenko Date: Thu, 2 Jul 2026 21:28:19 +0700 Subject: [PATCH 1/3] fix: recover stale checkpoint restarts --- ROADMAP.md | 77 ++--- crates/logex-node/src/runtime.rs | 227 ++++++++++-- crates/logex-sync/src/engine/anchored.rs | 417 +++++++++++++++++++++++ 3 files changed, 658 insertions(+), 63 deletions(-) diff --git a/ROADMAP.md b/ROADMAP.md index 624f859a..73314923 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -4,66 +4,63 @@ LogEx starts from a recent consensus checkpoint, tracks the live head, reverse-syncs execution history toward genesis, stores compressed verified logs, and serves the dashboard, SQL query API, JSON-RPC, gRPC, and live ERC20 transfer subscriptions. -Current branch after this completed work: `master`. +Current branch: `fix/auto-refresh-stale-checkpoint`. -The main dashboard now reports Execution Layer P2P download throughput instead of an ETA tile. The metric is measured from successful decoded block-body and receipt payloads returned by peers, then shown as Mbps in the UI. +The Mac mini is running this branch against the saved full-sync data directory. Public dashboard forwarding through `157.245.195.72:18683` is active, and the client is bridging the stale restart gap from the saved execution head to the refreshed consensus checkpoint. ## Completed Since Last Run -- Added execution P2P decoded payload byte-rate tracking for successful body and receipt responses. -- Exposed `p2p_download_bytes_per_sec` and `p2p_downloaded_payload_bytes` under `execution_network` in `/status`. -- Replaced the main dashboard ETA tile with a `P2P download` Mbps tile. -- Added the same throughput and cumulative decoded payload metric to advanced dashboard details. -- Removed obsolete dashboard ETA formatting helpers that no longer had call sites. +- Added automatic startup recovery for stale persisted consensus state and stale local execution progress. +- Added a checkpoint-gap bridge that validates the EL parent chain up to a fresh CL checkpoint anchor before ingesting the missing logs. +- Kept the existing EL/log storage and known peer cache intact during checkpoint refresh. +- Verified the remote dashboard path through the VPS and restarted the native Mac mini client with the fix. ## Remaining TODOs -No remaining TODOs for the dashboard P2P download-rate task. +No remaining TODOs for the stale-checkpoint restart recovery task. ## Design Decisions -- Use decoded P2P payload bytes rather than raw socket bytes. - - Why: Reth's current network handle does not expose per-peer transport byte counters, while decoded body/receipt responses are available at the request accounting layer. - - Alternatives considered: estimate throughput from block/log progress or add invasive transport hooks. Progress-derived rates are less direct, and transport hooks would be much higher risk for this UI change. - - Tradeoff: The metric represents useful Ethereum payload throughput, not encrypted TCP overhead or protocol framing bytes. +- Archive stale `cl/consensus_state.json` instead of deleting it. + - Why: The client can recover automatically while preserving a diagnostic copy of the old trusted state. + - Alternatives considered: fail startup and require manual data-dir surgery. That was operationally fragile and caused the dashboard outage. + - Tradeoff: The data directory may retain a small archived consensus snapshot after recovery. -- Expose bytes/sec in the API and format Mbps in the browser. - - Why: Integer bytes/sec keeps the Rust status type simple and precise while preserving the dashboard wording the user requested. - - Alternatives considered: expose floating-point Mbps directly. That would require weakening the existing `ExecutionNetworkStatus` equality semantics. - - Tradeoff: API consumers convert to their preferred network unit. - -- Decay stale P2P download rate to zero after a short freshness window. - - Why: The dashboard should not show an old high throughput number when no body or receipt payloads have arrived recently. - - Alternatives considered: keep the last EWMA indefinitely. That would be misleading during stalls or idle periods. - - Tradeoff: Very bursty request windows may briefly show zero between payload samples. +- Bridge long-offline gaps with EL parent-chain validation ending at a fresh CL checkpoint anchor. + - Why: CL historical sync is intentionally not required; the fresh checkpoint execution hash is enough to validate the canonical EL ancestor chain back to the saved head. + - Alternatives considered: require a fresh data directory, or reintroduce CL historical sync. Both add unnecessary operational cost for this restart case. + - Tradeoff: The bridge is correctness-first and less optimized than the normal historical pipeline. ## Challenges and Resolutions -- Challenge: The previous ETA field was computed from sync progress and did not reflect actual P2P download activity. - - Resolution: Added request-layer byte accounting for body and receipt responses and wired it into the dashboard. - - Remaining: No known issue for this task. +- Challenge: A full synced data directory failed to restart after being offline because the persisted consensus state and local execution head were older than the recent checkpoint window. + - Resolution: Startup now resolves a fresh checkpoint, archives stale CL state, and resumes using the existing EL/log storage. + - Remaining: None known for correctness. + +- Challenge: The first recovery attempt hit the consensus reorg guard because the fresh CL anchor was ahead of the persisted recent-header window. + - Resolution: Future-only anchors are classified as restart gaps, then bridged by fetching and validating the EL header chain to the CL anchor. + - Remaining: Gap-bridge throughput can be optimized later if this path becomes common. ## Dead Code and Obsolescence Cleanup -- Inspected the dashboard formatting helpers after replacing the ETA tile. -- Removed unused ETA and completion-time formatting functions from the dashboard script. -- Checked request accounting tuple usage and converted it to typed structs so payload bytes are carried consistently. -- No additional safe removal was identified. +- Inspected the stale startup guards and consensus reorg classifier. +- No obsolete code was safely removable in this pass; the new bridge reuses existing validation, peer, and storage primitives. +- Removed no files. ## Git Workflow -- Current branch after this completed work: `master`. -- Task branches created: - - `feature/dashboard-p2p-download-rate` - - `docs/finalize-p2p-download-roadmap` -- Commits made during this run: - - `9dc8b5e5 feat: show p2p download throughput` -- Pull request status: - - PR #99 was created, checks passed, and merged. - - PR #100 was created for the post-merge roadmap correction, checks passed, and merged. -- Merge status: dashboard P2P download-rate work is merged into `master`. -- Blockers: none known. +- Current branch: `fix/auto-refresh-stale-checkpoint`. +- New branch created from `master`. +- Commits made during this run: pending. +- Pull request status: pending after commit/push. +- Merge status: not merged yet. +- Validation run: + - `cargo test -p logex-node archived_consensus_state` + - `cargo test -p logex-node restart_guard` + - `cargo test -p logex-sync locate_consensus_reorg` + - `cargo clippy -p logex-node -p logex-sync --all-targets -- -D warnings` ## Known Issues or Risks -- `p2p_download_bytes_per_sec` measures decoded Ethereum body/receipt payloads from successful responses, not raw encrypted network interface throughput. +- The checkpoint-gap bridge prioritizes trustless recovery over throughput. Normal historical sync performance is unchanged. +- The local Codex sandbox cannot directly curl the public VPS dashboard, but the VPS can reach `10.66.0.2:18683` and packet capture showed public TCP/18683 traffic being forwarded and answered. diff --git a/crates/logex-node/src/runtime.rs b/crates/logex-node/src/runtime.rs index f3eda5fc..ea6811a5 100644 --- a/crates/logex-node/src/runtime.rs +++ b/crates/logex-node/src/runtime.rs @@ -221,16 +221,11 @@ pub async fn run_sync(options: RunSyncOptions) { let discovery_secret_file = discovery_secret_path(&data_dir); let known_peers_file = known_peers_path(&data_dir); let consensus_state_exists = data_dir.join("cl").join("consensus_state.json").exists(); - let checkpoint = if consensus_state_exists && checkpoint.is_none() { - checkpoint + let checkpoint_request = checkpoint; + let mut checkpoint = if consensus_state_exists && checkpoint_request.is_none() { + None } else { - match resolve_checkpoint(checkpoint, checkpoint_sync_url.as_deref()).await { - Ok(checkpoint) => checkpoint, - Err(error) => { - tracing::error!(%error, "failed to resolve weak-subjectivity checkpoint"); - std::process::exit(1); - } - } + resolve_checkpoint_or_exit(checkpoint_request, checkpoint_sync_url.as_deref()).await }; let mut storage = match PartitionManager::open(pm_config) { @@ -269,7 +264,9 @@ pub async fn run_sync(options: RunSyncOptions) { "storage ready" ); - let consensus = match maybe_open_consensus_store(&data_dir, &storage, checkpoint.as_deref()) { + let mut consensus_refreshed = false; + let mut consensus = match maybe_open_consensus_store(&data_dir, &storage, checkpoint.as_deref()) + { Ok(store) => store.map(Arc::new), Err(ConsensusStateError::MissingCheckpoint) => { tracing::error!( @@ -278,32 +275,107 @@ pub async fn run_sync(options: RunSyncOptions) { ); std::process::exit(1); } + Err(ConsensusStateError::StaleWeakSubjectivityCheckpoint { .. }) => { + tracing::warn!( + data_dir = %data_dir.display(), + "persisted consensus state is outside the weak-subjectivity window; refreshing from a recent checkpoint" + ); + match refresh_consensus_state_from_recent_checkpoint( + &data_dir, + &mut checkpoint, + checkpoint_sync_url.as_deref(), + "weak-subjectivity-stale", + ) + .await + { + Ok(store) => { + consensus_refreshed = true; + Some(Arc::new(store)) + } + Err(error) => { + tracing::error!(%error, "failed to refresh stale consensus state"); + std::process::exit(1); + } + } + } Err(error) => { tracing::error!(%error, "failed to initialize consensus state"); std::process::exit(1); } }; - if let Some(consensus) = consensus.as_ref() { - if let Some(staleness) = recent_consensus_state_staleness(consensus) { - tracing::error!( - trusted_slot = staleness.trusted_slot, - trusted_epoch = staleness.trusted_epoch, - current_epoch = staleness.current_epoch, - max_epochs = staleness.max_epochs, - "persisted consensus state is too stale; start from a fresh recent checkpoint" - ); - std::process::exit(1); + if let Some(store) = consensus.as_ref() + && let Some(staleness) = recent_consensus_state_staleness(store) + { + tracing::warn!( + trusted_slot = staleness.trusted_slot, + trusted_epoch = staleness.trusted_epoch, + current_epoch = staleness.current_epoch, + max_epochs = staleness.max_epochs, + "persisted consensus state is older than the recent checkpoint window; refreshing from checkpoint-sync source" + ); + match refresh_consensus_state_from_recent_checkpoint( + &data_dir, + &mut checkpoint, + checkpoint_sync_url.as_deref(), + "recent-checkpoint-stale", + ) + .await + { + Ok(store) => { + consensus = Some(Arc::new(store)); + consensus_refreshed = true; + } + Err(error) => { + tracing::error!(%error, "failed to refresh stale consensus state"); + std::process::exit(1); + } + } + } + + if let Some(store) = consensus.as_ref() + && let Some(staleness) = local_execution_progress_staleness(sync_head, store) + && !consensus_refreshed + { + tracing::warn!( + block_number = staleness.block_number, + timestamp = staleness.timestamp, + age_secs = staleness.age_secs, + max_age_secs = staleness.max_age_secs, + "local execution progress is older than the recent checkpoint window; refreshing consensus checkpoint before resuming" + ); + match refresh_consensus_state_from_recent_checkpoint( + &data_dir, + &mut checkpoint, + checkpoint_sync_url.as_deref(), + "execution-progress-stale", + ) + .await + { + Ok(store) => { + consensus = Some(Arc::new(store)); + consensus_refreshed = true; + } + Err(error) => { + tracing::error!( + %error, + "failed to refresh checkpoint for stale local execution progress" + ); + std::process::exit(1); + } } + } + + if let Some(consensus) = consensus.as_ref() { if let Some(staleness) = local_execution_progress_staleness(sync_head, consensus) { - tracing::error!( + tracing::warn!( block_number = staleness.block_number, timestamp = staleness.timestamp, age_secs = staleness.age_secs, max_age_secs = staleness.max_age_secs, - "local execution progress is too stale; start from a fresh recent checkpoint in a fresh data directory" + checkpoint_refreshed = consensus_refreshed, + "local execution progress is stale but startup has a recent checkpoint; CL anchors will bridge the gap before live EL sync advances" ); - std::process::exit(1); } let checkpoint = consensus.checkpoint(); let mut anchors = consensus.chain_anchors(); @@ -1250,6 +1322,91 @@ fn maybe_open_consensus_store( Ok(None) } +async fn resolve_checkpoint_or_exit( + checkpoint: Option, + checkpoint_sync_url: Option<&str>, +) -> Option { + match resolve_checkpoint(checkpoint, checkpoint_sync_url).await { + Ok(checkpoint) => checkpoint, + Err(error) => { + tracing::error!(%error, "failed to resolve weak-subjectivity checkpoint"); + std::process::exit(1); + } + } +} + +async fn refresh_consensus_state_from_recent_checkpoint( + data_dir: &Path, + checkpoint: &mut Option, + checkpoint_sync_url: Option<&str>, + reason: &str, +) -> Result { + let checkpoint = resolve_checkpoint_for_recovery(checkpoint, checkpoint_sync_url).await?; + if let Some(archive_path) = + archive_consensus_state(data_dir, reason).map_err(|error| error.to_string())? + { + tracing::warn!( + path = %archive_path.display(), + "archived stale consensus state before checkpoint refresh" + ); + } + ConsensusStore::open(data_dir, Some(&checkpoint)).map_err(|error| error.to_string()) +} + +async fn resolve_checkpoint_for_recovery( + checkpoint: &mut Option, + checkpoint_sync_url: Option<&str>, +) -> Result { + if let Some(checkpoint) = checkpoint.clone() { + return Ok(checkpoint); + } + + let resolved = resolve_checkpoint(None, checkpoint_sync_url) + .await + .map_err(|error| error.to_string())? + .ok_or_else(|| "checkpoint-sync source did not return a checkpoint".to_owned())?; + *checkpoint = Some(resolved.clone()); + Ok(resolved) +} + +fn archive_consensus_state( + data_dir: &Path, + reason: &str, +) -> Result, ConsensusStateError> { + let path = data_dir.join("cl").join("consensus_state.json"); + if !path.exists() { + return Ok(None); + } + + let parent = path.parent().unwrap_or(data_dir); + let timestamp = current_unix_timestamp(); + for suffix in 0..1000 { + let archive_path = if suffix == 0 { + parent.join(format!("consensus_state.{reason}.{timestamp}.json")) + } else { + parent.join(format!( + "consensus_state.{reason}.{timestamp}.{suffix}.json" + )) + }; + if archive_path.exists() { + continue; + } + fs::rename(&path, &archive_path).map_err(|source| ConsensusStateError::PersistState { + path: archive_path.clone(), + source, + })?; + return Ok(Some(archive_path)); + } + + Err(ConsensusStateError::PersistState { + path, + source: std::io::Error::new( + std::io::ErrorKind::AlreadyExists, + "could not choose a unique consensus state archive path", + ), + }) +} + #[derive(Debug, Clone, Copy)] struct RecentConsensusStateStaleness { trusted_slot: u64, @@ -1999,6 +2156,30 @@ mod tests { assert!(staleness.age_secs > staleness.max_age_secs); } + #[test] + fn archived_consensus_state_allows_fresh_checkpoint_reopen() { + let temp = tempfile::tempdir().unwrap(); + let old_checkpoint = checkpoint_at_slot(recent_checkpoint_slot(1)); + let _old_store = ConsensusStore::open(temp.path(), Some(&old_checkpoint)).unwrap(); + let state_path = temp.path().join("cl").join("consensus_state.json"); + + let archive_path = archive_consensus_state(temp.path(), "test-refresh") + .unwrap() + .unwrap(); + + assert!(!state_path.exists()); + assert!(archive_path.exists()); + + let new_slot = recent_checkpoint_slot(0); + let new_root = B256::repeat_byte(0x43); + let new_checkpoint = format!("{new_slot}@{new_root:#x}"); + let refreshed = ConsensusStore::open(temp.path(), Some(&new_checkpoint)).unwrap(); + + let checkpoint = refreshed.checkpoint(); + assert_eq!(checkpoint.beacon_slot, Some(new_slot)); + assert_eq!(checkpoint.beacon_root, new_root); + } + #[test] fn restart_guard_uses_recent_contiguous_consensus_anchor_for_progress() { let recent_slot = recent_checkpoint_slot(1); diff --git a/crates/logex-sync/src/engine/anchored.rs b/crates/logex-sync/src/engine/anchored.rs index fe0097e4..4ce16951 100644 --- a/crates/logex-sync/src/engine/anchored.rs +++ b/crates/logex-sync/src/engine/anchored.rs @@ -23,6 +23,7 @@ const CONSENSUS_ANCHOR_FORWARD_BODY_ATTEMPTS_DURING_HISTORICAL: usize = 2; const CONSENSUS_ANCHOR_FORWARD_RECEIPT_TIMEOUT_DURING_HISTORICAL: Duration = Duration::from_secs(2); const CONSENSUS_ANCHOR_FORWARD_RECEIPT_ATTEMPTS_DURING_HISTORICAL: usize = 2; const CONSENSUS_ANCHOR_FORWARD_RETRY_COOLDOWN_DURING_HISTORICAL: Duration = Duration::from_secs(8); +const CHECKPOINT_GAP_PEER_REFILL_TIMEOUT: Duration = Duration::from_secs(2); const CONSENSUS_READY_HISTORICAL_DRAIN_LIMIT: usize = 8; const HISTORICAL_VALIDATION_TASKS_PER_CPU: usize = 4; const HISTORICAL_VALIDATION_TASK_LIMIT: usize = 256; @@ -2117,6 +2118,32 @@ impl SyncEngine { .next_consensus_anchor_batch(current, &consensus, anchor_batch_limit) .await; if anchors.is_empty() { + if self + .ingest_gap_to_next_consensus_anchor( + current, + &consensus, + historical_backfill_active, + ) + .await? + { + continue; + } + if consensus.next_anchor_after(current).is_some() { + self.set_runtime_state(NodeState::Connecting); + if historical_pre_forward_progressed || historical_pre_anchor_progressed { + continue; + } + if cancelable( + &mut self.shutdown, + tokio::time::sleep(CONSENSUS_WAIT_INTERVAL), + ) + .await + .is_none() + { + return self.finish_shutdown(); + } + continue; + } if self.try_mark_synced("caught up to available consensus anchors") { self.sync_status_peers(); } else { @@ -2323,6 +2350,363 @@ impl SyncEngine { || prepare_progressed) } + async fn ingest_gap_to_next_consensus_anchor( + &mut self, + current: u64, + consensus: &ConsensusStore, + historical_backfill_active: bool, + ) -> Result { + let Some(anchor) = consensus.next_anchor_after(current) else { + return Ok(false); + }; + let Some(mut previous_header) = self.head_tracker.snapshot().last().cloned() else { + tracing::warn!( + current, + anchor_block = anchor.block_number, + "cannot bridge checkpoint gap without a persisted canonical header tip" + ); + return Ok(false); + }; + if anchor.block_number <= current.saturating_add(1) { + return Ok(false); + } + + let start_block = current.saturating_add(1); + let gap_blocks = anchor.block_number.saturating_sub(current); + tracing::info!( + start_block, + anchor_block = anchor.block_number, + gap_blocks, + "bridging stale restart gap to fresh consensus checkpoint" + ); + self.refill_checkpoint_gap_peers().await?; + + let mut headers = Vec::new(); + let mut header_peer = None; + let mut next_block = start_block; + while next_block <= anchor.block_number { + let remaining = anchor.block_number.saturating_sub(next_block) + 1; + let request_count = remaining.min(self.config.header_batch_size.max(1)); + let header_result = if historical_backfill_active { + cancelable( + &mut self.shutdown, + self.peers.get_headers_with_limits( + next_block, + request_count, + CONSENSUS_ANCHOR_FORWARD_HEADER_TIMEOUT_DURING_HISTORICAL, + CONSENSUS_ANCHOR_FORWARD_HEADER_ATTEMPTS_DURING_HISTORICAL, + ), + ) + .await + } else { + cancelable( + &mut self.shutdown, + self.peers.get_headers(next_block, request_count), + ) + .await + }; + let (peer_id, batch) = match header_result { + Some(Ok((peer_id, batch))) if batch.len() == request_count as usize => { + (peer_id, batch) + } + Some(Ok((peer_id, batch))) => { + tracing::warn!( + header_peer = %peer_id, + start_block = next_block, + requested = request_count, + returned = batch.len(), + "checkpoint gap header request returned an incomplete response" + ); + return Ok(false); + } + Some(Err(error)) => { + tracing::warn!( + error = %error, + start_block = next_block, + requested = request_count, + "checkpoint gap header request failed" + ); + return Ok(false); + } + None => { + self.finish_shutdown()?; + return Ok(false); + } + }; + + if let Err(error) = + validate_downloaded_headers(next_block, Some(&previous_header), &batch) + { + tracing::warn!( + start_block = next_block, + headers = batch.len(), + header_peer = %peer_id, + %error, + "checkpoint gap header validation failed" + ); + self.peers.report_invalid_block_data(peer_id, "headers"); + return Ok(false); + } + + previous_header = batch + .last() + .cloned() + .expect("non-empty checkpoint gap header batch"); + next_block = previous_header.number().saturating_add(1); + header_peer.get_or_insert(peer_id); + headers.extend(batch); + } + + let Some(terminal_header) = headers.last() else { + return Ok(false); + }; + let terminal_hash = terminal_header.hash_slow(); + if let Err(error) = validate_header_matches_anchor(&anchor, terminal_header, terminal_hash) + { + tracing::warn!( + anchor_block = anchor.block_number, + expected_hash = %anchor.block_hash, + got_hash = %terminal_hash, + %error, + "checkpoint gap terminal header did not match consensus anchor" + ); + return Ok(false); + } + + let header_peer = header_peer.expect("checkpoint gap contains at least one header batch"); + let hashes: Vec = headers.iter().map(|header| header.hash_slow()).collect(); + let mut newly_serving_peers = HashSet::new(); + let mut progressed = false; + let mut last_validated_header = None; + let mut last_head = None; + + for (chunk_headers, chunk_hashes) in headers + .chunks(self.config.fetch_batch_size) + .zip(hashes.chunks(self.config.fetch_batch_size)) + { + self.refill_checkpoint_gap_peers().await?; + let chunk_headers = chunk_headers.to_vec(); + let chunk_hashes = chunk_hashes.to_vec(); + let required_block = chunk_headers + .last() + .map(|header| header.number()) + .unwrap_or(anchor.block_number); + + let body_result = if historical_backfill_active { + cancelable( + &mut self.shutdown, + self.peers.get_bodies_prefer_peers_with_limits( + chunk_hashes.clone(), + required_block, + &[header_peer], + CONSENSUS_ANCHOR_FORWARD_BODY_TIMEOUT_DURING_HISTORICAL, + CONSENSUS_ANCHOR_FORWARD_BODY_ATTEMPTS_DURING_HISTORICAL, + ), + ) + .await + } else { + cancelable( + &mut self.shutdown, + self.peers.get_bodies_prefer_peers( + chunk_hashes.clone(), + required_block, + &[header_peer], + ), + ) + .await + }; + let bodies = match body_result { + Some(Ok(bodies)) if bodies.len() == chunk_headers.len() => bodies, + Some(Ok(bodies)) => { + tracing::warn!( + headers = chunk_headers.len(), + bodies = bodies.len(), + "checkpoint gap block body request returned an unexpected response" + ); + return Ok(progressed); + } + Some(Err(error)) => { + tracing::warn!(%error, "checkpoint gap block body request failed"); + return Ok(progressed); + } + None => { + self.finish_shutdown()?; + return Ok(progressed); + } + }; + + for (i, header) in chunk_headers.iter().enumerate() { + let block_number = header.number(); + let block_hash = chunk_hashes[i]; + let (body_peer, body) = &bodies[i]; + if let Err(error) = validate_block_pre_execution(header, block_hash, body) { + tracing::warn!( + block_number, + %block_hash, + body_peer = %body_peer, + %error, + "checkpoint gap block pre-execution validation failed" + ); + self.peers + .report_invalid_block_data(*body_peer, "block bodies"); + return Ok(progressed); + } + } + + let expected_receipt_counts: Vec = bodies + .iter() + .map(|(_peer_id, body)| body.transaction_count()) + .collect(); + let receipt_peer_preference = preferred_body_peers(&bodies, header_peer); + let receipt_result = if historical_backfill_active { + cancelable( + &mut self.shutdown, + self.peers + .get_receipts_matching_counts_prefer_peers_with_limits( + chunk_hashes.clone(), + required_block, + &expected_receipt_counts, + &receipt_peer_preference, + CONSENSUS_ANCHOR_FORWARD_RECEIPT_TIMEOUT_DURING_HISTORICAL, + CONSENSUS_ANCHOR_FORWARD_RECEIPT_ATTEMPTS_DURING_HISTORICAL, + ), + ) + .await + } else { + cancelable( + &mut self.shutdown, + self.peers.get_receipts_matching_counts_prefer_peers( + chunk_hashes.clone(), + required_block, + &expected_receipt_counts, + &receipt_peer_preference, + ), + ) + .await + }; + let (receipt_peer, receipts) = match receipt_result { + Some(Ok((peer_id, receipts))) if receipts.len() == chunk_headers.len() => { + (peer_id, receipts) + } + Some(Ok((peer_id, receipts))) => { + tracing::warn!( + headers = chunk_headers.len(), + receipt_peer = %peer_id, + returned_receipt_sets = receipts.len(), + "checkpoint gap receipt request returned an unexpected response" + ); + return Ok(progressed); + } + Some(Err(error)) => { + tracing::warn!(%error, "checkpoint gap receipt request failed"); + return Ok(progressed); + } + None => { + self.finish_shutdown()?; + return Ok(progressed); + } + }; + + for (i, header) in chunk_headers.iter().enumerate() { + let block_hash = chunk_hashes[i]; + let block_number = header.number(); + let (body_peer, body) = &bodies[i]; + + if !receipts_match_transaction_count(body, &receipts[i]) { + tracing::warn!( + block_number, + %block_hash, + receipt_peer = %receipt_peer, + transactions = body.transaction_count(), + receipts = receipts[i].len(), + "checkpoint gap block body / receipt count mismatch" + ); + self.peers + .report_invalid_block_data(receipt_peer, "receipts"); + return Ok(progressed); + } + + if let Err(error) = validate_receipts_for_header(header, &receipts[i]) { + tracing::warn!( + block_number, + %block_hash, + receipt_peer = %receipt_peer, + %error, + "checkpoint gap receipt validation failed" + ); + self.peers + .report_invalid_block_data(receipt_peer, "receipts"); + return Ok(progressed); + } + + let txs = assemble_txs(body, &receipts[i]); + if let Some(reorg) = self.head_tracker.track(header.clone()) { + self.handle_reorg(reorg).await?; + } + + let recent_headers = self.head_tracker.snapshot(); + self.peers + .cache_canonical_block(header.clone(), body.clone(), &receipts[i]); + let terminal_anchor = (block_number == anchor.block_number).then_some(&anchor); + let log_count = self + .ingest_block(header, block_hash, &txs, &recent_headers, terminal_anchor) + .await?; + self.progress.record_block(block_number, log_count); + self.note_serving_peer(header_peer, &mut newly_serving_peers); + self.note_serving_peer(*body_peer, &mut newly_serving_peers); + self.note_serving_peer(receipt_peer, &mut newly_serving_peers); + last_validated_header = Some(header.clone()); + last_head = Some(execution_head(block_number, block_hash, header.timestamp())); + progressed = true; + } + } + + self.last_validated_header = last_validated_header; + self.refresh_consensus_status().await; + self.refresh_historical_status().await; + + if let Some(head) = last_head { + self.peers.set_head(head); + } + + if self.try_mark_synced("bridged stale restart gap to consensus checkpoint") { + self.sync_status_peers(); + } else { + self.refresh_connectivity_state(); + } + + Ok(progressed) + } + + async fn refill_checkpoint_gap_peers(&mut self) -> Result<()> { + self.refresh_connectivity_state(); + let min_peers = peer_refill_goal( + self.peers.peer_count(), + self.peers.serving_peer_count(), + self.config.max_peers, + ) + .unwrap_or_else(|| refill_peer_floor(self.config.max_peers)); + match cancelable( + &mut self.shutdown, + tokio::time::timeout( + CHECKPOINT_GAP_PEER_REFILL_TIMEOUT, + self.peers.fill_peers(min_peers, self.config.max_peers), + ), + ) + .await + { + Some(Ok(())) => {} + Some(Err(_elapsed)) => { + self.peers.drain_events_now(); + } + None => { + self.finish_shutdown()?; + } + } + self.refresh_connectivity_state(); + Ok(()) + } + async fn refill_historical_fetch_pipeline_during_local_work(&mut self) -> Result { let child_header = { let storage = self.storage.read().await; @@ -6177,6 +6561,19 @@ fn locate_consensus_reorg( } } + let tip_number = tip.number(); + let has_anchor_at_or_before_tip = consensus + .ordered_anchors() + .iter() + .any(|record| record.anchor.block_number <= tip_number); + if !has_anchor_at_or_before_tip && consensus.next_anchor_after(tip_number).is_some() { + tracing::debug!( + tip_block = tip_number, + "consensus anchors are ahead of the persisted header window; treating as a restart gap" + ); + return Ok(None); + } + let first_block = recent_headers .first() .map(Header::number) @@ -7999,6 +8396,26 @@ mod tests { assert_eq!(reorg.reverted_hashes, vec![old_third.hash_slow()]); } + #[test] + fn locate_consensus_reorg_ignores_future_only_checkpoint_anchor() { + let temp = TempDir::new().unwrap(); + let first = header(100, B256::ZERO, 0x01); + let second = header(101, first.hash_slow(), 0x02); + let future = header(110, B256::repeat_byte(0x44), 0x10); + let store = ConsensusStore::open( + temp.path(), + Some("0xdddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"), + ) + .unwrap(); + store.append_anchors(vec![anchor_for(&future, 10)]).unwrap(); + + assert!( + locate_consensus_reorg(&store, &[first, second]) + .unwrap() + .is_none() + ); + } + #[test] fn locate_consensus_reorg_errors_when_window_has_no_common_ancestor() { let temp = TempDir::new().unwrap(); From af3a59a3573d4e21c033a0beeebd7615110e5bb5 Mon Sep 17 00:00:00 2001 From: Tim Denisenko Date: Thu, 2 Jul 2026 21:57:07 +0700 Subject: [PATCH 2/3] fix: pipeline stale checkpoint catchup --- ROADMAP.md | 25 +- crates/logex-sync/src/engine/anchored.rs | 576 +++++++++++++++++++---- crates/logex-sync/src/progress.rs | 27 +- 3 files changed, 532 insertions(+), 96 deletions(-) diff --git a/ROADMAP.md b/ROADMAP.md index 73314923..bcf3b32c 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -6,12 +6,14 @@ LogEx starts from a recent consensus checkpoint, tracks the live head, reverse-s Current branch: `fix/auto-refresh-stale-checkpoint`. -The Mac mini is running this branch against the saved full-sync data directory. Public dashboard forwarding through `157.245.195.72:18683` is active, and the client is bridging the stale restart gap from the saved execution head to the refreshed consensus checkpoint. +The Mac mini is running this branch against the saved full-sync data directory. Public dashboard forwarding through `157.245.195.72:18683` is active, and the client can bridge a stale restart gap from the saved execution head to a refreshed consensus checkpoint. ## Completed Since Last Run - Added automatic startup recovery for stale persisted consensus state and stale local execution progress. - Added a checkpoint-gap bridge that validates the EL parent chain up to a fresh CL checkpoint anchor before ingesting the missing logs. +- Pipelined checkpoint-gap body/receipt fetching with bounded lookahead so long-offline forward catch-up no longer waits for one chunk to fully ingest before requesting the next. +- Fixed checkpoint-gap progress accounting so dashboard logs/sec is sampled once per ingested batch instead of once per block after a batch write. - Kept the existing EL/log storage and known peer cache intact during checkpoint refresh. - Verified the remote dashboard path through the VPS and restarted the native Mac mini client with the fix. @@ -31,6 +33,11 @@ No remaining TODOs for the stale-checkpoint restart recovery task. - Alternatives considered: require a fresh data directory, or reintroduce CL historical sync. Both add unnecessary operational cost for this restart case. - Tradeoff: The bridge is correctness-first and less optimized than the normal historical pipeline. +- Account checkpoint-gap progress per contiguous batch. + - Why: The gap bridge ingests rows in batches; per-block progress updates after a batch write distort live logs/sec because each block update can be separated by only microseconds. + - Alternatives considered: keep per-block progress updates. That made the dashboard report impossible rates during catch-up. + - Tradeoff: The live forward metric is chunk-granular during restart recovery, which matches the actual batch-oriented work. + ## Challenges and Resolutions - Challenge: A full synced data directory failed to restart after being offline because the persisted consensus state and local execution head were older than the recent checkpoint window. @@ -39,11 +46,15 @@ No remaining TODOs for the stale-checkpoint restart recovery task. - Challenge: The first recovery attempt hit the consensus reorg guard because the fresh CL anchor was ahead of the persisted recent-header window. - Resolution: Future-only anchors are classified as restart gaps, then bridged by fetching and validating the EL header chain to the CL anchor. - - Remaining: Gap-bridge throughput can be optimized later if this path becomes common. + - Remaining: None known. + +- Challenge: The dashboard briefly reported impossible multi-billion logs/sec during checkpoint-gap catch-up. + - Resolution: Progress tracking now supports batched forward updates, and the gap bridge records one live rate sample per ingested chunk. + - Remaining: None known. ## Dead Code and Obsolescence Cleanup -- Inspected the stale startup guards and consensus reorg classifier. +- Inspected the stale startup guards, consensus reorg classifier, checkpoint-gap bridge, and progress tracker. - No obsolete code was safely removable in this pass; the new bridge reuses existing validation, peer, and storage primitives. - Removed no files. @@ -51,16 +62,20 @@ No remaining TODOs for the stale-checkpoint restart recovery task. - Current branch: `fix/auto-refresh-stale-checkpoint`. - New branch created from `master`. -- Commits made during this run: pending. +- Commits made during this run: + - `7dc7c6fd fix: recover stale checkpoint restarts` + - pending checkpoint-gap pipeline/progress commit. - Pull request status: pending after commit/push. - Merge status: not merged yet. - Validation run: - `cargo test -p logex-node archived_consensus_state` - `cargo test -p logex-node restart_guard` - `cargo test -p logex-sync locate_consensus_reorg` + - `cargo test -p logex-sync checkpoint_gap_pipeline_depth` + - `cargo test -p logex-sync forward_batch_progress_records_one_live_rate_sample` - `cargo clippy -p logex-node -p logex-sync --all-targets -- -D warnings` ## Known Issues or Risks -- The checkpoint-gap bridge prioritizes trustless recovery over throughput. Normal historical sync performance is unchanged. +- The checkpoint-gap bridge is only for long-offline restart recovery. Normal historical reverse sync and live tip-following paths are unchanged. - The local Codex sandbox cannot directly curl the public VPS dashboard, but the VPS can reach `10.66.0.2:18683` and packet capture showed public TCP/18683 traffic being forwarded and answered. diff --git a/crates/logex-sync/src/engine/anchored.rs b/crates/logex-sync/src/engine/anchored.rs index 4ce16951..0f2567ea 100644 --- a/crates/logex-sync/src/engine/anchored.rs +++ b/crates/logex-sync/src/engine/anchored.rs @@ -24,6 +24,8 @@ const CONSENSUS_ANCHOR_FORWARD_RECEIPT_TIMEOUT_DURING_HISTORICAL: Duration = Dur const CONSENSUS_ANCHOR_FORWARD_RECEIPT_ATTEMPTS_DURING_HISTORICAL: usize = 2; const CONSENSUS_ANCHOR_FORWARD_RETRY_COOLDOWN_DURING_HISTORICAL: Duration = Duration::from_secs(8); const CHECKPOINT_GAP_PEER_REFILL_TIMEOUT: Duration = Duration::from_secs(2); +const CHECKPOINT_GAP_PIPELINE_CHUNK_BLOCKS: usize = 128; +const CHECKPOINT_GAP_PIPELINE_DEPTH: usize = 8; const CONSENSUS_READY_HISTORICAL_DRAIN_LIMIT: usize = 8; const HISTORICAL_VALIDATION_TASKS_PER_CPU: usize = 4; const HISTORICAL_VALIDATION_TASK_LIMIT: usize = 256; @@ -122,6 +124,31 @@ struct HistoricalValidationJob { receipts: Vec::Receipt>>, } +struct ForwardGapFetchOutcome { + sequence: u64, + header_peer: PeerId, + headers: Vec
, + hashes: Vec, + required_block: u64, + body_receipt_elapsed: Duration, + outcome: crate::p2p::peer_manager::BodyReceiptRequestOutcome, +} + +struct ForwardGapFetchedChunk { + sequence: u64, + header_peer: PeerId, + headers: Vec
, + hashes: Vec, + body_receipt_elapsed: Duration, + blocks: Vec, +} + +struct ForwardGapCanonicalUpdate { + header: Header, + recent_headers: Vec
, + anchor: Option, +} + struct HistoricalValidationExtractedChunk { start_index: usize, peer_notes: Vec, @@ -1911,6 +1938,22 @@ fn historical_header_has_empty_body_and_receipts(header: &Header) -> bool { .is_none_or(|root| root == EMPTY_ROOT_HASH) } +fn checkpoint_gap_parallel_min_blocks() -> usize { + 64 +} + +fn checkpoint_gap_pipeline_depth(serving_peers: usize) -> usize { + if serving_peers >= 32 { + CHECKPOINT_GAP_PIPELINE_DEPTH + } else if serving_peers >= 16 { + 6 + } else if serving_peers >= 8 { + 4 + } else { + 2 + } +} + fn historical_header_batch_matches_child( batch: &HistoricalHeaderBatch, child_header: &Header, @@ -2475,7 +2518,269 @@ impl SyncEngine { let header_peer = header_peer.expect("checkpoint gap contains at least one header batch"); let hashes: Vec = headers.iter().map(|header| header.hash_slow()).collect(); - let mut newly_serving_peers = HashSet::new(); + let (progressed, last_validated_header, last_head) = self + .ingest_checkpoint_gap_fetch_pipeline( + header_peer, + headers, + hashes, + anchor, + historical_backfill_active, + ) + .await?; + + self.last_validated_header = last_validated_header; + self.refresh_consensus_status().await; + self.refresh_historical_status().await; + + if let Some(head) = last_head { + self.peers.set_head(head); + } + + if self.try_mark_synced("bridged stale restart gap to consensus checkpoint") { + self.sync_status_peers(); + } else { + self.refresh_connectivity_state(); + } + + Ok(progressed) + } + + async fn ingest_checkpoint_gap_fetch_pipeline( + &mut self, + header_peer: PeerId, + headers: Vec
, + hashes: Vec, + anchor: ExecutionAnchor, + historical_backfill_active: bool, + ) -> Result<(bool, Option
, Option)> { + if headers.is_empty() { + return Ok((false, None, None)); + } + if headers.len() != hashes.len() { + return Ok((false, None, None)); + } + + let mut active_fetches = JoinSet::new(); + let mut completed_fetches = BTreeMap::new(); + let mut next_sequence = 0u64; + let mut expected_sequence = 0u64; + let mut next_offset = 0usize; + let mut expected_offset = 0usize; + let mut progressed = false; + let mut last_validated_header = None; + let mut last_head = None; + + while expected_offset < headers.len() { + while active_fetches.len().saturating_add(completed_fetches.len()) + < checkpoint_gap_pipeline_depth(self.peers.serving_peer_count()) + && next_offset < headers.len() + { + let remaining = headers.len().saturating_sub(next_offset); + let chunk_len = remaining.min(CHECKPOINT_GAP_PIPELINE_CHUNK_BLOCKS); + if chunk_len < checkpoint_gap_parallel_min_blocks() { + break; + } + let chunk_end = next_offset.saturating_add(chunk_len); + let chunk_headers = headers[next_offset..chunk_end].to_vec(); + let chunk_hashes = hashes[next_offset..chunk_end].to_vec(); + if !self + .spawn_checkpoint_gap_fetch_task( + next_sequence, + header_peer, + chunk_headers, + chunk_hashes, + &mut active_fetches, + ) + .await? + { + break; + } + next_sequence = next_sequence.saturating_add(1); + next_offset = chunk_end; + tokio::task::yield_now().await; + self.drain_historical_request_accounting(); + } + + if let Some(outcome) = completed_fetches.remove(&expected_sequence) { + let Some(chunk) = self.materialize_checkpoint_gap_fetch_outcome(outcome)? else { + return Ok((progressed, last_validated_header, last_head)); + }; + let (chunk_progressed, chunk_last_header, chunk_last_head) = + self.ingest_forward_gap_fetched_chunk(chunk, anchor).await?; + if !chunk_progressed { + return Ok((progressed, last_validated_header, last_head)); + } + let chunk_block_count = chunk_last_header + .as_ref() + .and_then(|header| { + header + .number() + .checked_sub(headers[expected_offset].number()) + .map(|delta| delta.saturating_add(1) as usize) + }) + .unwrap_or(0); + expected_offset = expected_offset.saturating_add(chunk_block_count); + expected_sequence = expected_sequence.saturating_add(1); + progressed = true; + last_validated_header = chunk_last_header; + last_head = chunk_last_head; + continue; + } + + if active_fetches.is_empty() { + let (tail_progressed, tail_last_header, tail_last_head) = self + .ingest_checkpoint_gap_sequential_tail( + header_peer, + &headers[expected_offset..], + &hashes[expected_offset..], + anchor, + historical_backfill_active, + ) + .await?; + return Ok(( + progressed || tail_progressed, + tail_last_header.or(last_validated_header), + tail_last_head.or(last_head), + )); + } + + tokio::select! { + result = active_fetches.join_next() => { + match result { + Some(Ok(outcome)) => { + completed_fetches.insert(outcome.sequence, outcome); + self.drain_historical_request_accounting(); + } + Some(Err(error)) => { + tracing::warn!(%error, "checkpoint gap body/receipt pipeline worker failed"); + return Ok((progressed, last_validated_header, last_head)); + } + None => {} + } + } + changed = self.shutdown.changed() => { + if changed.is_ok() && self.shutdown_requested() { + active_fetches.abort_all(); + self.finish_shutdown()?; + return Ok((progressed, last_validated_header, last_head)); + } + } + } + } + + active_fetches.abort_all(); + Ok((progressed, last_validated_header, last_head)) + } + + async fn spawn_checkpoint_gap_fetch_task( + &mut self, + sequence: u64, + header_peer: PeerId, + headers: Vec
, + hashes: Vec, + active_fetches: &mut JoinSet, + ) -> Result { + let Some(required_block) = headers.last().map(|header| header.number()) else { + return Ok(false); + }; + let gas_used = headers.iter().map(|header| header.gas_used()).collect(); + self.drain_historical_request_accounting(); + let Some(plan) = self + .peers + .prepare_bodies_and_receipts_request_for_hashes_and_gas( + hashes.clone(), + gas_used, + self.historical_rows_per_block_ewma, + required_block, + &[header_peer], + ) + .await? + else { + return Ok(false); + }; + let plan = plan + .with_full_priority() + .with_peer_rotation_offset(sequence as usize) + .with_accounting_tx(self.historical_request_accounting_tx.clone()); + active_fetches.spawn(async move { + let started = std::time::Instant::now(); + let outcome = plan.execute().await; + ForwardGapFetchOutcome { + sequence, + header_peer, + headers, + hashes, + required_block, + body_receipt_elapsed: started.elapsed(), + outcome, + } + }); + Ok(true) + } + + fn materialize_checkpoint_gap_fetch_outcome( + &mut self, + outcome: ForwardGapFetchOutcome, + ) -> Result> { + let ForwardGapFetchOutcome { + sequence, + header_peer, + headers, + hashes, + required_block, + body_receipt_elapsed, + outcome, + } = outcome; + match self.peers.complete_bodies_and_receipts_request(outcome) { + Ok(Some(completion)) if completion.blocks.len() == headers.len() => { + tracing::debug!( + sequence, + required_block, + blocks = completion.blocks.len(), + body_receipt_ms = body_receipt_elapsed.as_millis(), + "checkpoint gap body/receipt pipeline chunk completed" + ); + Ok(Some(ForwardGapFetchedChunk { + sequence, + header_peer, + headers, + hashes, + body_receipt_elapsed, + blocks: completion.blocks, + })) + } + Ok(Some(completion)) => { + tracing::warn!( + sequence, + required_block, + expected = headers.len(), + returned = completion.blocks.len(), + "checkpoint gap body/receipt pipeline returned a partial chunk" + ); + Ok(None) + } + Ok(None) => Ok(None), + Err(error) => { + tracing::warn!( + sequence, + required_block, + %error, + "checkpoint gap body/receipt pipeline failed" + ); + self.refresh_connectivity_state(); + Ok(None) + } + } + } + + async fn ingest_checkpoint_gap_sequential_tail( + &mut self, + header_peer: PeerId, + headers: &[Header], + hashes: &[B256], + anchor: ExecutionAnchor, + historical_backfill_active: bool, + ) -> Result<(bool, Option
, Option)> { let mut progressed = false; let mut last_validated_header = None; let mut last_head = None; @@ -2487,10 +2792,9 @@ impl SyncEngine { self.refill_checkpoint_gap_peers().await?; let chunk_headers = chunk_headers.to_vec(); let chunk_hashes = chunk_hashes.to_vec(); - let required_block = chunk_headers - .last() - .map(|header| header.number()) - .unwrap_or(anchor.block_number); + let Some(required_block) = chunk_headers.last().map(|header| header.number()) else { + continue; + }; let body_result = if historical_backfill_active { cancelable( @@ -2521,38 +2825,20 @@ impl SyncEngine { tracing::warn!( headers = chunk_headers.len(), bodies = bodies.len(), - "checkpoint gap block body request returned an unexpected response" + "checkpoint gap tail body request returned an unexpected response" ); - return Ok(progressed); + return Ok((progressed, last_validated_header, last_head)); } Some(Err(error)) => { - tracing::warn!(%error, "checkpoint gap block body request failed"); - return Ok(progressed); + tracing::warn!(%error, "checkpoint gap tail body request failed"); + return Ok((progressed, last_validated_header, last_head)); } None => { self.finish_shutdown()?; - return Ok(progressed); + return Ok((progressed, last_validated_header, last_head)); } }; - for (i, header) in chunk_headers.iter().enumerate() { - let block_number = header.number(); - let block_hash = chunk_hashes[i]; - let (body_peer, body) = &bodies[i]; - if let Err(error) = validate_block_pre_execution(header, block_hash, body) { - tracing::warn!( - block_number, - %block_hash, - body_peer = %body_peer, - %error, - "checkpoint gap block pre-execution validation failed" - ); - self.peers - .report_invalid_block_data(*body_peer, "block bodies"); - return Ok(progressed); - } - } - let expected_receipt_counts: Vec = bodies .iter() .map(|(_peer_id, body)| body.transaction_count()) @@ -2593,89 +2879,186 @@ impl SyncEngine { headers = chunk_headers.len(), receipt_peer = %peer_id, returned_receipt_sets = receipts.len(), - "checkpoint gap receipt request returned an unexpected response" + "checkpoint gap tail receipt request returned an unexpected response" ); - return Ok(progressed); + return Ok((progressed, last_validated_header, last_head)); } Some(Err(error)) => { - tracing::warn!(%error, "checkpoint gap receipt request failed"); - return Ok(progressed); + tracing::warn!(%error, "checkpoint gap tail receipt request failed"); + return Ok((progressed, last_validated_header, last_head)); } None => { self.finish_shutdown()?; - return Ok(progressed); + return Ok((progressed, last_validated_header, last_head)); } }; - for (i, header) in chunk_headers.iter().enumerate() { - let block_hash = chunk_hashes[i]; - let block_number = header.number(); - let (body_peer, body) = &bodies[i]; - - if !receipts_match_transaction_count(body, &receipts[i]) { - tracing::warn!( - block_number, - %block_hash, - receipt_peer = %receipt_peer, - transactions = body.transaction_count(), - receipts = receipts[i].len(), - "checkpoint gap block body / receipt count mismatch" - ); - self.peers - .report_invalid_block_data(receipt_peer, "receipts"); - return Ok(progressed); - } + let blocks: Vec = bodies + .into_iter() + .zip( + receipts + .into_iter() + .map(|receipt_set| (receipt_peer, receipt_set)), + ) + .collect(); + let chunk = ForwardGapFetchedChunk { + sequence: 0, + header_peer, + headers: chunk_headers, + hashes: chunk_hashes, + body_receipt_elapsed: Duration::ZERO, + blocks, + }; + let (chunk_progressed, chunk_last_header, chunk_last_head) = + self.ingest_forward_gap_fetched_chunk(chunk, anchor).await?; + if !chunk_progressed { + return Ok((progressed, last_validated_header, last_head)); + } + progressed = true; + last_validated_header = chunk_last_header; + last_head = chunk_last_head; + } - if let Err(error) = validate_receipts_for_header(header, &receipts[i]) { - tracing::warn!( - block_number, - %block_hash, - receipt_peer = %receipt_peer, - %error, - "checkpoint gap receipt validation failed" - ); - self.peers - .report_invalid_block_data(receipt_peer, "receipts"); - return Ok(progressed); - } + Ok((progressed, last_validated_header, last_head)) + } - let txs = assemble_txs(body, &receipts[i]); - if let Some(reorg) = self.head_tracker.track(header.clone()) { - self.handle_reorg(reorg).await?; - } + async fn ingest_forward_gap_fetched_chunk( + &mut self, + chunk: ForwardGapFetchedChunk, + anchor: ExecutionAnchor, + ) -> Result<(bool, Option
, Option)> { + let ForwardGapFetchedChunk { + sequence, + header_peer, + headers, + hashes, + body_receipt_elapsed, + blocks, + } = chunk; + if headers.is_empty() { + return Ok((false, None, None)); + } + if blocks.len() != headers.len() { + tracing::warn!( + sequence, + headers = headers.len(), + blocks = blocks.len(), + "checkpoint gap fetched chunk has mismatched block count" + ); + return Ok((false, None, None)); + } - let recent_headers = self.head_tracker.snapshot(); + let validation_started = std::time::Instant::now(); + let validated = match validate_historical_blocks_parallel(&headers, &hashes, blocks).await? + { + Ok(validated) => validated, + Err(failure) => { + tracing::warn!( + block_number = failure.block_number, + block_hash = %failure.block_hash, + peer = %failure.peer, + error = %failure.message, + "checkpoint gap block validation failed" + ); self.peers - .cache_canonical_block(header.clone(), body.clone(), &receipts[i]); - let terminal_anchor = (block_number == anchor.block_number).then_some(&anchor); - let log_count = self - .ingest_block(header, block_hash, &txs, &recent_headers, terminal_anchor) - .await?; - self.progress.record_block(block_number, log_count); - self.note_serving_peer(header_peer, &mut newly_serving_peers); - self.note_serving_peer(*body_peer, &mut newly_serving_peers); - self.note_serving_peer(receipt_peer, &mut newly_serving_peers); - last_validated_header = Some(header.clone()); - last_head = Some(execution_head(block_number, block_hash, header.timestamp())); - progressed = true; + .report_invalid_block_data(failure.peer, failure.response_kind); + return Ok((false, None, None)); } + }; + + let mut newly_serving_peers = HashSet::new(); + let mut canonical_updates = Vec::with_capacity(validated.len()); + let mut rows = Vec::new(); + let mut last_validated_header = None; + let mut last_head = None; + let mut last_progress_block = None; + + for block in validated { + let HistoricalValidatedBlock { + header, + block_hash, + body_peer, + body, + receipt_peer, + receipts, + .. + } = block; + let block_number = header.number(); + if let Some(reorg) = self.head_tracker.track(header.clone()) { + self.handle_reorg(reorg).await?; + } + + let recent_headers = self.head_tracker.snapshot(); + self.peers + .cache_canonical_block(header.clone(), body.clone(), &receipts); + extract::append_from_body_receipts( + &mut rows, + block_number, + block_hash, + header.timestamp(), + &body, + &receipts, + ); + self.note_serving_peer(header_peer, &mut newly_serving_peers); + self.note_serving_peer(body_peer, &mut newly_serving_peers); + self.note_serving_peer(receipt_peer, &mut newly_serving_peers); + canonical_updates.push(ForwardGapCanonicalUpdate { + anchor: (block_number == anchor.block_number).then_some(anchor), + header: header.clone(), + recent_headers, + }); + last_head = Some(execution_head(block_number, block_hash, header.timestamp())); + last_validated_header = Some(header); + last_progress_block = Some(block_number); } - self.last_validated_header = last_validated_header; - self.refresh_consensus_status().await; - self.refresh_historical_status().await; + { + let mut storage = self.storage.write().await; + if !rows.is_empty() { + storage + .write_batch(&rows) + .map_err(|error| eyre::eyre!("storage write error: {error}"))?; - if let Some(head) = last_head { - self.peers.set_head(head); + if let Some(ref subs) = self.subscriptions { + subs.notify(&rows); + } + } + for update in &canonical_updates { + match update.anchor.as_ref() { + Some(anchor) => { + storage + .record_verified_canonical_state( + anchor, + &update.header, + &update.recent_headers, + ) + .map_err(|error| eyre::eyre!("storage metadata error: {error}"))?; + storage + .record_historical_floor(&update.header) + .map_err(|error| eyre::eyre!("historical metadata error: {error}"))?; + } + None => storage + .record_canonical_state(&update.header, &update.recent_headers) + .map_err(|error| eyre::eyre!("storage metadata error: {error}"))?, + } + } } - if self.try_mark_synced("bridged stale restart gap to consensus checkpoint") { - self.sync_status_peers(); - } else { - self.refresh_connectivity_state(); + if let Some(block_number) = last_progress_block { + self.progress + .record_blocks(block_number, canonical_updates.len() as u64, rows.len() as u64); } - Ok(progressed) + tracing::debug!( + sequence, + blocks = canonical_updates.len(), + logs = rows.len(), + body_receipt_ms = body_receipt_elapsed.as_millis(), + validation_ms = validation_started.elapsed().as_millis(), + "checkpoint gap chunk verified and ingested" + ); + + Ok((true, last_validated_header, last_head)) } async fn refill_checkpoint_gap_peers(&mut self) -> Result<()> { @@ -8132,6 +8515,19 @@ mod tests { ); } + #[test] + fn checkpoint_gap_pipeline_depth_scales_with_serving_peers() { + assert_eq!(checkpoint_gap_parallel_min_blocks(), 64); + assert_eq!(checkpoint_gap_pipeline_depth(0), 2); + assert_eq!(checkpoint_gap_pipeline_depth(7), 2); + assert_eq!(checkpoint_gap_pipeline_depth(8), 4); + assert_eq!(checkpoint_gap_pipeline_depth(16), 6); + assert_eq!( + checkpoint_gap_pipeline_depth(32), + CHECKPOINT_GAP_PIPELINE_DEPTH + ); + } + #[test] fn historical_advanced_fetch_position_walks_contiguous_materialized_sequences() { let h90 = header(90, B256::ZERO, 0x01); diff --git a/crates/logex-sync/src/progress.rs b/crates/logex-sync/src/progress.rs index 6d6898b9..7d72ac37 100644 --- a/crates/logex-sync/src/progress.rs +++ b/crates/logex-sync/src/progress.rs @@ -137,7 +137,16 @@ impl ProgressTracker { /// Record that a block has been ingested. pub fn record_block(&mut self, block_number: u64, log_count: u64) { - self.blocks_processed += 1; + self.record_blocks(block_number, 1, log_count); + } + + /// Record that a contiguous batch of forward blocks has been ingested. + pub fn record_blocks(&mut self, block_number: u64, block_count: u64, log_count: u64) { + if block_count == 0 { + return; + } + + self.blocks_processed += block_count; self.logs_ingested += log_count; let now = Instant::now(); @@ -453,6 +462,22 @@ mod tests { assert!(decayed < 100.0); } + #[test] + fn forward_batch_progress_records_one_live_rate_sample() { + let status = Arc::new(Mutex::new(SyncStatus::default())); + let mut tracker = ProgressTracker::new(Arc::clone(&status)); + + std::thread::sleep(std::time::Duration::from_millis(10)); + tracker.record_blocks(128, 128, 1_280); + + let status = status.lock().unwrap().clone(); + assert_eq!(status.current_block, 128); + assert_eq!(status.logs_ingested, 1_280); + assert!(status.blocks_per_sec > 0.0); + assert!(status.logs_per_sec > 0.0); + assert!(status.logs_per_sec < 1_000_000.0); + } + #[test] fn consensus_target_can_move_backwards_after_reorg() { let status = Arc::new(Mutex::new(SyncStatus { From a6f6c41980205d659f6d92a9716499b5d530df3f Mon Sep 17 00:00:00 2001 From: Tim Denisenko Date: Thu, 2 Jul 2026 21:58:18 +0700 Subject: [PATCH 3/3] docs: update stale checkpoint PR status --- ROADMAP.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/ROADMAP.md b/ROADMAP.md index bcf3b32c..521d0f53 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -64,8 +64,8 @@ No remaining TODOs for the stale-checkpoint restart recovery task. - New branch created from `master`. - Commits made during this run: - `7dc7c6fd fix: recover stale checkpoint restarts` - - pending checkpoint-gap pipeline/progress commit. -- Pull request status: pending after commit/push. + - `af3a59a3 fix: pipeline stale checkpoint catchup` +- Pull request status: draft PR #102 opened for `fix/auto-refresh-stale-checkpoint` into `master`. - Merge status: not merged yet. - Validation run: - `cargo test -p logex-node archived_consensus_state`