microsandbox_image/progress.rs
1//! Pull progress reporting.
2
3use std::sync::Arc;
4
5use tokio::sync::mpsc;
6
7//--------------------------------------------------------------------------------------------------
8// Constants
9//--------------------------------------------------------------------------------------------------
10
11/// Default channel capacity.
12const DEFAULT_PROGRESS_CHANNEL_CAPACITY: usize = 1024;
13
14//--------------------------------------------------------------------------------------------------
15// Types
16//--------------------------------------------------------------------------------------------------
17
18/// Progress events emitted during image pull and EROFS materialization.
19#[derive(Debug, Clone, serde::Serialize)]
20#[serde(tag = "kind", rename_all = "snake_case")]
21pub enum PullProgress {
22 /// Resolving the image reference.
23 Resolving {
24 /// The image reference being resolved.
25 reference: Arc<str>,
26 },
27
28 /// Manifest parsed. Layer count and total sizes now known.
29 Resolved {
30 /// The image reference.
31 reference: Arc<str>,
32 /// Resolved manifest digest.
33 manifest_digest: Arc<str>,
34 /// Number of layers.
35 layer_count: usize,
36 /// Sum of compressed layer sizes. `None` if manifest omits sizes.
37 total_download_bytes: Option<u64>,
38 },
39
40 /// Byte-level download progress for a single layer.
41 LayerDownloadProgress {
42 /// Layer index (0-based).
43 layer_index: usize,
44 /// Layer digest.
45 digest: Arc<str>,
46 /// Bytes downloaded so far.
47 downloaded_bytes: u64,
48 /// Total bytes (if known).
49 total_bytes: Option<u64>,
50 },
51
52 /// A single layer download completed and verified.
53 LayerDownloadComplete {
54 /// Layer index.
55 layer_index: usize,
56 /// Layer digest.
57 digest: Arc<str>,
58 /// Total downloaded bytes.
59 downloaded_bytes: u64,
60 },
61
62 /// Layer download completed and the blob is being verified.
63 LayerDownloadVerifying {
64 /// Layer index.
65 layer_index: usize,
66 /// Layer digest.
67 digest: Arc<str>,
68 },
69
70 /// Layer EROFS materialization started.
71 LayerMaterializeStarted {
72 /// Layer index.
73 layer_index: usize,
74 /// Layer diff ID.
75 diff_id: Arc<str>,
76 },
77
78 /// Byte-level materialization progress for a single layer.
79 LayerMaterializeProgress {
80 /// Layer index (0-based).
81 layer_index: usize,
82 /// Bytes read so far.
83 bytes_read: u64,
84 /// Total bytes.
85 total_bytes: u64,
86 },
87
88 /// Layer tar ingest is complete and the EROFS image is being written.
89 LayerMaterializeWriting {
90 /// Layer index.
91 layer_index: usize,
92 },
93
94 /// Layer EROFS materialization completed.
95 LayerMaterializeComplete {
96 /// Layer index.
97 layer_index: usize,
98 /// Layer diff ID.
99 diff_id: Arc<str>,
100 },
101
102 /// Merging per-layer trees into the unified rootfs view.
103 StitchMergingTrees {
104 /// Number of layers being merged.
105 layer_count: usize,
106 },
107
108 /// Writing the fsmeta EROFS image (metadata-only merged view).
109 StitchWritingFsmeta,
110
111 /// Writing the VMDK descriptor that stitches fsmeta + layer EROFSes.
112 StitchWritingVmdk,
113
114 /// Stitching phase finished — fsmeta + VMDK are on disk.
115 StitchComplete,
116
117 /// Entire image pull completed.
118 Complete {
119 /// The image reference.
120 reference: Arc<str>,
121 /// Number of layers.
122 layer_count: usize,
123 },
124}
125
126/// Receiver for progress events.
127pub struct PullProgressHandle {
128 rx: mpsc::Receiver<PullProgress>,
129}
130
131/// Emits progress events. Uses `try_send` — never blocks downloads.
132#[derive(Clone, Debug)]
133pub struct PullProgressSender {
134 tx: mpsc::Sender<PullProgress>,
135}
136
137//--------------------------------------------------------------------------------------------------
138// Methods
139//--------------------------------------------------------------------------------------------------
140
141impl PullProgressHandle {
142 /// Receive the next event. Returns `None` when the pull completes.
143 pub async fn recv(&mut self) -> Option<PullProgress> {
144 self.rx.recv().await
145 }
146
147 /// Convert into the underlying receiver for use with `tokio::select!`.
148 pub fn into_receiver(self) -> mpsc::Receiver<PullProgress> {
149 self.rx
150 }
151}
152
153impl PullProgressSender {
154 /// Emit a progress event. Silently discards if receiver is full or dropped.
155 pub fn send(&self, event: PullProgress) {
156 let _ = self.tx.try_send(event);
157 }
158}
159
160//--------------------------------------------------------------------------------------------------
161// Functions
162//--------------------------------------------------------------------------------------------------
163
164/// Create a progress channel pair.
165pub fn progress_channel() -> (PullProgressHandle, PullProgressSender) {
166 let (tx, rx) = mpsc::channel(DEFAULT_PROGRESS_CHANNEL_CAPACITY);
167 (PullProgressHandle { rx }, PullProgressSender { tx })
168}