file_engine/operations/
move_op.rs1use 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 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 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}