use anyhow::Result;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
pub const DEFAULT_BACKGROUND_PERMITS: usize = 3;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Permit {
pub pid: u32,
pub taken_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub what: Option<String>,
}
pub struct Permits {
dir: PathBuf,
capacity: usize,
}
pub struct Held {
path: PathBuf,
}
impl Drop for Held {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
impl Permits {
pub fn new(dir: impl Into<PathBuf>, capacity: usize) -> Self {
Permits {
dir: dir.into(),
capacity,
}
}
pub fn dir(&self) -> &Path {
&self.dir
}
pub fn live(&self) -> Vec<Permit> {
let Ok(entries) = std::fs::read_dir(&self.dir) else {
return Vec::new();
};
let mut held = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("permit") {
continue;
}
let parsed = std::fs::read_to_string(&path)
.ok()
.and_then(|t| serde_json::from_str::<Permit>(&t).ok());
match parsed {
Some(p) if crate::process_alive(p.pid) => held.push(p),
_ => {
let _ = std::fs::remove_file(&path);
}
}
}
held
}
pub fn take(&self, what: &str) -> Result<Result<Held, Vec<Permit>>> {
let held = self.live();
if held.len() >= self.capacity {
return Ok(Err(held));
}
crate::create_private_dir(&self.dir)?;
let permit = Permit {
pid: std::process::id(),
taken_at: Utc::now(),
what: (!what.is_empty()).then(|| what.to_string()),
};
let path = self.dir.join(format!("{}.permit", std::process::id()));
std::fs::write(&path, serde_json::to_string_pretty(&permit)?)?;
Ok(Ok(Held { path }))
}
pub fn capacity(&self) -> usize {
self.capacity
}
}
#[cfg(test)]
mod tests {
use super::*;
fn scratch(name: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"mecha-permit-{name}-{}-{:?}",
std::process::id(),
std::thread::current().id()
));
let _ = std::fs::remove_dir_all(&dir);
dir
}
#[test]
fn a_permit_is_held_until_dropped() {
let p = Permits::new(scratch("basic"), 2);
assert!(p.live().is_empty());
let one = p.take("task-a").unwrap().expect("a free seat");
assert_eq!(p.live().len(), 1);
drop(one);
assert!(p.live().is_empty(), "released on drop, not by a call");
}
#[test]
fn a_full_pool_refuses_and_says_who_has_it() {
let p = Permits::new(scratch("full"), 1);
let _one = p.take("task-a").unwrap().expect("a free seat");
match p.take("task-b").unwrap() {
Ok(_) => panic!("capacity ignored"),
Err(held) => {
assert_eq!(held.len(), 1);
assert_eq!(held[0].what.as_deref(), Some("task-a"));
}
}
}
#[test]
fn a_permit_whose_process_is_gone_is_reclaimed() {
let p = Permits::new(scratch("dead"), 1);
crate::create_private_dir(p.dir()).unwrap();
std::fs::write(
p.dir().join("0.permit"),
serde_json::json!({"pid": 0, "taken_at": Utc::now().to_rfc3339()}).to_string(),
)
.unwrap();
assert!(p.live().is_empty(), "a dead pid is not a holder");
assert!(
p.take("task-a").unwrap().is_ok(),
"and its seat is available again"
);
}
#[test]
fn an_unreadable_permit_is_swept_rather_than_counted() {
let p = Permits::new(scratch("junk"), 1);
crate::create_private_dir(p.dir()).unwrap();
std::fs::write(p.dir().join("x.permit"), "{not json").unwrap();
assert!(p.live().is_empty());
assert!(p.take("task-a").unwrap().is_ok());
}
#[test]
fn a_stray_file_is_not_a_holder() {
let p = Permits::new(scratch("stray"), 1);
crate::create_private_dir(p.dir()).unwrap();
std::fs::write(p.dir().join(".DS_Store"), "junk").unwrap();
assert!(p.live().is_empty());
assert!(p.take("task-a").unwrap().is_ok());
}
}