reth_cli_commands/download/
archive.rs1use 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
24pub(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
58struct ModularDownloadJob {
60 ctx: ArchiveProcessContext,
62 archive_concurrency: usize,
64}
65
66impl ModularDownloadJob {
67 const fn new(ctx: ArchiveProcessContext, archive_concurrency: usize) -> Self {
69 Self { ctx, archive_concurrency }
70 }
71
72 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107enum ArchiveAttemptState {
108 RunAttempt,
110 VerifyOutputs,
112 RetryAttempt,
114 Complete,
116 Fail,
118}
119
120struct ArchiveProcessor {
122 archive: PlannedArchive,
124 ctx: ArchiveProcessContext,
126}
127
128impl ArchiveProcessor {
129 fn new(archive: PlannedArchive, ctx: ArchiveProcessContext) -> Self {
131 Self { archive, ctx }
132 }
133
134 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 fn archive(&self) -> &SnapshotArchive {
219 &self.archive.archive
220 }
221
222 fn output_verifier(&self) -> OutputVerifier<'_> {
224 OutputVerifier::new(self.ctx.target_dir(), self.ctx.static_files_dir())
225 }
226
227 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 fn cleanup_outputs(&self) {
240 self.output_verifier().cleanup(&self.archive().output_files);
241 }
242
243 fn verify_outputs(&self) -> Result<bool> {
246 self.output_verifier().verify(&self.archive().output_files)
247 }
248
249 fn mark_complete(&self) {
251 self.ctx.session().record_reused_archive(self.archive().size, self.archive().output_size());
252 }
253
254 fn run_attempt(&self, mode: ArchiveMode, format: CompressionFormat) -> Result<()> {
256 mode.execute(self, format)
257 }
258
259 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 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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
345enum ArchiveMode {
346 Cached,
348 Streaming,
350}
351
352impl ArchiveMode {
353 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 const fn retries_fetch_errors(&self) -> bool {
365 matches!(self, Self::Cached)
366 }
367
368 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}