Skip to main content

aria2_core/engine/
download_progress.rs

1use std::sync::Arc;
2use std::time::Instant;
3use tokio::sync::mpsc;
4
5use crate::engine::command::ProgressUpdate;
6use crate::request::request_group::RequestGroup;
7use crate::util::perf_monitor::{AtomicMetrics, PerformanceMonitor};
8
9pub struct ProgressUpdater {
10    progress_sender: Option<mpsc::UnboundedSender<ProgressUpdate>>,
11    group: Arc<tokio::sync::RwLock<RequestGroup>>,
12    atomic_metrics: Arc<AtomicMetrics>,
13    perf_monitor: Option<Arc<PerformanceMonitor>>,
14    last_speed_update: Instant,
15    last_completed: u64,
16    last_progress_update: u64,
17}
18
19impl Clone for ProgressUpdater {
20    fn clone(&self) -> Self {
21        Self {
22            progress_sender: self.progress_sender.clone(),
23            group: Arc::clone(&self.group),
24            atomic_metrics: Arc::clone(&self.atomic_metrics),
25            perf_monitor: self.perf_monitor.clone(),
26            last_speed_update: self.last_speed_update,
27            last_completed: self.last_completed,
28            last_progress_update: self.last_progress_update,
29        }
30    }
31}
32
33impl ProgressUpdater {
34    pub fn new(
35        progress_sender: Option<mpsc::UnboundedSender<ProgressUpdate>>,
36        group: Arc<tokio::sync::RwLock<RequestGroup>>,
37        atomic_metrics: Arc<AtomicMetrics>,
38        perf_monitor: Option<Arc<PerformanceMonitor>>,
39    ) -> Self {
40        Self {
41            progress_sender,
42            group,
43            atomic_metrics,
44            perf_monitor,
45            last_speed_update: Instant::now(),
46            last_completed: 0,
47            last_progress_update: 0,
48        }
49    }
50
51    pub fn reset(&mut self, completed_bytes: u64) {
52        self.last_speed_update = Instant::now();
53        self.last_completed = completed_bytes;
54        self.last_progress_update = completed_bytes;
55    }
56
57    pub async fn update_progress(
58        &mut self,
59        completed_bytes: u64,
60        progress_update_threshold: u64,
61        speed_update_interval_ms: u64,
62    ) {
63        if completed_bytes - self.last_progress_update < progress_update_threshold {
64            return;
65        }
66
67        let elapsed = self.last_speed_update.elapsed();
68        let speed = if elapsed.as_millis() >= speed_update_interval_ms as u128 {
69            let delta = completed_bytes - self.last_completed;
70            let s = (delta as f64 / elapsed.as_secs_f64()) as u64;
71            self.last_speed_update = Instant::now();
72            self.last_completed = completed_bytes;
73            s
74        } else {
75            0
76        };
77
78        if let Some(ref sender) = self.progress_sender {
79            let _ = sender.send(ProgressUpdate {
80                completed_bytes,
81                download_speed: speed,
82                upload_speed: 0,
83            });
84        } else {
85            let g = self.group.write().await;
86            g.update_progress(completed_bytes).await;
87            g.set_completed_length(completed_bytes);
88            if speed > 0 {
89                g.update_speed(speed, 0).await;
90                g.set_download_speed_cached(speed);
91            }
92        }
93
94        if speed > 0 {
95            self.atomic_metrics.record_throughput(speed);
96            if let Some(ref monitor) = self.perf_monitor {
97                let metrics = crate::util::perf_monitor::Metrics::new(
98                    speed,
99                    elapsed.as_millis() as u64,
100                    0,
101                    0,
102                )
103                .with_label("download_speed");
104                monitor.record_metric("download_speed", metrics);
105            }
106        }
107
108        self.last_progress_update = completed_bytes;
109    }
110
111    pub fn last_progress_update(&self) -> u64 {
112        self.last_progress_update
113    }
114}