aria2_core/engine/
download_progress.rs1use 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}