Skip to main content

datui_lib/loading/
unfinished.rs

1//! The temporary files an open's workers write, so quitting removes them. Such a file
2//! removes itself when its last holder drops it, but the process can end before a
3//! stopped worker gets there. So each writer [creates](Writer::create) its file through
4//! the open, which claims it for the file's lifetime. On quit the app drops first, then
5//! [`ExitSweep`] gives workers a short grace and removes whatever is still claimed
6//! (waiting for files mid-creation; later ones are refused). On Windows a file mapped by
7//! a frame cannot be removed, so its claim keeps the path for later claims and the exit
8//! sweep to retry. SIGKILL skips all this.
9
10use std::path::{Path, PathBuf};
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{Arc, Condvar, Mutex, MutexGuard};
13use std::time::{Duration, Instant};
14
15/// How long the sweep at exit waits for writers to stop before removing their files.
16const SWEEP_GRACE: Duration = Duration::from_secs(1);
17
18/// How often the sweep tries a file that would not go again, within its grace.
19const RETRY_EVERY: Duration = Duration::from_millis(50);
20
21/// The files claimed by the writers of one loader's opens.
22#[derive(Clone, Debug, Default)]
23pub(crate) struct Unfinished(Arc<(Mutex<Files>, Condvar)>);
24
25#[derive(Debug, Default)]
26struct Files {
27    claimed: Vec<PathBuf>,
28    /// Writers creating a file, not yet claimed: the sweep waits for them too, since
29    /// the path is not known until the file exists.
30    creating: usize,
31    /// Files whose holder let go but could not remove them. Not waited for: nothing
32    /// is writing them.
33    left_over: Vec<PathBuf>,
34    swept: bool,
35    /// Set by a sweep as it starts waiting, so a test can tell it is blocked.
36    #[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/// What one open hands its writers: its stop flag, and where to claim their files.
47#[derive(Clone, Debug, Default)]
48pub(crate) struct Writer {
49    stop: Arc<AtomicBool>,
50    files: Unfinished,
51}
52
53/// A hold on a file a writer created, kept by whatever holds the file and dropped
54/// after it.
55#[derive(Debug)]
56#[must_use]
57pub(crate) struct Claim {
58    files: Unfinished,
59    path: PathBuf,
60}
61
62impl Unfinished {
63    /// The writer for an open whose stop flag is `stop`.
64    pub(crate) fn writer(&self, stop: Arc<AtomicBool>) -> Writer {
65        Writer {
66            stop,
67            files: self.clone(),
68        }
69    }
70
71    /// Whether a file is still claimed or being created.
72    #[cfg(test)]
73    pub(crate) fn writing(&self) -> bool {
74        self.lock().busy()
75    }
76
77    /// By `deadline`, remove every file still claimed; stopped writers remove their own,
78    /// and one busy at the deadline loses its file (and harmlessly fails to remove it). Files
79    /// mid-creation are waited for; later ones are refused.
80    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            // Set under the lock that `wait_timeout` releases, so whoever sees it next
89            // sees a sweep already waiting.
90            #[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    /// Whether a sweep starts waiting on a busy writer within `limit`. Bounded, so a
117    /// sweep that never waits fails the test instead of hanging it.
118    #[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    /// Whether the open was stopped: a writer gives up at its next chunk.
139    pub(crate) fn stopped(&self) -> bool {
140        self.stop.load(Ordering::Relaxed)
141    }
142
143    /// Create a file with `make` and claim it, as one step to the sweep (a sweep during
144    /// `make` waits, then removes it or finds it gone; claiming afterward would leave a
145    /// gap). `None`, with nothing on disk, once the open is stopped or swept.
146    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                // Removed before the sweep stops waiting for it, and not under the lock.
173                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/// Dropped after the app, removes what the app's opens were still writing. See the
188/// module.
189#[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        // Ones left over earlier may have been let go since: a frame dropped.
205        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        // Its holder removed it before letting go, unless that failed.
209        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
217/// Whether `path` is gone, removing it if it is there.
218fn 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;