Skip to main content

reth_snap_sync/
bootstrap.rs

1//! Runs snap synchronization from pivot selection until its state is handed to the trie rebuild.
2//!
3//! Every step commits before the next one starts and the attempt record is re-read on each pass,
4//! so a run stopped anywhere resumes from what the database holds.
5
6use crate::{
7    AccountRangeDownload, AccountRangeStep, BlockAccessListCatchUp, BytecodeDownload, BytecodeStep,
8    CatchUpStep, SnapAccountStore, SnapAttemptStore, SnapCatchUpStore, SnapGeneration,
9    SnapPivotPolicy, SnapStateVerifier, SnapSyncError, SnapSyncSession, SnapWrite,
10    StorageRangeDownload, StorageRangeStep, VerifiedRange, DEFAULT_SCAN_CHUNK,
11};
12use alloy_eips::BlockNumHash;
13use reth_db_api::transaction::DbTxMut;
14use reth_network_p2p::snap::client::SnapClient;
15use reth_primitives_traits::AlloyBlockHeader;
16use reth_storage_api::{
17    BlockHashReader, DBProvider, DatabaseProviderFactory, HeaderProvider, MetadataProvider,
18    MetadataWriter, StageCheckpointWriter, StateWriter,
19};
20use reth_storage_errors::provider::ProviderError;
21use reth_tasks::Runtime;
22use std::{fmt, future::Future};
23use tokio_util::sync::CancellationToken;
24use tracing::{debug, info};
25
26/// Default number of account ranges committed between pivot checks.
27///
28/// Bounds how long a pivot can lag unnoticed while ranges download.
29pub const DEFAULT_RANGES_PER_CHECK: usize = 64;
30
31/// Drives snap synchronization until its downloaded state is ready for the trie rebuild.
32///
33/// Each pass moves a lagging pivot forward, applies the lists carrying the state to it, refetches
34/// scheduled repairs there, then downloads account ranges with storage and code. A reorg across the
35/// pivot is repaired from the orphaned lists, while handed-off or expired state starts over.
36pub struct SnapBootstrap<C, F, X> {
37    factory: F,
38    // Blocking database work runs off the async worker.
39    runtime: Runtime,
40    // Head, finality and peer progress come from the node running the sync.
41    context: X,
42    policy: SnapPivotPolicy,
43    // Target of the attempt being driven, rebuilt from the attempt record on every pass.
44    session: SnapSyncSession,
45    ranges_per_check: usize,
46    // Stops the run at the next step boundary, leaving committed progress in place.
47    cancel: CancellationToken,
48    accounts: AccountRangeDownload<C, F>,
49    storage: StorageRangeDownload<C, F>,
50    bytecode: BytecodeDownload<C, F>,
51    catch_up: BlockAccessListCatchUp<C, F>,
52}
53
54impl<C: Clone, F: Clone, X> SnapBootstrap<C, F, X> {
55    /// Creates a run that has not touched the network or the database yet.
56    pub fn new(client: C, factory: F, runtime: Runtime, context: X) -> Self {
57        let policy = SnapPivotPolicy::default();
58        Self {
59            accounts: AccountRangeDownload::new(client.clone(), factory.clone(), runtime.clone()),
60            storage: StorageRangeDownload::new(client.clone(), factory.clone(), runtime.clone()),
61            bytecode: BytecodeDownload::new(client.clone(), factory.clone(), runtime.clone()),
62            catch_up: BlockAccessListCatchUp::new(client, factory.clone(), runtime.clone()),
63            factory,
64            runtime,
65            context,
66            policy,
67            session: SnapSyncSession::new(policy),
68            ranges_per_check: DEFAULT_RANGES_PER_CHECK,
69            cancel: CancellationToken::new(),
70        }
71    }
72}
73
74impl<C, F, X> SnapBootstrap<C, F, X> {
75    /// Returns this run anchoring and re-anchoring pivots with `policy`.
76    pub fn with_policy(mut self, policy: SnapPivotPolicy) -> Self {
77        self.policy = policy;
78        self.session = SnapSyncSession::new(policy);
79        self
80    }
81
82    /// Returns this run committing at most `ranges` account ranges between pivot checks, at
83    /// least one.
84    pub const fn with_ranges_per_check(mut self, ranges: usize) -> Self {
85        self.ranges_per_check = if ranges == 0 { 1 } else { ranges };
86        self
87    }
88
89    /// Returns this run stopping once `cancel` fires.
90    pub fn with_cancellation(mut self, cancel: CancellationToken) -> Self {
91        self.cancel = cancel;
92        self
93    }
94}
95
96impl<C, F, X> SnapBootstrap<C, F, X>
97where
98    C: SnapClient + Clone + Unpin,
99    F: DatabaseProviderFactory + Clone + 'static,
100    F::Provider: HeaderProvider + MetadataProvider + BlockHashReader,
101    F::ProviderRW: HeaderProvider
102        + BlockHashReader
103        + MetadataProvider
104        + MetadataWriter
105        + StageCheckpointWriter
106        + StateWriter
107        + DBProvider<Tx: DbTxMut>,
108    X: SnapSyncContext,
109{
110    /// Runs until the downloaded state is handed to the merkle stage, or the run stops.
111    ///
112    /// Peers that do not serve the state, or headers not downloaded yet, wait for the context to
113    /// report progress instead of failing the run.
114    pub async fn run(&mut self) -> Result<SnapBootstrapOutcome, SnapSyncError> {
115        loop {
116            if self.cancel.is_cancelled() {
117                return Ok(SnapBootstrapOutcome::Stopped)
118            }
119            let head = self.context.head()?;
120            let step = match self.resolve(head)? {
121                Resolved::Verified(pivot) => return Ok(SnapBootstrapOutcome::Verified { pivot }),
122                Resolved::Waiting => {
123                    debug!(target: "sync::snap", head, "No eligible snap pivot");
124                    Step::Wait
125                }
126                Resolved::Active(write) => match self.drive(write, head).await {
127                    Ok(step) => step,
128                    Err(error) if error.is_transient() => {
129                        debug!(target: "sync::snap", %error, "Waiting for peers or headers");
130                        Step::Wait
131                    }
132                    // The next pass finds the pivot orphaned and recovers or restarts it.
133                    Err(error) if error.is_reorg() => {
134                        debug!(target: "sync::snap", %error, "Snap pivot was reorged");
135                        Step::Continue
136                    }
137                    // Progress this build cannot read is unusable, so the attempt starts over.
138                    Err(error @ SnapSyncError::UnsupportedRecord { .. }) => {
139                        info!(target: "sync::snap", %error, "Snap progress is unreadable, restarting");
140                        Step::Restart
141                    }
142                    Err(SnapSyncError::Cancelled) => Step::Stop,
143                    Err(error) => return Err(error),
144                },
145            };
146            match step {
147                Step::HandedOff(write) => {
148                    let pivot = self.pivot()?;
149                    info!(target: "sync::snap", ?pivot, "Snap state handed to the trie rebuild");
150                    return Ok(SnapBootstrapOutcome::TrieRebuild { write, pivot })
151                }
152                Step::Continue => {}
153                Step::Wait => {
154                    if !self.wait(head).await {
155                        return Ok(SnapBootstrapOutcome::Stopped)
156                    }
157                }
158                Step::Restart => {
159                    let provider = self.factory.database_provider_rw()?;
160                    provider.abandon_snap_attempt()?;
161                    provider.commit()?;
162                }
163                Step::Stop => return Ok(SnapBootstrapOutcome::Stopped),
164            }
165        }
166    }
167
168    // Resumes the recorded attempt, even one whose pivot a reorg orphaned, otherwise starts one at
169    // the pivot under `head`. Starting one reads the kept headers and writes a few records, cheap
170    // enough for the async worker.
171    fn resolve(&mut self, head: u64) -> Result<Resolved, SnapSyncError> {
172        let provider = self.factory.database_provider_rw()?;
173        let mut session = SnapSyncSession::new(self.policy);
174        if let Some(attempt) = provider.snap_attempt()? {
175            if attempt.is_verified() {
176                return Ok(Resolved::Verified(attempt.pivot()))
177            }
178            if let Some(write) = provider.active_snap_write()? {
179                // An orphaned pivot is resumed too, so the pass can recover or restart it.
180                session.resume(SnapGeneration::new(attempt.pivot(), attempt.state_root()));
181                self.session = session;
182                return Ok(Resolved::Active(write))
183            }
184        }
185
186        session.select(&provider, head, self.context.finalized())?;
187        let Some((generation, _)) = session.start() else { return Ok(Resolved::Waiting) };
188        let write = provider.start_snap_attempt(generation)?;
189        provider.start_account_coverage(write)?;
190        provider.commit()?;
191        info!(
192            target: "sync::snap",
193            pivot = ?generation.target(),
194            state_root = %generation.state_root(),
195            "Started snap attempt"
196        );
197        self.session = session;
198        Ok(Resolved::Active(write))
199    }
200
201    // One pass over the attempt: reorg recovery, pivot, catch-up, repairs, then up to
202    // `ranges_per_check` account ranges.
203    async fn drive(&mut self, write: SnapWrite, head: u64) -> Result<Step, SnapSyncError> {
204        let (applied, complete) = {
205            let provider = self.factory.database_provider_ro()?;
206            let generation = self.session.target().copied().ok_or(SnapSyncError::NoAttempt)?;
207            let handed_off = provider.is_trie_rebuild_started(write)?;
208            if !generation.is_canonical(&provider)? {
209                // Handed-off state may be partly rebuilt, so it cannot be repaired in place.
210                if handed_off {
211                    info!(target: "sync::snap", pivot = ?generation.target(), "Handed-off snap pivot was reorged, restarting");
212                    return Ok(Step::Restart)
213                }
214                drop(provider);
215                return self.recover(write, head).await
216            }
217            if handed_off {
218                return Ok(Step::HandedOff(write))
219            }
220            let applied = provider
221                .catch_up_progress(write)?
222                .ok_or(SnapSyncError::NoCatchUpProgress)?
223                .applied();
224            // Repairs are fetched at the pivot, so pending ones still need it served.
225            let complete = provider.account_coverage(write)?.is_some_and(|c| c.is_complete()) &&
226                provider.snap_repairs(write)?.is_empty();
227            (applied, complete)
228        };
229        // Complete state carried to its pivot needs no further lists, only its trie rebuilt, so
230        // neither their retention nor a newer pivot applies to it.
231        if complete && applied == self.pivot()? {
232            return self.hand_off(write).await
233        }
234        // Catch-up continues from the last applied block, not the pivot, so once that block's
235        // successor is no longer served the state cannot be carried to any newer pivot.
236        if !self.policy.is_catchable_from(applied.number, head) {
237            info!(target: "sync::snap", ?applied, head, "Snap catch-up outlived the served block access lists, restarting");
238            return Ok(Step::Restart)
239        }
240        let write = self.advance_pivot(write, head)?;
241        // Lists only update accounts already downloaded, so coverage grows once they reach the
242        // pivot.
243        if let Some(step) = self.catch_up(write).await? {
244            return Ok(step)
245        }
246        // Repairs take the pivot's values, which the rest of the state holds once the lists reach
247        // it.
248        if let Some(step) = self.repair().await? {
249            return Ok(step)
250        }
251        self.download_ranges(write).await
252    }
253
254    // Repairs what a reorg across the pivot left in the downloaded state, as EIP-8189 describes:
255    // the orphaned blocks' lists schedule what they changed for repair, and catch-up continues from
256    // the last block both branches share. Unserved lists are waited for within the served state
257    // window.
258    async fn recover(&mut self, write: SnapWrite, head: u64) -> Result<Step, SnapSyncError> {
259        let pivot = self.pivot()?;
260        let provider = self.factory.database_provider_ro()?;
261        let Some(reorg) = provider.snap_reorg(write)? else {
262            info!(target: "sync::snap", ?pivot, "Snap pivot was reorged past its kept headers, restarting");
263            return Ok(Step::Restart)
264        };
265        let ancestor = reorg.ancestor();
266        // Lists activate by timestamp, which grows along a chain, so an ancestor committing to one
267        // means every block after it on either branch does too.
268        let committed = provider
269            .sealed_header(ancestor.number)?
270            .is_some_and(|header| header.block_access_list_hash().is_some());
271        if !committed {
272            info!(target: "sync::snap", ?pivot, ?ancestor, "Snap pivot was reorged across block access list activation, restarting");
273            return Ok(Step::Restart)
274        }
275        let resume = provider
276            .catch_up_progress(write)?
277            .ok_or(SnapSyncError::NoCatchUpProgress)?
278            .resume_after(ancestor)
279            .number;
280        if !self.policy.is_catchable_from(resume, head) {
281            info!(target: "sync::snap", ?pivot, ?ancestor, resume, head, "Orphaned block access lists expired, restarting");
282            return Ok(Step::Restart)
283        }
284        let generation = self.policy.select(&provider, head, self.context.finalized())?;
285        // The new pivot must descend from the ancestor, or the new branch is still too short.
286        // Checked first, so waiting for it does not fetch the orphaned lists on every pass.
287        let Some(generation) = generation.filter(|g| g.target().number >= ancestor.number) else {
288            return Ok(Step::Wait)
289        };
290        drop(provider);
291        // A peer lacking side-chain lists says nothing of the others, so another is asked later.
292        let Some(lists) = self.catch_up.orphaned_lists(reorg.orphaned()).await? else {
293            if !self.policy.awaits_orphaned_lists(ancestor.number, head) {
294                info!(target: "sync::snap", ?pivot, ?ancestor, head, "Orphaned block access lists stayed unavailable, restarting");
295                return Ok(Step::Restart)
296            }
297            debug!(target: "sync::snap", ?pivot, "Orphaned block access lists are unavailable");
298            return Ok(Step::Wait)
299        };
300
301        // Scheduling reads every orphaned list, so it runs on the blocking pool.
302        self.commit_blocking(move |provider| {
303            provider.commit_reorg_recovery(write, ancestor, &lists, generation).map(|_| ())
304        })
305        .await?;
306        info!(target: "sync::snap", ?pivot, ?ancestor, to = ?generation.target(), "Recovered snap state from a pivot reorg");
307        Ok(Step::Continue)
308    }
309
310    // Moves a lagging pivot forward, since a few lists cost less than downloading state peers no
311    // longer serve. Returns the write the attempt accepts afterwards.
312    fn advance_pivot(&mut self, write: SnapWrite, head: u64) -> Result<SnapWrite, SnapSyncError> {
313        let provider = self.factory.database_provider_rw()?;
314        let from = self.pivot()?;
315        let Some(advanced) =
316            self.session.advance(&provider, write, head, self.context.finalized())?
317        else {
318            return Ok(write)
319        };
320        // The database takes one writer at a time, so this commits before any download does.
321        provider.commit()?;
322        info!(target: "sync::snap", ?from, to = ?self.pivot()?, "Advanced snap pivot");
323        Ok(advanced)
324    }
325
326    // Applies lists until the downloaded state reaches the pivot. `Some` ends the pass.
327    async fn catch_up(&mut self, write: SnapWrite) -> Result<Option<Step>, SnapSyncError> {
328        let pivot = self.pivot()?.number;
329        loop {
330            match self.catch_up.next(write, pivot).await? {
331                CatchUpStep::Complete => return Ok(None),
332                CatchUpStep::Applied { progress, .. } => {
333                    debug!(target: "sync::snap", applied = ?progress.applied(), pivot, "Applied block access lists");
334                }
335                CatchUpStep::Unavailable { peer_id, .. } => {
336                    debug!(target: "sync::snap", ?peer_id, pivot, "Peer does not serve the pivot's block access lists");
337                    return Ok(Some(Step::Wait))
338                }
339            }
340            if self.cancel.is_cancelled() {
341                return Ok(Some(Step::Stop))
342            }
343        }
344    }
345
346    // Fetches scheduled repairs again at the pivot, with the slots and code they need, committing
347    // each batch. `Some` ends the pass, at the latest after `ranges_per_check` commits so the pivot
348    // is checked again.
349    async fn repair(&mut self) -> Result<Option<Step>, SnapSyncError> {
350        for _ in 0..self.ranges_per_check {
351            let range = match self.accounts.next_repair().await? {
352                None => return Ok(None),
353                Some(AccountRangeStep::Verified(range)) => range,
354                Some(AccountRangeStep::Unavailable { .. }) => return Ok(Some(Step::Wait)),
355            };
356            let Some(slots) = self.storage.repair_slots(&range).await? else {
357                return Ok(Some(Step::Wait))
358            };
359            if let Some(step) = self.download_code(&range).await? {
360                return Ok(Some(step))
361            }
362            let hashed_address = range.origin();
363            let remaining = self.accounts.commit_repair(range, slots).await?;
364            debug!(target: "sync::snap", %hashed_address, remaining, "Committed snap repair batch");
365            if self.cancel.is_cancelled() {
366                return Ok(Some(Step::Stop))
367            }
368        }
369        Ok(Some(Step::Continue))
370    }
371
372    // Commits up to `ranges_per_check` account ranges, handing the state off once none remain.
373    async fn download_ranges(&mut self, write: SnapWrite) -> Result<Step, SnapSyncError> {
374        for _ in 0..self.ranges_per_check {
375            if self.cancel.is_cancelled() {
376                return Ok(Step::Stop)
377            }
378            let range = match self.accounts.next().await? {
379                Some(AccountRangeStep::Verified(range)) => range,
380                Some(AccountRangeStep::Unavailable { origin, peer_id }) => {
381                    debug!(target: "sync::snap", ?peer_id, %origin, "Peer does not serve the pivot state");
382                    return Ok(Step::Wait)
383                }
384                None => return self.hand_off(write).await,
385            };
386            if let Some(step) = self.download_storage_and_code(&range).await? {
387                return Ok(step)
388            }
389            self.accounts.commit(range, Default::default(), Vec::new()).await?;
390        }
391        Ok(Step::Continue)
392    }
393
394    // Persists the storage and code `range` needs. `Some` ends the pass. Both commit as they
395    // arrive, so a range dropped here is fetched again without repeating them.
396    async fn download_storage_and_code(
397        &mut self,
398        range: &VerifiedRange,
399    ) -> Result<Option<Step>, SnapSyncError> {
400        loop {
401            match self.storage.next(range).await? {
402                StorageRangeStep::Complete => break,
403                StorageRangeStep::Committed(_) => {}
404                StorageRangeStep::Unavailable { peer_id, .. } => {
405                    debug!(target: "sync::snap", ?peer_id, "Peer does not serve the pivot's storage");
406                    return Ok(Some(Step::Wait))
407                }
408            }
409            // A large contract takes many responses, each committed, so any of them is a
410            // resumable place to stop.
411            if self.cancel.is_cancelled() {
412                return Ok(Some(Step::Stop))
413            }
414        }
415        self.download_code(range).await
416    }
417
418    // Persists the code `range` references. `Some` ends the pass.
419    async fn download_code(
420        &mut self,
421        range: &VerifiedRange,
422    ) -> Result<Option<Step>, SnapSyncError> {
423        loop {
424            match self.bytecode.next(range).await? {
425                BytecodeStep::Complete => return Ok(None),
426                BytecodeStep::Committed { .. } => {}
427                BytecodeStep::Unavailable { peer_id, .. } => {
428                    debug!(target: "sync::snap", ?peer_id, "Peer does not serve the pivot's code");
429                    return Ok(Some(Step::Wait))
430                }
431            }
432            if self.cancel.is_cancelled() {
433                return Ok(Some(Step::Stop))
434            }
435        }
436    }
437
438    // Checks the downloaded state is complete and hands it to the merkle stage. The scan reads
439    // every account, so it runs on the blocking pool.
440    async fn hand_off(&self, write: SnapWrite) -> Result<Step, SnapSyncError> {
441        let cancel = self.cancel.clone();
442        self.commit_blocking(move |provider| {
443            provider.start_trie_rebuild(write, DEFAULT_SCAN_CHUNK, &cancel)
444        })
445        .await?;
446        Ok(Step::HandedOff(write))
447    }
448
449    // Runs `apply` in one transaction on the blocking pool, committing only if it succeeds.
450    async fn commit_blocking(
451        &self,
452        apply: impl FnOnce(&F::ProviderRW) -> Result<(), SnapSyncError> + Send + 'static,
453    ) -> Result<(), SnapSyncError> {
454        let factory = self.factory.clone();
455        self.runtime
456            .spawn_blocking(move || -> Result<(), SnapSyncError> {
457                let provider = factory.database_provider_rw()?;
458                apply(&provider)?;
459                provider.commit()?;
460                Ok(())
461            })
462            .await
463            .map_err(|error| SnapSyncError::Provider(ProviderError::other(error)))?
464    }
465
466    // Pivot of the attempt being driven.
467    fn pivot(&self) -> Result<BlockNumHash, SnapSyncError> {
468        self.session.target().map(SnapGeneration::target).ok_or(SnapSyncError::NoAttempt)
469    }
470
471    // Whether the context reported progress before the run was cancelled.
472    async fn wait(&mut self, head: u64) -> bool {
473        let Self { cancel, context, .. } = self;
474        cancel.run_until_cancelled(context.wait_for_progress(head)).await.unwrap_or(false)
475    }
476}
477
478impl<C, F, X> fmt::Debug for SnapBootstrap<C, F, X> {
479    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
480        f.debug_struct("SnapBootstrap")
481            .field("policy", &self.policy)
482            .field("session", &self.session)
483            .field("ranges_per_check", &self.ranges_per_check)
484            .finish_non_exhaustive()
485    }
486}
487
488/// How a [`SnapBootstrap`] run ended.
489#[derive(Clone, Copy, Debug, Eq, PartialEq)]
490pub enum SnapBootstrapOutcome {
491    /// The downloaded state is complete and handed to the merkle stage; once the stage reaches
492    /// `pivot`, [`SnapStateVerifier::verify_state_root`] accepts it under `write`.
493    TrieRebuild {
494        /// Write the state was downloaded under.
495        write: SnapWrite,
496        /// Block the state is anchored to.
497        pivot: BlockNumHash,
498    },
499    /// An earlier run's state was already verified, so there is nothing to download.
500    Verified {
501        /// Block the verified state is anchored to.
502        pivot: BlockNumHash,
503    },
504    /// The run was cancelled or the context reported no further progress. Committed progress is
505    /// kept for the next run.
506    Stopped,
507}
508
509/// What a [`SnapBootstrap`] needs to know about the chain and peers it synchronizes from.
510pub trait SnapSyncContext: Send {
511    /// Highest block whose canonical header is downloaded.
512    fn head(&self) -> Result<u64, SnapSyncError>;
513
514    /// Latest finalized block, preferred as the pivot when recent enough.
515    fn finalized(&self) -> Option<u64> {
516        None
517    }
518
519    /// Waits until the head moves past `head` or new peers connect, returning `false` once no
520    /// further progress will come.
521    fn wait_for_progress(&mut self, head: u64) -> impl Future<Output = bool> + Send;
522}
523
524// Whether the attempt being driven can continue, and how.
525enum Resolved {
526    // An unfinished attempt, resumed or just started, accepting this write.
527    Active(SnapWrite),
528    // An earlier run's state was verified at this pivot, so nothing is left to download.
529    Verified(BlockNumHash),
530    // No block is eligible as a pivot yet.
531    Waiting,
532}
533
534// What one pass over the attempt left to do.
535enum Step {
536    HandedOff(SnapWrite),
537    // The range budget ran out; the pivot is checked again before more ranges.
538    Continue,
539    // Peers or headers are missing until the node progresses.
540    Wait,
541    // The attempt's state can no longer be carried forward.
542    Restart,
543    // The run was cancelled; committed progress stays for the next one.
544    Stop,
545}
546
547#[cfg(test)]
548mod tests {
549    use super::*;
550    use crate::{
551        test_utils::{
552            account, account_range, hashed_factory, header, key, policy, state_root,
553            storage_ranges, storage_root_of, stored_slots, verified_range, ReorgFactoryExt,
554            ScriptedSnapClient,
555        },
556        StateRepairs, VerifiedSnapState,
557    };
558    use alloy_eip7928::{
559        compute_block_access_list_hash, AccountChanges, BalanceChange, BlockAccessIndex,
560        NonceChange,
561    };
562    use alloy_primitives::{keccak256, Address, B256, U256};
563    use reth_eth_wire_types::{
564        snap::{AccountRangeMessage, BlockAccessListsMessage},
565        BlockAccessLists,
566    };
567    use reth_network_p2p::{
568        error::{PeerRequestResult, RequestError},
569        snap::client::SnapResponse,
570    };
571    use reth_network_peers::{PeerId, WithPeerId};
572    use reth_primitives_traits::{Account, AlloyBlockHeader, SealedHeader};
573    use reth_provider::{
574        test_utils::{insert_headers, MockNodeTypesWithDB},
575        ProviderFactory,
576    };
577    use reth_stages::stages::MerkleStage;
578    use reth_stages_api::{ExecInput, Stage};
579    use reth_stages_types::StageId;
580    use reth_storage_api::{SnapAttemptId, StageCheckpointReader};
581    use reth_trie_common::{HashedPostState, TrieAccount};
582    use std::{cell::RefCell, collections::VecDeque, sync::Arc};
583
584    type Factory = ProviderFactory<MockNodeTypesWithDB>;
585    type Bootstrap = SnapBootstrap<Arc<ScriptedSnapClient>, Factory, TestContext>;
586
587    const FAR: B256 = B256::repeat_byte(0xaa);
588    // Changed on the orphaned branch, which credited it.
589    const STALE: Address = Address::repeat_byte(0x51);
590    // Untouched by the orphaned branch.
591    const KEPT: Address = Address::repeat_byte(0x52);
592
593    fn accounts() -> Vec<(B256, TrieAccount)> {
594        vec![(key(1), account(1)), (key(2), account(2)), (FAR, account(3))]
595    }
596
597    // Blocks `0..=tip`, each committing to an empty list and to `root`, so any of them anchors
598    // the same state.
599    fn insert_chain(factory: &Factory, tip: u64, root: B256) {
600        insert_headers(factory, &chain(tip, root));
601    }
602
603    // A peer serving `blocks` empty lists.
604    fn empty_lists(request_id: u64, blocks: usize) -> PeerRequestResult<SnapResponse> {
605        lists(request_id, &vec![Vec::new(); blocks], true)
606    }
607
608    // A peer holding no list for the first block it is asked for.
609    fn no_lists(request_id: u64) -> PeerRequestResult<SnapResponse> {
610        lists(request_id, &[Vec::new()], false)
611    }
612
613    fn scripted(
614        factory: &Factory,
615        responses: impl IntoIterator<Item = PeerRequestResult<SnapResponse>>,
616        heads: impl IntoIterator<Item = u64>,
617    ) -> (Arc<ScriptedSnapClient>, Bootstrap) {
618        let client = Arc::new(ScriptedSnapClient::new(responses));
619        let context = TestContext { heads: RefCell::new(heads.into_iter().collect()), waits: 0 };
620        let bootstrap =
621            SnapBootstrap::new(Arc::clone(&client), factory.clone(), Runtime::test(), context)
622                .with_policy(policy());
623        (client, bootstrap)
624    }
625
626    fn attempt_id(factory: &Factory) -> SnapAttemptId {
627        factory.database_provider_ro().unwrap().snap_attempt().unwrap().unwrap().id()
628    }
629
630    // Serves heads in order, repeating the last, and ends the run at the first wait.
631    struct TestContext {
632        heads: RefCell<VecDeque<u64>>,
633        waits: usize,
634    }
635
636    impl SnapSyncContext for TestContext {
637        fn head(&self) -> Result<u64, SnapSyncError> {
638            let mut heads = self.heads.borrow_mut();
639            let head = if heads.len() > 1 { heads.pop_front() } else { heads.front().copied() };
640            Ok(head.expect("the context is given at least one head"))
641        }
642
643        fn wait_for_progress(&mut self, _head: u64) -> impl Future<Output = bool> + Send {
644            self.waits += 1;
645            std::future::ready(false)
646        }
647    }
648
649    // Blocks `0..=tip`, each committing to an empty list and to `root`.
650    fn chain(tip: u64, root: B256) -> Vec<SealedHeader> {
651        let commitment = compute_block_access_list_hash(&Vec::<AccountChanges>::new());
652        let mut parent = B256::ZERO;
653        (0..=tip)
654            .map(|number| {
655                let mut header = header(number, parent, Some(commitment));
656                header.state_root = root;
657                let sealed = SealedHeader::seal_slow(header);
658                parent = sealed.hash();
659                sealed
660            })
661            .collect()
662    }
663
664    // Blocks continuing `parent`, one per list, each committing to its list and to `root`.
665    fn branch(
666        parent: &SealedHeader,
667        lists: &[Vec<AccountChanges>],
668        root: B256,
669    ) -> Vec<SealedHeader> {
670        let mut parent = parent.clone();
671        lists
672            .iter()
673            .map(|list| {
674                let commitment = compute_block_access_list_hash(list);
675                let mut header = header(parent.number() + 1, parent.hash(), Some(commitment));
676                header.state_root = root;
677                parent = SealedHeader::seal_slow(header);
678                parent.clone()
679            })
680            .collect()
681    }
682
683    // A peer serving `lists` in order, or none of them.
684    fn lists(
685        request_id: u64,
686        lists: &[Vec<AccountChanges>],
687        served: bool,
688    ) -> PeerRequestResult<SnapResponse> {
689        let block_access_lists =
690            lists.iter().map(|list| served.then(|| alloy_rlp::encode(list).into())).collect();
691        let message = BlockAccessListsMessage {
692            request_id,
693            block_access_lists: BlockAccessLists(block_access_lists),
694        };
695        Ok(WithPeerId::new(PeerId::random(), SnapResponse::BlockAccessLists(message)))
696    }
697
698    // Rebuilds the trie of `write`'s state up to `target` and checks it against that header.
699    fn rebuild_and_verify(factory: &Factory, write: SnapWrite, target: u64) -> VerifiedSnapState {
700        let provider = factory.database_provider_rw().unwrap();
701        let mut stage = MerkleStage::default_execution();
702        loop {
703            let checkpoint = provider.get_stage_checkpoint(StageId::MerkleExecute).unwrap();
704            let output =
705                stage.execute(&provider, ExecInput { target: Some(target), checkpoint }).unwrap();
706            provider.save_stage_checkpoint(StageId::MerkleExecute, output.checkpoint).unwrap();
707            if output.done {
708                break
709            }
710        }
711        provider.verify_state_root(write).unwrap()
712    }
713
714    // The accounts at every block of the new branch, which never credits `STALE`.
715    fn reorg_accounts() -> Vec<(B256, TrieAccount)> {
716        let mut accounts = vec![(keccak256(STALE), account(1)), (keccak256(KEPT), account(2))];
717        accounts.sort_by_key(|(hashed_address, _)| *hashed_address);
718        accounts
719    }
720
721    fn stale_changes() -> AccountChanges {
722        AccountChanges::new(STALE)
723    }
724
725    // An attempt that downloaded every account at pivot 4 of a branch crediting `STALE` in blocks
726    // 3 and 4, then a reorg to a branch forking after block 2 whose first lists are `new_lists`.
727    // Returns the orphaned headers and the new branch's.
728    fn reorged(
729        factory: &Factory,
730        new_lists: [Vec<AccountChanges>; 2],
731    ) -> (SnapAttemptId, Vec<SealedHeader>, Vec<SealedHeader>) {
732        let accounts = reorg_accounts();
733        let root = state_root(&accounts);
734        let shared = chain(2, root);
735        let orphaned = branch(&shared[2], &orphaned_lists(), root);
736        insert_headers(factory, &shared);
737        insert_headers(factory, &orphaned);
738        let provider = factory.database_provider_rw().unwrap();
739        let write =
740            provider.start_snap_attempt(SnapGeneration::new(orphaned[1].num_hash(), root)).unwrap();
741        provider.start_account_coverage(write).unwrap();
742        let range = verified_range(&accounts, 0..2, B256::ZERO, &[]);
743        provider.commit_account_range(write, &range, Default::default(), Vec::new()).unwrap();
744        // The range was served at the orphaned pivot, holding its credit.
745        let mut credited = Account::from(account(1));
746        credited.balance = U256::from(1_000);
747        let stale = HashedPostState::default().with_accounts([(keccak256(STALE), Some(credited))]);
748        provider.write_hashed_state(&stale.into_sorted()).unwrap();
749        provider.commit().unwrap();
750
751        let [first, second] = new_lists;
752        let new = branch(&shared[2], &[first, second, Vec::new()], root);
753        factory.replace_headers_after(2, &new);
754        (attempt_id(factory), orphaned, new)
755    }
756
757    // The orphaned branch's lists for blocks 3 and 4, crediting `STALE` in each.
758    fn orphaned_lists() -> Vec<Vec<AccountChanges>> {
759        [999, 1_000]
760            .map(|balance| {
761                vec![stale_changes().with_balance_change(BalanceChange::new(
762                    BlockAccessIndex::new(1),
763                    U256::from(balance),
764                ))]
765            })
766            .to_vec()
767    }
768
769    // A peer serving no account range.
770    fn unserved_range(request_id: u64) -> PeerRequestResult<SnapResponse> {
771        let message = AccountRangeMessage { request_id, accounts: Vec::new(), proof: Vec::new() };
772        Ok(WithPeerId::new(PeerId::random(), SnapResponse::AccountRange(message)))
773    }
774
775    // A peer serving the account at `hashed_address` alone.
776    fn served_account(
777        request_id: u64,
778        accounts: &[(B256, TrieAccount)],
779        hashed_address: B256,
780    ) -> PeerRequestResult<SnapResponse> {
781        let index = accounts.iter().position(|(key, _)| *key == hashed_address).unwrap();
782        account_range(request_id, accounts, index..index + 1, &[hashed_address])
783    }
784
785    fn hashes(headers: &[SealedHeader]) -> Vec<B256> {
786        headers.iter().map(SealedHeader::hash).collect()
787    }
788
789    // Runs a bootstrap through a reorg to a branch whose first lists are `new_lists`, returning
790    // the account ranges it fetched again.
791    async fn recover(new_lists: [Vec<AccountChanges>; 2], repaired: bool) -> Vec<B256> {
792        let accounts = reorg_accounts();
793        let factory = hashed_factory();
794        let (attempt, orphaned, new) = reorged(&factory, new_lists.clone());
795        let mut responses = vec![
796            lists(1, &orphaned_lists(), true),
797            // Blocks 3 and 4 of the new branch carry the state to its pivot.
798            lists(2, &new_lists, true),
799        ];
800        if repaired {
801            responses.push(served_account(1, &accounts, keccak256(STALE)));
802        }
803        let (client, mut bootstrap) = scripted(&factory, responses, [5]);
804
805        let outcome = bootstrap.run().await.unwrap();
806
807        let SnapBootstrapOutcome::TrieRebuild { write, pivot } = outcome else {
808            panic!("the state is repaired: {outcome:?}")
809        };
810        assert_eq!(pivot, new[1].num_hash());
811        assert_eq!(attempt_id(&factory), attempt);
812        assert_eq!(*client.block_requests(), [hashes(&orphaned), hashes(&new[..2])]);
813        let origins = client.origins().clone();
814        // The repaired state is the new pivot's, as its header commits to.
815        let verified = rebuild_and_verify(&factory, write, pivot.number);
816        assert_eq!(verified.state_root(), state_root(&accounts));
817        origins
818    }
819
820    #[tokio::test]
821    async fn downloads_the_state_and_hands_it_to_the_trie_rebuild() {
822        let accounts = accounts();
823        let factory = hashed_factory();
824        insert_chain(&factory, 3, state_root(&accounts));
825        let (client, mut bootstrap) =
826            scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
827
828        let outcome = bootstrap.run().await.unwrap();
829
830        let SnapBootstrapOutcome::TrieRebuild { write, pivot } = outcome else {
831            panic!("the state is complete: {outcome:?}")
832        };
833        assert_eq!(pivot.number, 2);
834        assert_eq!(*client.origins(), [B256::ZERO]);
835        let provider = factory.database_provider_ro().unwrap();
836        assert!(provider.is_trie_rebuild_started(write).unwrap());
837    }
838
839    #[tokio::test]
840    async fn pending_repairs_follow_the_pivot_once_the_accounts_are_complete() {
841        let accounts = accounts();
842        let root = state_root(&accounts);
843        let factory = hashed_factory();
844        insert_chain(&factory, 8, root);
845        let provider = factory.database_provider_rw().unwrap();
846        let pivot = provider.sealed_header(2).unwrap().unwrap().num_hash();
847        let write = provider.start_snap_attempt(SnapGeneration::new(pivot, root)).unwrap();
848        provider.start_account_coverage(write).unwrap();
849        let range = verified_range(&accounts, 0..3, B256::ZERO, &[]);
850        provider.commit_account_range(write, &range, Default::default(), Vec::new()).unwrap();
851        let mut repairs = StateRepairs::default();
852        repairs.insert_account(key(2));
853        provider.schedule_snap_repairs(write, repairs).unwrap();
854        provider.commit().unwrap();
855        // No peer serves the repair at pivot 2 any more.
856        let (_, mut bootstrap) = scripted(&factory, [unserved_range(1)], [3]);
857        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
858
859        // Once the head moves on, the pivot advances and the repair is fetched there.
860        let responses = [empty_lists(1, 5), account_range(1, &accounts, 1..2, &[key(2)])];
861        let (client, mut bootstrap) = scripted(&factory, responses, [8]);
862        let outcome = bootstrap.run().await.unwrap();
863
864        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
865            panic!("the state is complete: {outcome:?}")
866        };
867        assert_eq!(pivot.number, 7);
868        assert_eq!(*client.origins(), [key(2)]);
869    }
870
871    #[tokio::test]
872    async fn scheduled_repairs_are_fetched_at_the_pivot_before_ranges() {
873        let accounts = accounts();
874        let factory = hashed_factory();
875        insert_chain(&factory, 3, state_root(&accounts));
876        let provider = factory.database_provider_rw().unwrap();
877        let pivot = provider.sealed_header(2).unwrap().unwrap().num_hash();
878        let write =
879            provider.start_snap_attempt(SnapGeneration::new(pivot, state_root(&accounts))).unwrap();
880        provider.start_account_coverage(write).unwrap();
881        let mut repairs = StateRepairs::default();
882        repairs.insert_account(key(2));
883        provider.schedule_snap_repairs(write, repairs).unwrap();
884        provider.commit().unwrap();
885        // One download fetches the repair and then the ranges, numbering both requests.
886        let responses =
887            [account_range(1, &accounts, 1..2, &[key(2)]), account_range(2, &accounts, 0..3, &[])];
888        let (client, mut bootstrap) = scripted(&factory, responses, [3]);
889
890        let outcome = bootstrap.run().await.unwrap();
891
892        assert!(matches!(outcome, SnapBootstrapOutcome::TrieRebuild { .. }), "{outcome:?}");
893        assert_eq!(*client.origins(), [key(2), B256::ZERO]);
894        let provider = factory.database_provider_ro().unwrap();
895        assert!(provider.snap_repairs(write).unwrap().is_empty());
896    }
897
898    #[tokio::test]
899    async fn repairs_yield_to_pivot_checks_between_batches() {
900        let accounts = accounts();
901        let root = state_root(&accounts);
902        let factory = hashed_factory();
903        insert_chain(&factory, 8, root);
904        let provider = factory.database_provider_rw().unwrap();
905        let pivot = provider.sealed_header(2).unwrap().unwrap().num_hash();
906        let write = provider.start_snap_attempt(SnapGeneration::new(pivot, root)).unwrap();
907        provider.start_account_coverage(write).unwrap();
908        let range = verified_range(&accounts, 0..3, B256::ZERO, &[]);
909        provider.commit_account_range(write, &range, Default::default(), Vec::new()).unwrap();
910        let mut repairs = StateRepairs::default();
911        repairs.insert_account(key(1));
912        repairs.insert_account(key(2));
913        provider.schedule_snap_repairs(write, repairs).unwrap();
914        provider.commit().unwrap();
915        let responses = [
916            account_range(1, &accounts, 0..1, &[key(1)]),
917            // The head moved past the advance window meanwhile, so the pivot moves to block 7.
918            empty_lists(1, 5),
919            account_range(2, &accounts, 1..2, &[key(2)]),
920        ];
921        let (client, bootstrap) = scripted(&factory, responses, [3, 8]);
922        let mut bootstrap = bootstrap.with_ranges_per_check(1);
923
924        let outcome = bootstrap.run().await.unwrap();
925
926        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
927            panic!("the state is repaired: {outcome:?}")
928        };
929        assert_eq!(pivot.number, 7);
930        assert_eq!(*client.origins(), [key(1), key(2)]);
931        assert_eq!(client.block_requests().len(), 1);
932    }
933
934    #[tokio::test]
935    async fn a_lagging_pivot_advances_before_more_ranges_download() {
936        let accounts = accounts();
937        let root = state_root(&accounts);
938        let factory = hashed_factory();
939        insert_chain(&factory, 8, root);
940        let responses = [
941            account_range(1, &accounts, 0..1, &[key(1)]),
942            // Blocks 3 through 7 carry the first account from pivot 2 to pivot 7.
943            empty_lists(1, 5),
944            account_range(2, &accounts, 1..3, &[key(2), FAR]),
945        ];
946        // The head moves past the advance window after the first range.
947        let (client, bootstrap) = scripted(&factory, responses, [3, 8]);
948        let mut bootstrap = bootstrap.with_ranges_per_check(1);
949
950        let outcome = bootstrap.run().await.unwrap();
951
952        let SnapBootstrapOutcome::TrieRebuild { write, pivot } = outcome else {
953            panic!("the state is complete: {outcome:?}")
954        };
955        assert_eq!(pivot.number, 7);
956        assert_eq!(*client.origins(), [B256::ZERO, key(2)]);
957        assert_eq!(client.block_requests().len(), 1);
958
959        // The merkle stage rebuilds the trie, which the pivot's header then accepts.
960        let verified = rebuild_and_verify(&factory, write, 7);
961        assert_eq!((verified.target(), verified.state_root()), (pivot, root));
962    }
963
964    #[tokio::test]
965    async fn peers_without_the_state_wait_and_the_next_run_resumes() {
966        let accounts = accounts();
967        let factory = hashed_factory();
968        insert_chain(&factory, 3, state_root(&accounts));
969        let responses = [
970            account_range(1, &accounts, 0..1, &[key(1)]),
971            Err(RequestError::UnsupportedCapability),
972        ];
973        let (_, mut bootstrap) = scripted(&factory, responses, [3]);
974
975        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
976        assert_eq!(bootstrap.context.waits, 1);
977        let attempt = attempt_id(&factory);
978
979        let (client, mut resumed) =
980            scripted(&factory, [account_range(1, &accounts, 1..3, &[key(2), FAR])], [3]);
981        let outcome = resumed.run().await.unwrap();
982
983        assert!(matches!(outcome, SnapBootstrapOutcome::TrieRebuild { .. }));
984        assert_eq!(attempt_id(&factory), attempt);
985        assert_eq!(*client.origins(), [key(2)]);
986    }
987
988    #[tokio::test]
989    async fn no_eligible_pivot_waits_without_starting_an_attempt() {
990        let factory = hashed_factory();
991        insert_chain(&factory, 0, state_root(&accounts()));
992        let (client, mut bootstrap) = scripted(&factory, [], [0]);
993
994        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
995
996        assert_eq!(bootstrap.context.waits, 1);
997        assert!(client.origins().is_empty());
998        assert!(factory.database_provider_ro().unwrap().snap_attempt().unwrap().is_none());
999    }
1000
1001    #[tokio::test]
1002    async fn a_reorged_pivot_without_kept_headers_restarts_the_attempt() {
1003        let accounts = accounts();
1004        let factory = hashed_factory();
1005        insert_chain(&factory, 3, state_root(&accounts));
1006        // A previous run anchored to a block the canonical chain no longer holds.
1007        let provider = factory.database_provider_rw().unwrap();
1008        let orphaned = SnapGeneration::new(
1009            BlockNumHash::new(2, B256::repeat_byte(0xee)),
1010            state_root(&accounts),
1011        );
1012        provider.start_snap_attempt(orphaned).unwrap();
1013        provider.commit().unwrap();
1014        let orphaned = attempt_id(&factory);
1015        let (_, mut bootstrap) = scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
1016
1017        let outcome = bootstrap.run().await.unwrap();
1018
1019        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1020            panic!("the state is complete: {outcome:?}")
1021        };
1022        assert_ne!(pivot.hash, B256::repeat_byte(0xee));
1023        assert_ne!(attempt_id(&factory), orphaned);
1024    }
1025
1026    #[tokio::test]
1027    async fn a_handed_off_attempt_is_not_downloaded_again() {
1028        let accounts = accounts();
1029        let factory = hashed_factory();
1030        insert_chain(&factory, 3, state_root(&accounts));
1031        let (_, mut bootstrap) = scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
1032        let first = bootstrap.run().await.unwrap();
1033
1034        let (client, mut resumed) = scripted(&factory, [], [3]);
1035
1036        assert_eq!(resumed.run().await.unwrap(), first);
1037        assert!(client.origins().is_empty());
1038    }
1039
1040    #[tokio::test]
1041    async fn a_cancelled_run_stops_before_any_work() {
1042        let factory = hashed_factory();
1043        insert_chain(&factory, 3, state_root(&accounts()));
1044        let cancel = CancellationToken::new();
1045        let (client, bootstrap) = scripted(&factory, [], [3]);
1046        let mut bootstrap = bootstrap.with_cancellation(cancel.clone());
1047        cancel.cancel();
1048
1049        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1050        assert!(client.origins().is_empty());
1051    }
1052
1053    #[tokio::test]
1054    async fn a_catch_up_stalled_past_the_served_lists_restarts() {
1055        let accounts = accounts();
1056        let factory = hashed_factory();
1057        insert_chain(&factory, 11, state_root(&accounts));
1058        // The pivot moves from 2 to 7, but no peer serves block 3's list.
1059        let responses = [account_range(1, &accounts, 0..1, &[key(1)]), no_lists(1)];
1060        let (_, stalled) = scripted(&factory, responses, [3, 8]);
1061        let mut stalled = stalled.with_ranges_per_check(1);
1062        assert_eq!(stalled.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1063        let attempt = attempt_id(&factory);
1064
1065        // Pivot 7 is still recent under head 11, but block 3's list is past the served history.
1066        let (client, mut restarted) =
1067            scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [11]);
1068        let outcome = restarted.run().await.unwrap();
1069
1070        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1071            panic!("the state is complete: {outcome:?}")
1072        };
1073        assert_eq!(pivot.number, 10);
1074        assert_ne!(attempt_id(&factory), attempt);
1075        assert!(client.block_requests().is_empty());
1076    }
1077
1078    #[tokio::test]
1079    async fn complete_state_is_handed_off_past_the_served_lists() {
1080        let accounts = accounts();
1081        let factory = hashed_factory();
1082        insert_chain(&factory, 11, state_root(&accounts));
1083        // A run committed every range at pivot 2, then stopped before the hand-off.
1084        let provider = factory.database_provider_rw().unwrap();
1085        let pivot = provider.sealed_header(2).unwrap().unwrap();
1086        let generation = SnapGeneration::new(pivot.num_hash(), pivot.state_root);
1087        let downloaded = provider.start_snap_attempt(generation).unwrap();
1088        provider.start_account_coverage(downloaded).unwrap();
1089        let range = verified_range(&accounts, 0..3, B256::ZERO, &[]);
1090        provider.commit_account_range(downloaded, &range, Default::default(), Vec::new()).unwrap();
1091        provider.commit().unwrap();
1092
1093        // Block 3's list is past the served history under head 11, but nothing needs it.
1094        let (client, mut resumed) = scripted(&factory, [], [11]);
1095        let outcome = resumed.run().await.unwrap();
1096
1097        assert_eq!(
1098            outcome,
1099            SnapBootstrapOutcome::TrieRebuild { write: downloaded, pivot: pivot.num_hash() }
1100        );
1101        assert!(client.origins().is_empty());
1102        assert!(client.block_requests().is_empty());
1103    }
1104
1105    #[tokio::test]
1106    async fn unreadable_progress_restarts_the_attempt() {
1107        let accounts = accounts();
1108        let factory = hashed_factory();
1109        insert_chain(&factory, 3, state_root(&accounts));
1110        let (_, mut stopped) = scripted(&factory, [Err(RequestError::UnsupportedCapability)], [3]);
1111        assert_eq!(stopped.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1112        // Another build rewrote the coverage the attempt resumes from.
1113        let provider = factory.database_provider_rw().unwrap();
1114        let stale = provider.active_snap_write().unwrap().unwrap();
1115        provider.write_metadata("snap_account_coverage", br#"{"version":999}"#.to_vec()).unwrap();
1116        provider.commit().unwrap();
1117
1118        let (_, mut restarted) = scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
1119        let outcome = restarted.run().await.unwrap();
1120
1121        let SnapBootstrapOutcome::TrieRebuild { write, .. } = outcome else {
1122            panic!("the state is complete: {outcome:?}")
1123        };
1124        assert_ne!(write.attempt(), stale.attempt());
1125        let provider = factory.database_provider_ro().unwrap();
1126        assert!(matches!(
1127            provider.authorize_snap_write(stale),
1128            Err(SnapSyncError::StaleWrite { .. })
1129        ));
1130    }
1131
1132    #[tokio::test]
1133    async fn cancellation_between_storage_responses_stops_the_run() {
1134        let slots = vec![(key(1), U256::from(11)), (key(2), U256::from(12))];
1135        let mut contract = account(1);
1136        contract.storage_root = storage_root_of(&slots);
1137        let accounts = vec![(key(1), contract)];
1138        let factory = hashed_factory();
1139        insert_chain(&factory, 3, state_root(&accounts));
1140        let cancel = CancellationToken::new();
1141        let on_request = cancel.clone();
1142        // The contract's storage takes two responses, and shutdown fires during the first.
1143        let client = Arc::new(
1144            ScriptedSnapClient::new([
1145                account_range(1, &accounts, 0..1, &[]),
1146                storage_ranges(1, &[&slots[..1]], &slots, &[B256::ZERO, key(1)]),
1147            ])
1148            .on_storage_request(move || on_request.cancel()),
1149        );
1150        let context = TestContext { heads: RefCell::new(VecDeque::from([3])), waits: 0 };
1151        let mut bootstrap =
1152            SnapBootstrap::new(Arc::clone(&client), factory.clone(), Runtime::test(), context)
1153                .with_policy(policy())
1154                .with_cancellation(cancel);
1155
1156        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1157        assert_eq!(client.storage_requests().len(), 1);
1158    }
1159
1160    #[tokio::test]
1161    async fn a_new_branch_too_short_for_a_pivot_waits_before_fetching_lists() {
1162        let factory = hashed_factory();
1163        let (attempt, _, _) = reorged(&factory, [Vec::new(), Vec::new()]);
1164        // Head 2 puts the pivot at block 1, below the ancestor at block 2.
1165        let (client, mut bootstrap) = scripted(&factory, [], [2]);
1166
1167        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1168
1169        assert!(client.block_requests().is_empty());
1170        assert_eq!(attempt_id(&factory), attempt);
1171    }
1172
1173    #[tokio::test]
1174    async fn a_handed_off_pivot_reorged_before_verification_restarts_the_attempt() {
1175        let accounts = accounts();
1176        let root = state_root(&accounts);
1177        let factory = hashed_factory();
1178        let shared = chain(3, root);
1179        insert_headers(&factory, &shared);
1180        let (_, mut bootstrap) = scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
1181        let SnapBootstrapOutcome::TrieRebuild { pivot: orphaned, .. } =
1182            bootstrap.run().await.unwrap()
1183        else {
1184            panic!("the state is complete")
1185        };
1186        let attempt = attempt_id(&factory);
1187        // Blocks 2 and 3 are replaced before the rebuild verifies the orphaned pivot.
1188        let new = branch(&shared[1], &[vec![stale_changes()], Vec::new()], root);
1189        factory.replace_headers_after(1, &new);
1190        let (_, mut resumed) = scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [3]);
1191
1192        let outcome = resumed.run().await.unwrap();
1193
1194        let SnapBootstrapOutcome::TrieRebuild { write, pivot } = outcome else {
1195            panic!("the new branch's state is complete: {outcome:?}")
1196        };
1197        assert_ne!(pivot, orphaned);
1198        assert_eq!(pivot, new[0].num_hash());
1199        assert_ne!(attempt_id(&factory), attempt);
1200        let verified = rebuild_and_verify(&factory, write, pivot.number);
1201        assert_eq!(verified.state_root(), root);
1202    }
1203
1204    #[tokio::test]
1205    async fn a_field_only_the_orphaned_branch_changed_is_fetched_again() {
1206        let origins = recover([Vec::new(), Vec::new()], true).await;
1207
1208        // The untouched account is kept, never fetched again.
1209        assert_eq!(origins, [keccak256(STALE)]);
1210    }
1211
1212    #[tokio::test]
1213    async fn another_field_changed_on_the_new_branch_leaves_the_account_to_repair() {
1214        let nonce =
1215            stale_changes().with_nonce_change(NonceChange::new(BlockAccessIndex::new(1), 1));
1216
1217        let origins = recover([vec![nonce], Vec::new()], true).await;
1218
1219        assert_eq!(origins, [keccak256(STALE)]);
1220    }
1221
1222    #[tokio::test]
1223    async fn a_field_the_new_branch_changes_too_needs_no_repair() {
1224        let balance = stale_changes()
1225            .with_balance_change(BalanceChange::new(BlockAccessIndex::new(1), U256::from(1)));
1226
1227        let origins = recover([vec![balance], Vec::new()], false).await;
1228
1229        assert!(origins.is_empty());
1230    }
1231
1232    #[tokio::test]
1233    async fn a_second_reorg_repairs_what_the_first_new_branch_changed() {
1234        let accounts = reorg_accounts();
1235        let factory = hashed_factory();
1236        let kept = AccountChanges::new(KEPT)
1237            .with_balance_change(BalanceChange::new(BlockAccessIndex::new(1), U256::from(7)));
1238        let first_lists = [vec![kept], Vec::new()];
1239        let (attempt, _, first) = reorged(&factory, first_lists.clone());
1240        // Catch-up applies the first new branch, then no peer serves the repair at its pivot.
1241        let responses =
1242            [lists(1, &orphaned_lists(), true), lists(2, &first_lists, true), unserved_range(1)];
1243        let (_, mut bootstrap) = scripted(&factory, responses, [5]);
1244        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1245
1246        // After a restart, a second reorg orphans the re-anchored pivot too.
1247        let shared = factory.database_provider_ro().unwrap().sealed_header(2).unwrap().unwrap();
1248        let second = branch(&shared, &[Vec::new(), Vec::new(), Vec::new()], state_root(&accounts));
1249        factory.replace_headers_after(2, &second);
1250        let mut responses =
1251            vec![lists(1, &first_lists, true), lists(2, &[Vec::new(), Vec::new()], true)];
1252        // Both the first orphaned branch's and the first new branch's changes are fetched again.
1253        let mut repaired = [keccak256(STALE), keccak256(KEPT)];
1254        repaired.sort();
1255        for (id, hashed_address) in (1..).zip(repaired) {
1256            responses.push(served_account(id, &accounts, hashed_address));
1257        }
1258        let (client, mut bootstrap) = scripted(&factory, responses, [5]);
1259
1260        let outcome = bootstrap.run().await.unwrap();
1261
1262        let SnapBootstrapOutcome::TrieRebuild { write, pivot } = outcome else {
1263            panic!("the state is repaired: {outcome:?}")
1264        };
1265        assert_eq!(pivot, second[1].num_hash());
1266        assert_eq!(attempt_id(&factory), attempt);
1267        assert_eq!(client.block_requests()[0], hashes(&first[..2]));
1268        assert_eq!(*client.origins(), repaired);
1269        let verified = rebuild_and_verify(&factory, write, pivot.number);
1270        assert_eq!(verified.state_root(), state_root(&accounts));
1271    }
1272
1273    #[tokio::test]
1274    async fn unserved_orphaned_lists_are_waited_for() {
1275        let accounts = reorg_accounts();
1276        let factory = hashed_factory();
1277        let (attempt, orphaned, new) = reorged(&factory, [Vec::new(), Vec::new()]);
1278        let (_, mut bootstrap) = scripted(&factory, [lists(1, &orphaned_lists(), false)], [5]);
1279        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1280        assert_eq!(attempt_id(&factory), attempt);
1281
1282        // A restarted node resumes the orphaned attempt, and another peer serves them.
1283        let responses = [
1284            lists(1, &orphaned_lists(), true),
1285            lists(2, &[Vec::new(), Vec::new()], true),
1286            served_account(1, &accounts, keccak256(STALE)),
1287        ];
1288        let (client, mut bootstrap) = scripted(&factory, responses, [5]);
1289        let outcome = bootstrap.run().await.unwrap();
1290
1291        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1292            panic!("the state is repaired: {outcome:?}")
1293        };
1294        assert_eq!(pivot, new[1].num_hash());
1295        assert_eq!(attempt_id(&factory), attempt);
1296        assert_eq!(client.block_requests()[0], hashes(&orphaned));
1297    }
1298
1299    #[tokio::test]
1300    async fn a_reorg_across_block_access_list_activation_restarts_the_attempt() {
1301        let accounts = accounts();
1302        let root = state_root(&accounts);
1303        let factory = hashed_factory();
1304        let mut shared = chain(1, root);
1305        // The ancestor predates block access lists, so the branches may hold blocks without them.
1306        let mut ancestor = header(2, shared[1].hash(), None);
1307        ancestor.state_root = root;
1308        shared.push(SealedHeader::seal_slow(ancestor));
1309        let orphaned = branch(&shared[2], &orphaned_lists(), root);
1310        insert_headers(&factory, &shared);
1311        insert_headers(&factory, &orphaned);
1312        let provider = factory.database_provider_rw().unwrap();
1313        provider.start_snap_attempt(SnapGeneration::new(orphaned[1].num_hash(), root)).unwrap();
1314        provider.commit().unwrap();
1315        let attempt = attempt_id(&factory);
1316        let new = branch(&shared[2], &[vec![stale_changes()], Vec::new(), Vec::new()], root);
1317        factory.replace_headers_after(2, &new);
1318        let (client, mut bootstrap) =
1319            scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [5]);
1320
1321        let outcome = bootstrap.run().await.unwrap();
1322
1323        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1324            panic!("the new branch's state is complete: {outcome:?}")
1325        };
1326        assert_eq!(pivot, new[1].num_hash());
1327        assert_ne!(attempt_id(&factory), attempt);
1328        assert!(client.block_requests().is_empty());
1329    }
1330
1331    #[tokio::test]
1332    async fn a_catch_up_stalled_below_the_ancestor_restarts_before_fetching_lists() {
1333        let accounts = accounts();
1334        let root = state_root(&accounts);
1335        let factory = hashed_factory();
1336        let shared = chain(2, root);
1337        let orphaned = branch(&shared[2], &orphaned_lists(), root);
1338        insert_headers(&factory, &shared);
1339        insert_headers(&factory, &orphaned);
1340        // The pivot advanced to block 4 while catch-up stayed at block 1, below the ancestor.
1341        let provider = factory.database_provider_rw().unwrap();
1342        let write =
1343            provider.start_snap_attempt(SnapGeneration::new(shared[1].num_hash(), root)).unwrap();
1344        provider
1345            .advance_snap_pivot(write, SnapGeneration::new(orphaned[1].num_hash(), root))
1346            .unwrap();
1347        provider.commit().unwrap();
1348        let attempt = attempt_id(&factory);
1349        let mut lists = vec![vec![stale_changes()]];
1350        lists.resize(8, Vec::new());
1351        let new = branch(&shared[2], &lists, root);
1352        factory.replace_headers_after(2, &new);
1353        // Head 10 still serves the list after the ancestor, but not the one after block 1.
1354        let (client, mut bootstrap) =
1355            scripted(&factory, [account_range(1, &accounts, 0..3, &[])], [10]);
1356
1357        let outcome = bootstrap.run().await.unwrap();
1358
1359        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1360            panic!("the new branch's state is complete: {outcome:?}")
1361        };
1362        assert_eq!(pivot, new[6].num_hash());
1363        assert_ne!(attempt_id(&factory), attempt);
1364        assert!(client.block_requests().is_empty());
1365    }
1366
1367    #[tokio::test]
1368    async fn orphaned_lists_unserved_past_the_state_window_restart_the_attempt() {
1369        let accounts = accounts();
1370        let root = state_root(&accounts);
1371        let factory = hashed_factory();
1372        let shared = chain(2, root);
1373        let orphaned = branch(&shared[2], &orphaned_lists(), root);
1374        insert_headers(&factory, &shared);
1375        insert_headers(&factory, &orphaned);
1376        let provider = factory.database_provider_rw().unwrap();
1377        provider.start_snap_attempt(SnapGeneration::new(orphaned[1].num_hash(), root)).unwrap();
1378        provider.commit().unwrap();
1379        let attempt = attempt_id(&factory);
1380        let mut new_lists = vec![vec![stale_changes()]];
1381        new_lists.resize(129, Vec::new());
1382        let new = branch(&shared[2], &new_lists, root);
1383        factory.replace_headers_after(2, &new);
1384        // Head 131 leaves the ancestor past the served state window, but its lists are still
1385        // served.
1386        let responses =
1387            [lists(1, &orphaned_lists(), false), account_range(1, &accounts, 0..3, &[])];
1388        let (_, bootstrap) = scripted(&factory, responses, [131]);
1389        let mut bootstrap = bootstrap.with_policy(policy().with_history(256));
1390
1391        let outcome = bootstrap.run().await.unwrap();
1392
1393        let SnapBootstrapOutcome::TrieRebuild { pivot, .. } = outcome else {
1394            panic!("the new branch's state is complete: {outcome:?}")
1395        };
1396        assert_eq!(pivot, new[127].num_hash());
1397        assert_ne!(attempt_id(&factory), attempt);
1398    }
1399
1400    #[tokio::test]
1401    async fn a_cancelled_repair_resumes_at_its_pending_slots() {
1402        let slots = vec![(key(1), U256::from(11)), (key(2), U256::from(12))];
1403        let mut contract = account(1);
1404        contract.storage_root = storage_root_of(&slots);
1405        let accounts = vec![(key(1), contract)];
1406        let factory = hashed_factory();
1407        insert_chain(&factory, 3, state_root(&accounts));
1408        let provider = factory.database_provider_rw().unwrap();
1409        let pivot = provider.sealed_header(2).unwrap().unwrap().num_hash();
1410        let write =
1411            provider.start_snap_attempt(SnapGeneration::new(pivot, state_root(&accounts))).unwrap();
1412        provider.start_account_coverage(write).unwrap();
1413        let mut repairs = StateRepairs::default();
1414        repairs.insert_slot(key(1), key(1));
1415        repairs.insert_slot(key(1), key(2));
1416        provider.schedule_snap_repairs(write, repairs).unwrap();
1417        provider.commit().unwrap();
1418        let cancel = CancellationToken::new();
1419        let on_request = cancel.clone();
1420        // Cancellation is checked between repair batches, so the batch ends first, with
1421        // `key(2)` left pending by a peer that does not serve it.
1422        let client = Arc::new(
1423            ScriptedSnapClient::new([
1424                account_range(1, &accounts, 0..1, &[key(1)]),
1425                storage_ranges(1, &[&slots[..1]], &slots, &[key(1)]),
1426                storage_ranges(2, &[], &[], &[]),
1427            ])
1428            .on_storage_request(move || on_request.cancel()),
1429        );
1430        let context = TestContext { heads: RefCell::new(VecDeque::from([3])), waits: 0 };
1431        let mut bootstrap =
1432            SnapBootstrap::new(Arc::clone(&client), factory.clone(), Runtime::test(), context)
1433                .with_policy(policy())
1434                .with_cancellation(cancel);
1435        assert_eq!(bootstrap.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1436        assert_eq!(client.storage_requests().len(), 2);
1437        let provider = factory.database_provider_ro().unwrap();
1438        let pending = provider.snap_repairs(write).unwrap().slots(key(1)).collect::<Vec<_>>();
1439        assert_eq!(pending, [key(2)]);
1440        assert_eq!(stored_slots(&provider, key(1)), slots[..1]);
1441        drop(provider);
1442
1443        let (client, mut resumed) = scripted(
1444            &factory,
1445            [
1446                account_range(1, &accounts, 0..1, &[key(1)]),
1447                storage_ranges(1, &[&slots[1..]], &slots, &[key(2)]),
1448                Err(RequestError::UnsupportedCapability),
1449            ],
1450            [3],
1451        );
1452        assert_eq!(resumed.run().await.unwrap(), SnapBootstrapOutcome::Stopped);
1453        assert_eq!(*client.storage_requests(), [(vec![key(1)], key(2))]);
1454        let provider = factory.database_provider_ro().unwrap();
1455        assert!(provider.snap_repairs(write).unwrap().is_empty());
1456        assert_eq!(stored_slots(&provider, key(1)), slots);
1457    }
1458}