datui_lib/loading/
unfinished.rs1use std::path::{Path, PathBuf};
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{Arc, Condvar, Mutex, MutexGuard};
13use std::time::{Duration, Instant};
14
15const SWEEP_GRACE: Duration = Duration::from_secs(1);
17
18const RETRY_EVERY: Duration = Duration::from_millis(50);
20
21#[derive(Clone, Debug, Default)]
23pub(crate) struct Unfinished(Arc<(Mutex<Files>, Condvar)>);
24
25#[derive(Debug, Default)]
26struct Files {
27 claimed: Vec<PathBuf>,
28 creating: usize,
31 left_over: Vec<PathBuf>,
34 swept: bool,
35 #[cfg(test)]
37 sweep_waiting: bool,
38}
39
40impl Files {
41 fn busy(&self) -> bool {
42 !self.claimed.is_empty() || self.creating > 0
43 }
44}
45
46#[derive(Clone, Debug, Default)]
48pub(crate) struct Writer {
49 stop: Arc<AtomicBool>,
50 files: Unfinished,
51}
52
53#[derive(Debug)]
56#[must_use]
57pub(crate) struct Claim {
58 files: Unfinished,
59 path: PathBuf,
60}
61
62impl Unfinished {
63 pub(crate) fn writer(&self, stop: Arc<AtomicBool>) -> Writer {
65 Writer {
66 stop,
67 files: self.clone(),
68 }
69 }
70
71 #[cfg(test)]
73 pub(crate) fn writing(&self) -> bool {
74 self.lock().busy()
75 }
76
77 pub(crate) fn sweep(&self, deadline: Instant) {
81 let (_, released) = &*self.0;
82 let mut files = self.lock();
83 while files.busy() {
84 let left = deadline.saturating_duration_since(Instant::now());
85 if left.is_zero() {
86 break;
87 }
88 #[cfg(test)]
91 {
92 files.sweep_waiting = true;
93 released.notify_all();
94 }
95 files = released
96 .wait_timeout(files, left)
97 .unwrap_or_else(|e| e.into_inner())
98 .0;
99 }
100 files.swept = true;
101 let mut left = std::mem::take(&mut files.claimed);
102 left.append(&mut files.left_over);
103 drop(files);
104 loop {
105 left.retain(|path| !removed(path));
106 if left.is_empty() || Instant::now() + RETRY_EVERY > deadline {
107 break;
108 }
109 std::thread::sleep(RETRY_EVERY);
110 }
111 for path in left {
112 log::warn!("Could not remove the temporary file {}", path.display());
113 }
114 }
115
116 #[cfg(test)]
119 fn a_sweep_waits_within(&self, limit: Duration) -> bool {
120 let (_, released) = &*self.0;
121 released
122 .wait_timeout_while(self.lock(), limit, |files| !files.sweep_waiting)
123 .unwrap_or_else(|e| e.into_inner())
124 .0
125 .sweep_waiting
126 }
127
128 fn lock(&self) -> MutexGuard<'_, Files> {
129 self.0.0.lock().unwrap_or_else(|e| e.into_inner())
130 }
131
132 fn released(&self) {
133 self.0.1.notify_all();
134 }
135}
136
137impl Writer {
138 pub(crate) fn stopped(&self) -> bool {
140 self.stop.load(Ordering::Relaxed)
141 }
142
143 pub(crate) fn create<F: AsRef<Path>, E>(
147 &self,
148 make: impl FnOnce() -> Result<F, E>,
149 ) -> Result<Option<(F, Claim)>, E> {
150 {
151 let mut files = self.files.lock();
152 if files.swept || self.stopped() {
153 return Ok(None);
154 }
155 files.creating += 1;
156 }
157 let made = make();
158 let mut files = self.files.lock();
159 let claimed = match made {
160 Ok(file) if !files.swept && !self.stopped() => {
161 let path = file.as_ref().to_path_buf();
162 files.claimed.push(path.clone());
163 Ok(Some((
164 file,
165 Claim {
166 files: self.files.clone(),
167 path,
168 },
169 )))
170 }
171 Ok(file) => {
172 drop(files);
174 drop(file);
175 files = self.files.lock();
176 Ok(None)
177 }
178 Err(e) => Err(e),
179 };
180 files.creating -= 1;
181 drop(files);
182 self.files.released();
183 claimed
184 }
185}
186
187#[must_use]
190pub struct ExitSweep(pub(crate) Unfinished);
191
192impl Drop for ExitSweep {
193 fn drop(&mut self) {
194 self.0.sweep(Instant::now() + SWEEP_GRACE);
195 }
196}
197
198impl Drop for Claim {
199 fn drop(&mut self) {
200 let mut files = self.files.lock();
201 if let Some(at) = files.claimed.iter().position(|p| *p == self.path) {
202 files.claimed.swap_remove(at);
203 }
204 let earlier = std::mem::take(&mut files.left_over);
206 drop(files);
207 let mut left: Vec<PathBuf> = earlier.into_iter().filter(|p| !removed(p)).collect();
208 if std::fs::symlink_metadata(&self.path).is_ok() {
210 left.push(self.path.clone());
211 }
212 self.files.lock().left_over.append(&mut left);
213 self.files.released();
214 }
215}
216
217fn removed(path: &Path) -> bool {
219 match std::fs::remove_file(path) {
220 Ok(()) => true,
221 Err(e) => e.kind() == std::io::ErrorKind::NotFound,
222 }
223}
224
225#[cfg(test)]
226mod tests;