use std::path::{Path, PathBuf};
use std::time::Duration;
const EPOCH_MARKER: &str = "pushkin-epoch:";
const PROBE_HEAD_BYTES: usize = 512;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RegenOutcome {
Regenerated,
TimedOut,
Failed(String),
}
pub fn probe_stale(dir: &Path, expected_epoch: u32) -> std::io::Result<Vec<PathBuf>> {
let mut stale = Vec::new();
let mut pending = vec![dir.to_path_buf()];
while let Some(current) = pending.pop() {
for entry in std::fs::read_dir(¤t)? {
let entry = entry?;
let path = entry.path();
if path.is_dir() {
pending.push(path);
} else if file_epoch(&path)? != Some(expected_epoch) {
if let Ok(relative) = path.strip_prefix(dir) {
stale.push(relative.to_path_buf());
}
}
}
}
Ok(stale)
}
fn file_epoch(path: &Path) -> std::io::Result<Option<u32>> {
use std::io::Read;
let mut head = vec![0_u8; PROBE_HEAD_BYTES];
let mut file = std::fs::File::open(path)?;
let read = file.read(&mut head)?;
head.truncate(read);
let text = String::from_utf8_lossy(&head);
Ok(text.lines().find_map(|line| {
let (_, rest) = line.split_once(EPOCH_MARKER)?;
rest.trim().parse::<u32>().ok()
}))
}
pub fn run_queue<W>(items: &[PathBuf], timeout: Duration, worker: W) -> Vec<(PathBuf, RegenOutcome)>
where
W: Fn(&Path) -> Result<(), String> + Send + Sync + 'static,
{
let worker = std::sync::Arc::new(worker);
let mut outcomes = Vec::with_capacity(items.len());
for item in items {
outcomes.push((item.clone(), run_one(item, timeout, &worker)));
}
outcomes
}
fn run_one<W>(item: &Path, timeout: Duration, worker: &std::sync::Arc<W>) -> RegenOutcome
where
W: Fn(&Path) -> Result<(), String> + Send + Sync + 'static,
{
let (sender, receiver) = std::sync::mpsc::channel();
let worker = std::sync::Arc::clone(worker);
let path = item.to_path_buf();
std::thread::spawn(move || {
let _ = sender.send(worker(&path));
});
match receiver.recv_timeout(timeout) {
Ok(Ok(())) => RegenOutcome::Regenerated,
Ok(Err(message)) => RegenOutcome::Failed(message),
Err(_) => RegenOutcome::TimedOut,
}
}