1use std::path::{Path, PathBuf};
23use std::sync::atomic::{AtomicBool, Ordering};
24use std::sync::{Arc, Condvar, Mutex, MutexGuard};
25use std::time::{Duration, Instant};
26
27const SWEEP_GRACE: Duration = Duration::from_secs(1);
29
30const RETRY_EVERY: Duration = Duration::from_millis(50);
32
33#[derive(Clone, Debug, Default)]
35pub(crate) struct Unfinished(Arc<(Mutex<Files>, Condvar)>);
36
37#[derive(Debug, Default)]
38struct Files {
39 claimed: Vec<PathBuf>,
40 creating: usize,
43 left_over: Vec<PathBuf>,
46 swept: bool,
47 #[cfg(test)]
49 sweep_waiting: bool,
50}
51
52impl Files {
53 fn busy(&self) -> bool {
54 !self.claimed.is_empty() || self.creating > 0
55 }
56}
57
58#[derive(Clone, Debug, Default)]
60pub(crate) struct Writer {
61 stop: Arc<AtomicBool>,
62 files: Unfinished,
63}
64
65#[derive(Debug)]
68#[must_use]
69pub(crate) struct Claim {
70 files: Unfinished,
71 path: PathBuf,
72}
73
74impl Unfinished {
75 pub(crate) fn writer(&self, stop: Arc<AtomicBool>) -> Writer {
77 Writer {
78 stop,
79 files: self.clone(),
80 }
81 }
82
83 #[cfg(test)]
85 pub(crate) fn writing(&self) -> bool {
86 self.lock().busy()
87 }
88
89 pub(crate) fn sweep(&self, deadline: Instant) {
94 let (_, released) = &*self.0;
95 let mut files = self.lock();
96 while files.busy() {
97 let left = deadline.saturating_duration_since(Instant::now());
98 if left.is_zero() {
99 break;
100 }
101 #[cfg(test)]
104 {
105 files.sweep_waiting = true;
106 released.notify_all();
107 }
108 files = released
109 .wait_timeout(files, left)
110 .unwrap_or_else(|e| e.into_inner())
111 .0;
112 }
113 files.swept = true;
114 let mut left = std::mem::take(&mut files.claimed);
115 left.append(&mut files.left_over);
116 drop(files);
117 loop {
118 left.retain(|path| !removed(path));
119 if left.is_empty() || Instant::now() + RETRY_EVERY > deadline {
120 break;
121 }
122 std::thread::sleep(RETRY_EVERY);
123 }
124 for path in left {
125 log::warn!("Could not remove the temporary file {}", path.display());
126 }
127 }
128
129 #[cfg(test)]
132 fn a_sweep_waits_within(&self, limit: Duration) -> bool {
133 let (_, released) = &*self.0;
134 released
135 .wait_timeout_while(self.lock(), limit, |files| !files.sweep_waiting)
136 .unwrap_or_else(|e| e.into_inner())
137 .0
138 .sweep_waiting
139 }
140
141 fn lock(&self) -> MutexGuard<'_, Files> {
142 self.0.0.lock().unwrap_or_else(|e| e.into_inner())
143 }
144
145 fn released(&self) {
146 self.0.1.notify_all();
147 }
148}
149
150impl Writer {
151 pub(crate) fn stopped(&self) -> bool {
153 self.stop.load(Ordering::Relaxed)
154 }
155
156 pub(crate) fn create<F: AsRef<Path>, E>(
164 &self,
165 make: impl FnOnce() -> Result<F, E>,
166 ) -> Result<Option<(F, Claim)>, E> {
167 {
168 let mut files = self.files.lock();
169 if files.swept || self.stopped() {
170 return Ok(None);
171 }
172 files.creating += 1;
173 }
174 let made = make();
175 let mut files = self.files.lock();
176 let claimed = match made {
177 Ok(file) if !files.swept && !self.stopped() => {
178 let path = file.as_ref().to_path_buf();
179 files.claimed.push(path.clone());
180 Ok(Some((
181 file,
182 Claim {
183 files: self.files.clone(),
184 path,
185 },
186 )))
187 }
188 Ok(file) => {
189 drop(files);
191 drop(file);
192 files = self.files.lock();
193 Ok(None)
194 }
195 Err(e) => Err(e),
196 };
197 files.creating -= 1;
198 drop(files);
199 self.files.released();
200 claimed
201 }
202}
203
204#[must_use]
207pub struct ExitSweep(pub(crate) Unfinished);
208
209impl Drop for ExitSweep {
210 fn drop(&mut self) {
211 self.0.sweep(Instant::now() + SWEEP_GRACE);
212 }
213}
214
215impl Drop for Claim {
216 fn drop(&mut self) {
217 let mut files = self.files.lock();
218 if let Some(at) = files.claimed.iter().position(|p| *p == self.path) {
219 files.claimed.swap_remove(at);
220 }
221 let earlier = std::mem::take(&mut files.left_over);
223 drop(files);
224 let mut left: Vec<PathBuf> = earlier.into_iter().filter(|p| !removed(p)).collect();
225 if std::fs::symlink_metadata(&self.path).is_ok() {
227 left.push(self.path.clone());
228 }
229 self.files.lock().left_over.append(&mut left);
230 self.files.released();
231 }
232}
233
234fn removed(path: &Path) -> bool {
236 match std::fs::remove_file(path) {
237 Ok(()) => true,
238 Err(e) => e.kind() == std::io::ErrorKind::NotFound,
239 }
240}
241
242#[cfg(test)]
243mod tests {
244 use super::*;
245
246 fn file_in(dir: &Path) -> std::io::Result<tempfile::NamedTempFile> {
247 tempfile::NamedTempFile::new_in(dir)
248 }
249
250 fn files_in(dir: &Path) -> usize {
251 std::fs::read_dir(dir).unwrap().count()
252 }
253
254 #[test]
257 fn a_sweep_waits_for_a_writer_that_stops() {
258 let dir = tempfile::tempdir().unwrap();
259 let unfinished = Unfinished::default();
260 let stop = Arc::new(AtomicBool::new(false));
261 let writer = unfinished.writer(stop.clone());
262 let (claimed, wait) = std::sync::mpsc::channel();
263 let worker = {
264 let dir = dir.path().to_path_buf();
265 std::thread::spawn(move || {
266 let (file, claim) = writer
267 .create(|| file_in(&dir))
268 .unwrap()
269 .expect("not stopped yet");
270 claimed.send(()).unwrap();
271 while !writer.stopped() {
272 std::thread::sleep(Duration::from_millis(5));
273 }
274 drop(file);
276 drop(claim);
277 })
278 };
279 wait.recv().unwrap();
280 stop.store(true, Ordering::Relaxed);
281 let began = Instant::now();
282 unfinished.sweep(began + Duration::from_secs(30));
283 assert!(
284 began.elapsed() < Duration::from_secs(10),
285 "not the whole grace"
286 );
287 assert_eq!(files_in(dir.path()), 0);
288 assert!(!unfinished.writing());
289 worker.join().unwrap();
290 }
291
292 #[test]
295 fn a_sweep_removes_what_a_writer_still_holds() {
296 let dir = tempfile::tempdir().unwrap();
297 let unfinished = Unfinished::default();
298 let writer = unfinished.writer(Arc::default());
299 let (file, claim) = writer
300 .create(|| file_in(dir.path()))
301 .unwrap()
302 .expect("not stopped yet");
303 let path = file.path().to_path_buf();
304 unfinished.sweep(Instant::now() + Duration::from_millis(20));
305 assert!(!path.exists());
306 assert!(!unfinished.writing());
307 drop(file);
308 drop(claim);
309 }
310
311 #[test]
315 fn a_sweep_waits_for_a_file_being_created() {
316 let dir = tempfile::tempdir().unwrap();
317 let unfinished = Unfinished::default();
318 let stop = Arc::new(AtomicBool::new(false));
319 let writer = unfinished.writer(stop.clone());
320 let (creating, wait) = std::sync::mpsc::channel();
321 let (go, gate) = std::sync::mpsc::channel::<()>();
322 let worker = {
323 let dir = dir.path().to_path_buf();
324 std::thread::spawn(move || {
325 writer
326 .create(|| {
327 creating.send(()).unwrap();
328 gate.recv().unwrap();
329 file_in(&dir)
330 })
331 .unwrap()
332 .is_none()
333 })
334 };
335 wait.recv().unwrap();
336 stop.store(true, Ordering::Relaxed);
337 let sweeper = {
338 let unfinished = unfinished.clone();
339 std::thread::spawn(move || unfinished.sweep(Instant::now() + Duration::from_secs(30)))
340 };
341 assert!(
342 unfinished.a_sweep_waits_within(Duration::from_secs(30)),
343 "the sweep waits for the file"
344 );
345 assert!(unfinished.writing());
346 assert!(!sweeper.is_finished());
347 go.send(()).unwrap();
348 sweeper.join().unwrap();
349 assert_eq!(files_in(dir.path()), 0, "the sweep waited for the file");
352 assert!(!unfinished.writing());
353 assert!(worker.join().unwrap(), "a stopped open's file is refused");
354 }
355
356 #[test]
360 fn a_file_its_holder_could_not_remove_is_removed_later() {
361 let dir = tempfile::tempdir().unwrap();
362 let unfinished = Unfinished::default();
363 let writer = unfinished.writer(Arc::default());
364 let left_over = |name: &str| {
366 let path = dir.path().join(name);
367 let (path, claim) = writer
368 .create(|| std::fs::write(&path, b"x").map(|()| path.clone()))
369 .unwrap()
370 .expect("not stopped yet");
371 drop(claim);
372 assert!(path.exists());
373 path
374 };
375 let first = left_over("first");
376 assert!(!unfinished.writing(), "nothing is writing it");
377 let (file, claim) = writer
378 .create(|| file_in(dir.path()))
379 .unwrap()
380 .expect("not stopped yet");
381 drop(file);
382 drop(claim);
383 assert!(!first.exists(), "the next claim tried it again");
384
385 let second = left_over("second");
386 let began = Instant::now();
387 unfinished.sweep(began + Duration::from_secs(30));
388 assert!(began.elapsed() < Duration::from_secs(10), "not waited on");
389 assert!(!second.exists(), "the sweep removed it");
390 assert_eq!(files_in(dir.path()), 0);
391 }
392
393 #[test]
395 fn a_stopped_or_swept_open_refuses_new_files() {
396 let dir = tempfile::tempdir().unwrap();
397 let unfinished = Unfinished::default();
398 let made = |writer: &Writer| {
399 writer
400 .create(|| file_in(dir.path()))
401 .unwrap()
402 .map(|(file, _)| file)
403 };
404 let stopped = unfinished.writer(Arc::new(AtomicBool::new(true)));
405 assert!(made(&stopped).is_none());
406
407 let writer = unfinished.writer(Arc::default());
408 unfinished.sweep(Instant::now());
409 assert!(made(&writer).is_none());
410 assert_eq!(files_in(dir.path()), 0);
411 }
412}