Skip to main content

zng_ext_setup/task/
extract_tar.rs

1use std::{
2    collections::HashSet,
3    fmt, fs, io,
4    path::{Path, PathBuf},
5    sync::Arc,
6};
7
8use zng_task::Progress;
9use zng_txt::{ToTxt as _, Txt, formatx};
10use zng_unit::ByteUnits as _;
11use zng_var::{expr_var, var};
12
13use crate::task::{InstallTaskError, SetupTaskError};
14
15/// Setup task that extracts TAR container to a new or existing directory
16/// on install and removes these files on uninstall.
17pub enum ExtractTar {}
18impl ExtractTar {
19    fn prepare_install_blocking(args: super::PrepareInstallArgs<Self>) -> Result<PrepareInstallData, SetupTaskError> {
20        let (parent_dir, dir_name) = match (args.config.target_dir.parent(), args.config.target_dir.file_name()) {
21            (Some(p), Some(n)) if let Some(n) = n.to_str() => (p, n),
22            _ => {
23                return Err(SetupTaskError::io(
24                    args.config.target_dir,
25                    io::Error::new(io::ErrorKind::InvalidInput, "invalid target"),
26                ));
27            }
28        };
29        if let Err(e) = fs::create_dir_all(parent_dir) {
30            return Err(SetupTaskError::io(parent_dir.to_owned(), e));
31        }
32
33        // make a temp path beside the target_dir
34        let mut retries = 0;
35        let (tmp, tmp_state) = loop {
36            let tmp = parent_dir.join(format!("{dir_name}-temp{retries}"));
37            let state = tmp.join(TEMP_STATE);
38
39            if tmp.is_dir() {
40                if state.exists() {
41                    if let Ok(s) = fs::read_to_string(&state)
42                        && s != "prepared"
43                        && fs::remove_dir_all(&tmp).is_ok()
44                    {
45                        // existed, but was leftover from a failed cleanup
46                        break (tmp, state);
47                    }
48                } else if fs::remove_dir(&tmp).is_ok() {
49                    // existed, but empty
50                    break (tmp, state);
51                }
52            }
53
54            retries += 1;
55            if retries == 1000 {
56                return Err(SetupTaskError::io(
57                    args.config.target_dir,
58                    io::Error::new(io::ErrorKind::QuotaExceeded, "cannot create temp target"),
59                ));
60            }
61        };
62        if let Err(e) = fs::create_dir(&tmp) {
63            return Err(SetupTaskError::io(tmp, e));
64        }
65        if let Err(e) = fs::write(&tmp_state, "preparing") {
66            let _ = fs::remove_dir(&tmp);
67            return Err(SetupTaskError::io(tmp_state, e));
68        }
69
70        macro_rules! error {
71            ($e:expr) => {{
72                let _ = fs::write(&tmp_state, "error");
73                let _ = fs::remove_dir(&tmp);
74                return Err($e);
75            }};
76        }
77
78        let mut entries = HashSet::new();
79
80        // prepare progress reporting
81        let tar = zng_task::io::Measure::new(args.config.tar, args.config.tar_len.bytes(), 0.bytes());
82        let progress_name = var(Txt::default());
83        let progress = expr_var! {
84            let metrics = #{tar.metrics()};
85            let (n, total) = metrics.read_progress;
86            if n <= total {
87                Progress::from_n_of(n.0, total.0)
88            } else {
89                Progress::indeterminate()
90            }
91            .with_msg(formatx!("{}\n{}", #{progress_name.clone()}, metrics))
92            .with_meta_mut(|mut m| {
93                m.set(*zng_task::io::METRICS_ID, metrics.clone());
94            })
95        };
96        progress.set_bind(&args.progress).perm();
97
98        // extract
99        let mut tar = tar::Archive::new(tar);
100        let tar_entries = match tar.entries() {
101            Ok(e) => e,
102            Err(e) => error!(SetupTaskError::io(":tar/entries".into(), e)),
103        };
104        entries.insert(PathBuf::new()); // target_dir
105        for entry in tar_entries {
106            let mut entry = match entry {
107                Ok(e) => e,
108                Err(e) => error!(SetupTaskError::io(":tar/entry".into(), e)),
109            };
110            let path = match entry.path() {
111                Ok(p) => p.into_owned(),
112                Err(e) => error!(SetupTaskError::io(":tar/entry/path".into(), e)),
113            };
114
115            let display_name = path.as_os_str().to_string_lossy().replace('\\', "/").to_txt();
116
117            let entry_type = entry.header().entry_type();
118            if entry_type.is_dir() {
119                entries.insert(path);
120            } else if entry_type.is_file() {
121                for p in path.ancestors() {
122                    if !entries.contains(p) {
123                        entries.insert(p.to_owned());
124                    }
125                }
126                if !entries.insert(path) {
127                    error!(SetupTaskError::io(
128                        ":tar/entry/path".into(),
129                        io::Error::new(io::ErrorKind::InvalidData, "repeated file in tar")
130                    ))
131                }
132            } else {
133                if args.config.strict {
134                    error!(SetupTaskError::io(
135                        ":tar/entry/entry_type".into(),
136                        io::Error::new(
137                            io::ErrorKind::InvalidData,
138                            format!("found entry {entry_type:?}, only directory and files are allowed")
139                        )
140                    ))
141                } else {
142                    continue;
143                }
144            }
145
146            progress_name.set(display_name);
147
148            if let Err(e) = entry.unpack_in(&tmp) {
149                error!(SetupTaskError::io(":tar/entry/unpack_in".into(), e))
150            }
151
152            if args.cancel.get() {
153                break;
154            }
155        }
156        let _ = tar.into_inner().finish();
157
158        // if is updating find orphan entries
159        let mut remove = vec![];
160        if let Some(prev) = args.update
161            && !args.cancel.get()
162        {
163            args.progress.set(Progress::indeterminate());
164
165            if prev.target_dir != args.config.target_dir {
166                remove = prev
167                    .entries
168                    .into_iter()
169                    .filter_map(|p| {
170                        let p = prev.target_dir.join(p);
171                        if p.exists() { Some(p) } else { None }
172                    })
173                    .collect();
174            } else {
175                remove = prev
176                    .entries
177                    .into_iter()
178                    .filter_map(|p| {
179                        if entries.contains(&p) {
180                            None
181                        } else {
182                            let p = prev.target_dir.join(p);
183                            if p.exists() { Some(p) } else { None }
184                        }
185                    })
186                    .collect();
187            }
188        }
189
190        let _ = fs::write(&tmp_state, "prepared");
191
192        let mut entries: Vec<_> = entries.into_iter().collect();
193        entries.sort(); // root first, this allows renaming entire dirs when possible during `install`
194
195        Ok(PrepareInstallData {
196            temp_dir: tmp,
197            target_dir: args.config.target_dir,
198            add: entries,
199            remove,
200        })
201    }
202
203    fn install_blocking(args: super::InstallArgs<Self>) -> Result<InstallData, InstallTaskError<InstallData>> {
204        let mut errors = vec![];
205
206        let mut entries = args.data.add;
207        let mut moved_dirs = HashSet::<&Path>::new();
208        for add in &entries {
209            if add.ancestors().any(|p| moved_dirs.contains(p)) {
210                // already renamed parent dir
211                continue;
212            }
213            let from = args.data.temp_dir.join(add);
214            let to = args.data.target_dir.join(add);
215
216            if from.is_dir() {
217                // if `to` is existing dir
218                if to.is_dir() {
219                    // can still move entire dir if is empty
220                    let empty = match fs::read_dir(&to) {
221                        Ok(mut d) => d.next().is_none(),
222                        Err(_) => false,
223                    };
224                    if !empty {
225                        // otherwise needs to merge per entry
226                        continue;
227                    }
228                }
229
230                // if `to` is existing file, remove it
231                if let Err(e) = fs::remove_file(&to)
232                    && !matches!(e.kind(), io::ErrorKind::NotFound)
233                {
234                    errors.push((to, Arc::new(e)));
235                    continue;
236                }
237
238                // move dir
239                if let Err(e) = fs::rename(from, &to) {
240                    errors.push((to, Arc::new(e)));
241                    continue;
242                }
243
244                // moved entire dir
245                if add.as_os_str().is_empty() {
246                    // already moved root dir
247                    break;
248                }
249                moved_dirs.insert(add);
250            } else if from.is_file() {
251                // if is existing dir, remove it all
252                if let Err(e) = fs::remove_dir_all(&to)
253                    && !matches!(e.kind(), io::ErrorKind::NotADirectory)
254                {
255                    errors.push((to, Arc::new(e)));
256                    continue;
257                }
258
259                // move file
260                if let Err(e) = fs::rename(from, &to) {
261                    errors.push((to, Arc::new(e)));
262                }
263            } else {
264                let e = io::Error::new(io::ErrorKind::NotFound, "expected dir or file");
265                errors.push((from, Arc::new(e)));
266            }
267        }
268
269        entries.reverse(); // uninstall removes depth first to cleanup empty dirs as it goes
270        let data = InstallData {
271            target_dir: args.data.target_dir,
272            entries,
273        };
274
275        if errors.is_empty() {
276            Ok(data)
277        } else {
278            // tasks must cleanup prepared data in case of error
279            if let Err(e) = fs::remove_dir_all(&args.data.temp_dir)
280                && !matches!(e.kind(), io::ErrorKind::NotFound)
281            {
282                errors.push((args.data.temp_dir, Arc::new(e)));
283            }
284
285            Err(InstallTaskError {
286                error: SetupTaskError::Io(errors),
287                clean_data: Some(data),
288            })
289        }
290    }
291
292    fn validate_uninstall_blocking(args: super::ValidateUninstallArgs<Self>) -> Result<InstallData, SetupTaskError> {
293        let mut data = args.data;
294        if !data.target_dir.exists() {
295            data.entries.clear();
296            return Ok(data);
297        }
298
299        // retain only entries that exist
300        // expand entries to full path
301        // check if entries are really in the target_dir
302        let len = data.entries.len() as u64;
303        let mut i = 0u64;
304        let mut invalid_entry = false;
305        data.entries.retain_mut(|e| {
306            let path = data.target_dir.join(&e);
307            let retain = if path.starts_with(&data.target_dir) {
308                path.exists()
309            } else {
310                invalid_entry = true;
311                false
312            };
313            *e = path;
314            i += 1;
315            if i.is_multiple_of(50) {
316                // avoid too many progress notifications
317                args.progress.set(Progress::from_n_of(i, len));
318            }
319            retain
320        });
321        if invalid_entry {
322            #[derive(Debug)]
323            struct EntryPathNotInTargetDir;
324            impl fmt::Display for EntryPathNotInTargetDir {
325                fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
326                    write!(f, "entry path not in target dir")
327                }
328            }
329            impl std::error::Error for EntryPathNotInTargetDir {}
330            return Err(super::SetupTaskError::CorruptedTaskData(Arc::new(EntryPathNotInTargetDir)));
331        }
332
333        if i > 50 {
334            // go back to indeterminate if notified
335            args.progress.set(Progress::indeterminate());
336        }
337
338        // verify order, uninstall must consume depth first
339        data.entries.sort_by(|a, b| b.cmp(a));
340
341        Ok(data)
342    }
343
344    fn uninstall_blocking(args: super::UninstallArgs<Self>) -> Result<(), SetupTaskError> {
345        let mut errors = vec![];
346        for entry in args.data.entries {
347            let entry = args.data.target_dir.join(entry);
348            if entry.is_dir() {
349                if let Err(e) = fs::remove_dir(&entry)
350                    && !matches!(
351                        e.kind(),
352                        io::ErrorKind::NotFound | io::ErrorKind::NotADirectory | io::ErrorKind::DirectoryNotEmpty
353                    )
354                {
355                    errors.push((entry, Arc::new(e)));
356                }
357            } else if entry.is_file()
358                && let Err(e) = fs::remove_file(&entry)
359                && !matches!(e.kind(), io::ErrorKind::NotFound)
360            {
361                errors.push((entry, Arc::new(e)));
362            }
363        }
364        if errors.is_empty() {
365            Ok(())
366        } else {
367            Err(super::SetupTaskError::Io(errors))
368        }
369    }
370}
371impl super::SetupTask for ExtractTar {
372    type InstallConfig = ExtractTarConfig;
373
374    type PrepareInstall = PrepareInstallData;
375
376    type Install = InstallData;
377
378    fn task_type_id() -> super::TaskTypeId {
379        "zng-setup/ExtractTar".into()
380    }
381
382    async fn prepare_install(args: super::PrepareInstallArgs<Self>) -> Result<Self::PrepareInstall, SetupTaskError> {
383        zng_task::wait(move || Self::prepare_install_blocking(args)).await
384    }
385
386    async fn install(args: super::InstallArgs<Self>) -> Result<Self::Install, InstallTaskError<Self::Install>> {
387        zng_task::wait(move || Self::install_blocking(args)).await
388    }
389
390    async fn cancel_install(args: super::CancelInstallArgs<Self>) -> Result<(), SetupTaskError> {
391        if let Err(e) = zng_task::fs::remove_dir_all(&args.data.temp_dir).await
392            && !matches!(e.kind(), io::ErrorKind::NotFound)
393        {
394            return Err(SetupTaskError::io(args.data.temp_dir, e));
395        }
396        Ok(())
397    }
398
399    async fn validate_uninstall(args: super::ValidateUninstallArgs<Self>) -> Result<Self::Install, SetupTaskError> {
400        zng_task::wait(move || Self::validate_uninstall_blocking(args)).await
401    }
402
403    async fn uninstall(args: super::UninstallArgs<Self>) -> Result<(), SetupTaskError> {
404        zng_task::wait(move || Self::uninstall_blocking(args)).await
405    }
406}
407
408const TEMP_STATE: &str = ".zng-setup-ExtractTar";
409
410/// Config for [`ExtractTar`]
411pub struct ExtractTarConfig {
412    tar_len: u64,
413    tar: Box<dyn io::Read + Send>,
414    target_dir: PathBuf,
415    strict: bool,
416}
417
418impl ExtractTarConfig {
419    /// New config.
420    ///
421    /// * `tar_len` estimated length of `tar`, used for progress reporting only. Pass `0` for indeterminate.
422    /// * `tar` must read only a TAR stream from start to end. Must contain only directory and file entries.
423    /// * `target_dir` Directory that will be created or merged with the `tar` root directory.
424    ///
425    /// If all entries share a common first path component (for example, `my-app/bin/app.exe` and `my-app/README.md`),
426    /// that directory is created inside `target_dir`. To extract files directly into `target_dir`, the entries must
427    /// not have a common leading directory component.
428    pub fn new(tar_len: u64, tar: Box<dyn io::Read + Send>, target_dir: PathBuf) -> Self {
429        Self {
430            tar_len,
431            tar,
432            target_dir,
433            strict: false,
434        }
435    }
436
437    /// New config from a [`SfxClient`] read stream.
438    ///
439    /// [`SfxClient`]: crate::SfxClient
440    pub fn from_sfx(tar: crate::SfxReadBlocking, target_dir: PathBuf) -> Self {
441        Self::new(tar.exact_len().unwrap_or(0), Box::new(tar), target_dir)
442    }
443
444    /// Enable strict errors.
445    ///
446    /// When enabled:
447    ///
448    /// * Error on TAR entry that is not directory nor file.
449    pub fn strict(mut self) -> Self {
450        self.strict = true;
451        self
452    }
453}
454
455#[doc(hidden)]
456#[derive(Debug, PartialEq, Clone, serde::Serialize, serde::Deserialize)]
457pub struct PrepareInstallData {
458    temp_dir: PathBuf,
459    target_dir: PathBuf,
460    /// relative to the temp_dir, must move to temp_dir
461    add: Vec<PathBuf>,
462    /// absolute paths from previous install
463    remove: Vec<PathBuf>,
464}
465
466#[doc(hidden)]
467#[derive(Debug, PartialEq, Clone, serde::Serialize, serde::Deserialize)]
468pub struct InstallData {
469    target_dir: PathBuf,
470    // entries relative to the `target_dir` that where created by the task.
471    entries: Vec<PathBuf>,
472}