use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard};
use std::time::{Duration, Instant};
const SWEEP_GRACE: Duration = Duration::from_secs(1);
const RETRY_EVERY: Duration = Duration::from_millis(50);
#[derive(Clone, Debug, Default)]
pub(crate) struct Unfinished(Arc<(Mutex<Files>, Condvar)>);
#[derive(Debug, Default)]
struct Files {
claimed: Vec<PathBuf>,
creating: usize,
left_over: Vec<PathBuf>,
swept: bool,
#[cfg(test)]
sweep_waiting: bool,
}
impl Files {
fn busy(&self) -> bool {
!self.claimed.is_empty() || self.creating > 0
}
}
#[derive(Clone, Debug, Default)]
pub(crate) struct Writer {
stop: Arc<AtomicBool>,
files: Unfinished,
}
#[derive(Debug)]
#[must_use]
pub(crate) struct Claim {
files: Unfinished,
path: PathBuf,
}
impl Unfinished {
pub(crate) fn writer(&self, stop: Arc<AtomicBool>) -> Writer {
Writer {
stop,
files: self.clone(),
}
}
#[cfg(test)]
pub(crate) fn writing(&self) -> bool {
self.lock().busy()
}
pub(crate) fn sweep(&self, deadline: Instant) {
let (_, released) = &*self.0;
let mut files = self.lock();
while files.busy() {
let left = deadline.saturating_duration_since(Instant::now());
if left.is_zero() {
break;
}
#[cfg(test)]
{
files.sweep_waiting = true;
released.notify_all();
}
files = released
.wait_timeout(files, left)
.unwrap_or_else(|e| e.into_inner())
.0;
}
files.swept = true;
let mut left = std::mem::take(&mut files.claimed);
left.append(&mut files.left_over);
drop(files);
loop {
left.retain(|path| !removed(path));
if left.is_empty() || Instant::now() + RETRY_EVERY > deadline {
break;
}
std::thread::sleep(RETRY_EVERY);
}
for path in left {
log::warn!("Could not remove the temporary file {}", path.display());
}
}
#[cfg(test)]
fn a_sweep_waits_within(&self, limit: Duration) -> bool {
let (_, released) = &*self.0;
released
.wait_timeout_while(self.lock(), limit, |files| !files.sweep_waiting)
.unwrap_or_else(|e| e.into_inner())
.0
.sweep_waiting
}
fn lock(&self) -> MutexGuard<'_, Files> {
self.0.0.lock().unwrap_or_else(|e| e.into_inner())
}
fn released(&self) {
self.0.1.notify_all();
}
}
impl Writer {
pub(crate) fn stopped(&self) -> bool {
self.stop.load(Ordering::Relaxed)
}
pub(crate) fn create<F: AsRef<Path>, E>(
&self,
make: impl FnOnce() -> Result<F, E>,
) -> Result<Option<(F, Claim)>, E> {
{
let mut files = self.files.lock();
if files.swept || self.stopped() {
return Ok(None);
}
files.creating += 1;
}
let made = make();
let mut files = self.files.lock();
let claimed = match made {
Ok(file) if !files.swept && !self.stopped() => {
let path = file.as_ref().to_path_buf();
files.claimed.push(path.clone());
Ok(Some((
file,
Claim {
files: self.files.clone(),
path,
},
)))
}
Ok(file) => {
drop(files);
drop(file);
files = self.files.lock();
Ok(None)
}
Err(e) => Err(e),
};
files.creating -= 1;
drop(files);
self.files.released();
claimed
}
}
#[must_use]
pub struct ExitSweep(pub(crate) Unfinished);
impl Drop for ExitSweep {
fn drop(&mut self) {
self.0.sweep(Instant::now() + SWEEP_GRACE);
}
}
impl Drop for Claim {
fn drop(&mut self) {
let mut files = self.files.lock();
if let Some(at) = files.claimed.iter().position(|p| *p == self.path) {
files.claimed.swap_remove(at);
}
let earlier = std::mem::take(&mut files.left_over);
drop(files);
let mut left: Vec<PathBuf> = earlier.into_iter().filter(|p| !removed(p)).collect();
if std::fs::symlink_metadata(&self.path).is_ok() {
left.push(self.path.clone());
}
self.files.lock().left_over.append(&mut left);
self.files.released();
}
}
fn removed(path: &Path) -> bool {
match std::fs::remove_file(path) {
Ok(()) => true,
Err(e) => e.kind() == std::io::ErrorKind::NotFound,
}
}
#[cfg(test)]
mod tests;