Skip to main content

file_engine/operations/
move_op.rs

1use std::path::PathBuf;
2
3use tokio_stream::wrappers::UnboundedReceiverStream;
4use tokio_util::sync::CancellationToken;
5
6use crate::error::{from_io, FileEngineError, Result};
7use crate::handle::Handle;
8#[cfg(feature = "permissions")]
9use crate::operations::copy::preserve_permissions_recursive;
10use crate::operations::copy::{copy_dir, copy_file, copy_symlink};
11
12pub struct MoveBuilder {
13    src: PathBuf,
14    dst: PathBuf,
15    overwrite: bool,
16    buffer_size: usize,
17    follow_symlinks: bool,
18    #[cfg(feature = "permissions")]
19    pub(crate) preserve_permissions: bool,
20    cancel_token: Option<CancellationToken>,
21}
22
23impl MoveBuilder {
24    pub(crate) fn new(
25        src: PathBuf,
26        dst: PathBuf,
27        buffer_size: usize,
28        follow_symlinks: bool,
29    ) -> Self {
30        Self {
31            src,
32            dst,
33            overwrite: false,
34            buffer_size,
35            follow_symlinks,
36            #[cfg(feature = "permissions")]
37            preserve_permissions: false,
38            cancel_token: None,
39        }
40    }
41
42    pub fn overwrite(mut self, enabled: bool) -> Self {
43        self.overwrite = enabled;
44        self
45    }
46
47    pub fn cancellation_token(mut self, token: CancellationToken) -> Self {
48        self.cancel_token = Some(token);
49        self
50    }
51
52    pub fn start(self) -> Result<Handle<()>> {
53        let cancel_token = self.cancel_token.unwrap_or_default();
54        let (progress_tx, progress_rx) = tokio::sync::mpsc::unbounded_channel();
55
56        let task_cancel_token = cancel_token.clone();
57        let buffer_size = self.buffer_size;
58        let follow_symlinks = self.follow_symlinks;
59        let overwrite = self.overwrite;
60        let src = self.src;
61        let dst = self.dst;
62        #[cfg(feature = "permissions")]
63        let preserve_permissions = self.preserve_permissions;
64
65        let join = tokio::spawn(async move {
66            if !overwrite && tokio::fs::try_exists(&dst).await.unwrap_or(false) {
67                return Err(FileEngineError::DestinationExists(dst));
68            }
69
70            // Fast path: same-filesystem rename is atomic and correct for
71            // both files and directories. Only fall back to copy+remove on
72            // failure (e.g. a cross-device move, `EXDEV`).
73            if tokio::fs::rename(&src, &dst).await.is_ok() {
74                return Ok(());
75            }
76
77            let src_metadata = tokio::fs::symlink_metadata(&src)
78                .await
79                .map_err(|e| from_io(src.clone(), e))?;
80
81            // Permissions must be applied (when requested) before `src` is
82            // removed below — `copy_file`/`copy_dir` create `dst` with
83            // default permissions, not `src`'s.
84            if src_metadata.is_dir() {
85                copy_dir(
86                    &src,
87                    &dst,
88                    buffer_size,
89                    follow_symlinks,
90                    &progress_tx,
91                    &task_cancel_token,
92                )
93                .await?;
94                #[cfg(feature = "permissions")]
95                if preserve_permissions {
96                    preserve_permissions_recursive(&src, &dst).await?;
97                }
98                tokio::fs::remove_dir_all(&src)
99                    .await
100                    .map_err(|e| from_io(src.clone(), e))?;
101            } else if src_metadata.is_symlink() && !follow_symlinks {
102                copy_symlink(&src, &dst).await?;
103                tokio::fs::remove_file(&src)
104                    .await
105                    .map_err(|e| from_io(src.clone(), e))?;
106            } else {
107                copy_file(
108                    &src,
109                    &dst,
110                    buffer_size,
111                    0,
112                    1,
113                    &progress_tx,
114                    &task_cancel_token,
115                )
116                .await?;
117                #[cfg(feature = "permissions")]
118                if preserve_permissions {
119                    preserve_permissions_recursive(&src, &dst).await?;
120                }
121                tokio::fs::remove_file(&src)
122                    .await
123                    .map_err(|e| from_io(src.clone(), e))?;
124            }
125
126            Ok(())
127        });
128
129        Ok(Handle {
130            join,
131            progress_rx: UnboundedReceiverStream::new(progress_rx),
132            cancel_token,
133        })
134    }
135}