deps_engine/progress.rs
1//! Progress-reporting port for registry fetch tasks.
2//!
3//! [`ProgressSender`] lets a fetch task report `fetched`/`total` progress without knowing
4//! whether, or how, anything renders it — a driving adapter drains the paired
5//! `mpsc::Receiver<ProgressUpdate>` returned by [`channel`] however it sees fit (`deps-lsp`'s
6//! `RegistryProgress` drives the LSP work-done-progress protocol from it; `deps-cli` may render
7//! a terminal progress bar the same way, or drop the receiver and pass `None` as the sender).
8
9use tokio::sync::mpsc;
10
11/// Channel capacity for progress updates.
12/// Small buffer is sufficient since updates are coalesced by the consumer.
13const PROGRESS_CHANNEL_CAPACITY: usize = 8;
14
15/// Non-blocking sender for progress updates from fetch tasks.
16///
17/// Cheap to clone and safe to use from multiple concurrent futures.
18/// Dropped messages are acceptable — progress is best-effort UI feedback.
19#[derive(Clone)]
20pub struct ProgressSender {
21 tx: mpsc::Sender<ProgressUpdate>,
22 total: usize,
23}
24
25/// A single fetch-progress observation: `fetched` out of `total` packages processed so far.
26#[derive(Debug, Clone, Copy)]
27pub struct ProgressUpdate {
28 /// Number of packages fetched so far.
29 pub fetched: usize,
30 /// Total number of packages being fetched in this run.
31 pub total: usize,
32}
33
34impl ProgressSender {
35 /// Send a progress update without blocking.
36 ///
37 /// Uses `try_send` — if the channel is full, the update is silently dropped.
38 /// This is intentional: progress is best-effort UI feedback, and dropping
39 /// updates is always preferable to blocking fetch tasks.
40 ///
41 /// # Examples
42 ///
43 /// ```
44 /// let (sender, mut receiver) = deps_engine::progress::channel(10);
45 /// sender.send(3);
46 /// let update = receiver.try_recv().unwrap();
47 /// assert_eq!((update.fetched, update.total), (3, 10));
48 /// ```
49 pub fn send(&self, fetched: usize) {
50 let _ = self.tx.try_send(ProgressUpdate {
51 fetched,
52 total: self.total,
53 });
54 }
55}
56
57/// Creates a progress port: a [`ProgressSender`] fetch tasks report through, paired with the
58/// `Receiver` a driving adapter drains to render progress however it sees fit.
59///
60/// `total` is the total unit count (e.g. dependency count) every [`ProgressUpdate`] sent
61/// through the returned sender will carry.
62///
63/// # Examples
64///
65/// ```
66/// let (sender, mut receiver) = deps_engine::progress::channel(5);
67/// sender.send(1);
68/// sender.send(2);
69/// assert_eq!(receiver.try_recv().unwrap().fetched, 1);
70/// assert_eq!(receiver.try_recv().unwrap().fetched, 2);
71/// ```
72#[must_use]
73pub fn channel(total: usize) -> (ProgressSender, mpsc::Receiver<ProgressUpdate>) {
74 let (tx, rx) = mpsc::channel(PROGRESS_CHANNEL_CAPACITY);
75 (ProgressSender { tx, total }, rx)
76}
77
78#[cfg(test)]
79mod tests {
80 use super::*;
81
82 #[test]
83 fn test_percentage_calculation() {
84 let calculate = |fetched: usize, total: usize| -> u32 {
85 if total == 0 {
86 return 0;
87 }
88 ((fetched as f64 / total as f64) * 100.0) as u32
89 };
90
91 assert_eq!(calculate(0, 10), 0);
92 assert_eq!(calculate(5, 10), 50);
93 assert_eq!(calculate(10, 10), 100);
94 assert_eq!(calculate(7, 10), 70);
95 assert_eq!(calculate(0, 0), 0);
96 }
97
98 #[tokio::test]
99 async fn test_progress_sender_try_send_on_closed_channel() {
100 let (sender, rx) = channel(10);
101
102 drop(rx);
103
104 sender.send(5);
105 }
106
107 #[tokio::test]
108 async fn test_progress_sender_try_send_on_full_channel() {
109 let (tx, _rx) = mpsc::channel(1);
110 let sender = ProgressSender { tx, total: 10 };
111
112 sender.send(1);
113 sender.send(2);
114 sender.send(3);
115 }
116}