Skip to main content

reth_cli_commands/download/
archive.rs

1use super::{
2    extract::{extract_archive_raw, streaming_download_and_extract, CompressionFormat},
3    fetch::ArchiveFetcher,
4    manifest::SnapshotArchive,
5    planning::{PlannedArchive, PlannedDownloads},
6    progress::{
7        spawn_progress_display, ArchiveDownloadProgress, ArchiveExtractionProgress,
8        ArchiveVerificationProgress, DownloadRequestLimiter, SharedProgress,
9    },
10    session::{ArchiveProcessContext, DownloadSession},
11    verify::OutputVerifier,
12    MAX_DOWNLOAD_RETRIES, RETRY_BACKOFF_SECS,
13};
14use eyre::Result;
15use futures::stream::{self, StreamExt};
16use reth_cli_util::cancellation::CancellationToken;
17use reth_fs_util as fs;
18use std::{path::Path, sync::Arc, time::Duration};
19use tokio::task;
20use tracing::{debug, info, warn};
21
22const DOWNLOAD_CACHE_DIR: &str = ".download-cache";
23
24/// Runs all planned modular archive downloads for one command invocation.
25pub(crate) async fn run_modular_downloads(
26    planned_downloads: PlannedDownloads,
27    target_dir: &Path,
28    static_files_dir: Option<&Path>,
29    download_concurrency: usize,
30    cancel_token: CancellationToken,
31    retry_backoff: Option<Duration>,
32) -> Result<()> {
33    let download_cache_dir = target_dir.join(DOWNLOAD_CACHE_DIR);
34    fs::create_dir_all(&download_cache_dir)?;
35
36    let shared = SharedProgress::new(
37        planned_downloads.total_download_size,
38        planned_downloads.total_output_size,
39        planned_downloads.total_archives() as u64,
40        cancel_token.clone(),
41    );
42    let session = DownloadSession::new(
43        Some(Arc::clone(&shared)),
44        Some(DownloadRequestLimiter::new(download_concurrency)),
45        cancel_token,
46    )
47    .with_retry_backoff(retry_backoff);
48    let ctx = ArchiveProcessContext::new(
49        target_dir.to_path_buf(),
50        static_files_dir.map(Path::to_path_buf),
51        Some(download_cache_dir),
52        session,
53    );
54
55    ModularDownloadJob::new(ctx, download_concurrency).run(planned_downloads).await
56}
57
58/// Schedules modular archive work for one run of `reth download`.
59struct ModularDownloadJob {
60    /// Shared paths and session state for each archive in this job.
61    ctx: ArchiveProcessContext,
62    /// Maximum number of archives processed at once.
63    archive_concurrency: usize,
64}
65
66impl ModularDownloadJob {
67    /// Creates the modular download job for one command run.
68    const fn new(ctx: ArchiveProcessContext, archive_concurrency: usize) -> Self {
69        Self { ctx, archive_concurrency }
70    }
71
72    /// Runs all planned archives and waits for the shared progress task to finish.
73    async fn run(self, planned_downloads: PlannedDownloads) -> Result<()> {
74        let shared = Arc::clone(
75            self.ctx.session().progress().expect("modular downloads always use shared progress"),
76        );
77        let progress_handle = spawn_progress_display(Arc::clone(&shared));
78        let ctx = self.ctx.clone();
79        let results: Vec<Result<()>> = stream::iter(planned_downloads.archives)
80            .map(move |archive| {
81                let ctx = ctx.clone();
82                async move { Self::process_archive(ctx, archive).await }
83            })
84            .buffer_unordered(self.archive_concurrency)
85            .collect()
86            .await;
87
88        shared.done.notify_one();
89        let _ = progress_handle.await;
90
91        for result in results {
92            result?;
93        }
94
95        Ok(())
96    }
97
98    /// Runs one archive on the blocking pool so fetch and extraction stay off the async executor.
99    async fn process_archive(ctx: ArchiveProcessContext, archive: PlannedArchive) -> Result<()> {
100        task::spawn_blocking(move || ArchiveProcessor::new(archive, ctx).run()).await??;
101        Ok(())
102    }
103}
104
105/// Explicit retry states for one modular archive.
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107enum ArchiveAttemptState {
108    /// Start or restart one full archive attempt.
109    RunAttempt,
110    /// Check whether the extracted outputs verify.
111    VerifyOutputs,
112    /// Wait and decide whether another full attempt should run.
113    RetryAttempt,
114    /// Finish successfully.
115    Complete,
116    /// Stop with an error after retries are exhausted.
117    Fail,
118}
119
120/// Processes one modular archive from reuse check through extraction and verification.
121struct ArchiveProcessor {
122    /// The concrete archive and component being processed.
123    archive: PlannedArchive,
124    /// Shared paths and session state for this archive attempt.
125    ctx: ArchiveProcessContext,
126}
127
128impl ArchiveProcessor {
129    /// Creates a processor for one archive and the shared download context.
130    fn new(archive: PlannedArchive, ctx: ArchiveProcessContext) -> Self {
131        Self { archive, ctx }
132    }
133
134    /// Runs the archive retry state machine until outputs are verified or retries are exhausted.
135    fn run(self) -> Result<()> {
136        let archive = self.archive();
137        if self.try_reuse_outputs()? {
138            info!(target: "reth::cli", file = %archive.file_name, component = %self.archive.component, "Skipping already verified plain files");
139            return Ok(());
140        }
141
142        let mode = ArchiveMode::new(&self.ctx)?;
143        let format = CompressionFormat::from_url(&archive.file_name)?;
144        let mut attempt = 1;
145        let mut last_error: Option<eyre::Error> = None;
146        let mut state = ArchiveAttemptState::RunAttempt;
147
148        loop {
149            match state {
150                ArchiveAttemptState::RunAttempt => {
151                    self.cleanup_outputs();
152
153                    if attempt > 1 {
154                        info!(target: "reth::cli",
155                            file = %archive.file_name,
156                            component = %self.archive.component,
157                            attempt,
158                            max = MAX_DOWNLOAD_RETRIES,
159                            "Retrying archive from scratch"
160                        );
161                    }
162
163                    match self.run_attempt(mode, format) {
164                        Ok(()) => state = ArchiveAttemptState::VerifyOutputs,
165                        Err(error) if mode.retries_fetch_errors() => {
166                            warn!(target: "reth::cli",
167                                file = %archive.file_name,
168                                component = %self.archive.component,
169                                attempt,
170                                err = %format_args!("{error:#}"),
171                                "Archive attempt failed, retrying from scratch"
172                            );
173                            last_error = Some(error);
174                            state = ArchiveAttemptState::RetryAttempt;
175                        }
176                        Err(error) => return Err(error),
177                    }
178                }
179                ArchiveAttemptState::VerifyOutputs => {
180                    if self.verify_outputs_with_progress()? {
181                        state = ArchiveAttemptState::Complete;
182                    } else {
183                        warn!(target: "reth::cli", file = %archive.file_name, component = %self.archive.component, attempt, "Archive extracted, but output verification failed, retrying");
184                        state = ArchiveAttemptState::RetryAttempt;
185                    }
186                }
187                ArchiveAttemptState::RetryAttempt => {
188                    if attempt >= MAX_DOWNLOAD_RETRIES {
189                        state = ArchiveAttemptState::Fail;
190                    } else {
191                        std::thread::sleep(
192                            self.ctx.session().retry_delay(Duration::from_secs(RETRY_BACKOFF_SECS)),
193                        );
194                        attempt += 1;
195                        state = ArchiveAttemptState::RunAttempt;
196                    }
197                }
198                ArchiveAttemptState::Complete => return Ok(()),
199                ArchiveAttemptState::Fail => {
200                    if let Some(error) = last_error {
201                        return Err(error.wrap_err(format!(
202                            "Failed after {} attempts for {}",
203                            MAX_DOWNLOAD_RETRIES, archive.file_name
204                        )));
205                    }
206
207                    eyre::bail!(
208                        "Failed integrity validation after {} attempts for {}",
209                        MAX_DOWNLOAD_RETRIES,
210                        archive.file_name
211                    );
212                }
213            }
214        }
215    }
216
217    /// Returns the concrete archive being fetched or verified.
218    fn archive(&self) -> &SnapshotArchive {
219        &self.archive.archive
220    }
221
222    /// Returns the verifier for this archive's output files.
223    fn output_verifier(&self) -> OutputVerifier<'_> {
224        OutputVerifier::new(self.ctx.target_dir(), self.ctx.static_files_dir())
225    }
226
227    /// Returns `true` if this archive can be reused from existing verified outputs.
228    /// Returns `false` if a fresh archive attempt is still needed.
229    fn try_reuse_outputs(&self) -> Result<bool> {
230        if self.verify_outputs()? {
231            self.mark_complete();
232            return Ok(true);
233        }
234
235        Ok(false)
236    }
237
238    /// Removes any partial outputs before a fresh archive attempt.
239    fn cleanup_outputs(&self) {
240        self.output_verifier().cleanup(&self.archive().output_files);
241    }
242
243    /// Returns `true` if all declared plain outputs verify.
244    /// Returns `false` if any output is missing or does not match.
245    fn verify_outputs(&self) -> Result<bool> {
246        self.output_verifier().verify(&self.archive().output_files)
247    }
248
249    /// Records archive completion in shared progress once outputs verify.
250    fn mark_complete(&self) {
251        self.ctx.session().record_reused_archive(self.archive().size, self.archive().output_size());
252    }
253
254    /// Executes one archive attempt according to the selected cache-vs-stream mode.
255    fn run_attempt(&self, mode: ArchiveMode, format: CompressionFormat) -> Result<()> {
256        mode.execute(self, format)
257    }
258
259    /// Downloads the archive into the cache, then extracts from the cached file.
260    fn run_cached_attempt(&self, format: CompressionFormat) -> Result<()> {
261        let cache_dir =
262            self.ctx.cache_dir().ok_or_else(|| eyre::eyre!("Missing download cache directory"))?;
263        let fetcher =
264            ArchiveFetcher::new(self.archive().url.clone(), cache_dir, self.ctx.session().clone());
265
266        if self.archive.ty == super::manifest::SnapshotComponentType::State {
267            debug!(target: "reth::cli", url = %self.archive().url, "Downloading state snapshot archive");
268        }
269
270        let download_result = {
271            let mut download_progress = ArchiveDownloadProgress::new(self.ctx.session().progress());
272            let result = fetcher.download(Some(&mut download_progress));
273            if let Ok(ref downloaded) = result &&
274                download_progress.has_tracked_bytes()
275            {
276                download_progress.complete(downloaded.size);
277            }
278            result
279        };
280
281        let downloaded = match download_result {
282            Ok(downloaded) => downloaded,
283            Err(error) => {
284                fetcher.cleanup_downloaded_files();
285                return Err(error);
286            }
287        };
288
289        info!(target: "reth::cli",
290            file = %self.archive().file_name,
291            component = %self.archive.component,
292            size = %super::progress::DownloadProgress::format_size(downloaded.size),
293            "Archive download complete"
294        );
295
296        let extract_result = self.extract_cached_archive(&downloaded.path, format);
297        fetcher.cleanup_downloaded_files();
298        extract_result
299    }
300
301    /// Streams the archive directly into extraction without keeping a cached copy.
302    fn run_streaming_attempt(&self, format: CompressionFormat) -> Result<()> {
303        let _download_progress = ArchiveDownloadProgress::new(self.ctx.session().progress());
304        streaming_download_and_extract(
305            &self.archive().url,
306            format,
307            self.ctx.target_dir(),
308            self.ctx.static_files_dir(),
309            self.ctx.session(),
310        )
311    }
312
313    /// Extracts a cached archive file while updating shared extraction activity.
314    fn extract_cached_archive(&self, archive_path: &Path, format: CompressionFormat) -> Result<()> {
315        let mut extraction_progress = ArchiveExtractionProgress::new(self.ctx.session().progress());
316        let file = fs::open(archive_path)?;
317        let result = extract_archive_raw(
318            file,
319            format,
320            self.ctx.target_dir(),
321            self.ctx.static_files_dir(),
322            Some(&mut extraction_progress),
323        );
324        extraction_progress.finish();
325        result
326    }
327
328    /// Returns `true` if all declared plain outputs verify while updating shared verification
329    /// progress.
330    fn verify_outputs_with_progress(&self) -> Result<bool> {
331        let mut verification_progress =
332            ArchiveVerificationProgress::new(self.ctx.session().progress());
333        let verified = self
334            .output_verifier()
335            .verify_with_progress(&self.archive().output_files, Some(&mut verification_progress))?;
336        if verified {
337            verification_progress.complete(self.archive().output_size());
338        }
339        Ok(verified)
340    }
341}
342
343/// Chooses whether an archive attempt uses the cache or streams directly.
344#[derive(Debug, Clone, Copy, PartialEq, Eq)]
345enum ArchiveMode {
346    /// Download the archive to the cache, then extract it.
347    Cached,
348    /// Stream the archive directly into extraction.
349    Streaming,
350}
351
352impl ArchiveMode {
353    /// Picks the archive mode from the process context.
354    fn new(ctx: &ArchiveProcessContext) -> Result<Self> {
355        if ctx.cache_dir().is_some() {
356            ctx.session().require_request_limiter()?;
357            return Ok(Self::Cached)
358        }
359
360        Ok(Self::Streaming)
361    }
362
363    /// Returns `true` when fetch failures should retry the whole archive attempt.
364    const fn retries_fetch_errors(&self) -> bool {
365        matches!(self, Self::Cached)
366    }
367
368    /// Runs the selected archive mode for a single attempt.
369    fn execute(&self, processor: &ArchiveProcessor, format: CompressionFormat) -> Result<()> {
370        match self {
371            Self::Cached => processor.run_cached_attempt(format),
372            Self::Streaming => processor.run_streaming_attempt(format),
373        }
374    }
375}