use std::path::{Path, PathBuf};
pub const STORE_LOCK: &str = ".locks/store";
pub fn store_dir(worker: &str) -> PathBuf {
crate::paths::worker_store_dir(crate::paths::short_worker(worker))
}
pub fn provision(worker: &str) -> Result<PathBuf, String> {
let dir = store_dir(worker);
std::fs::create_dir_all(dir.join(".locks"))
.map_err(|e| format!("store {}: {e}", dir.display()))?;
Ok(dir)
}
#[derive(Debug)]
pub struct SnapshotOutcome {
pub dir: PathBuf,
pub changes: usize,
}
fn snapshots_dir(worker: &str) -> PathBuf {
crate::paths::worker_store_snapshots_dir(crate::paths::short_worker(worker))
}
pub fn list_snapshots(worker: &str) -> Vec<String> {
let mut names: Vec<String> = std::fs::read_dir(snapshots_dir(worker))
.into_iter()
.flatten()
.flatten()
.filter(|e| e.path().is_dir())
.map(|e| e.file_name().to_string_lossy().into_owned())
.filter(|n| !n.ends_with(".tmp"))
.collect();
names.sort();
names.reverse();
names
}
pub fn snapshot(
worker: &str,
run_id: Option<&str>,
keep: usize,
) -> Result<Option<SnapshotOutcome>, String> {
if keep == 0 {
return Ok(None);
}
let store = store_dir(worker);
if !store.is_dir() {
return Ok(None);
}
let _lock = store_lock(worker)?;
let outcome = snapshot_locked(worker, &store, run_id)?;
prune(worker, keep);
Ok(outcome)
}
fn store_lock(worker: &str) -> Result<crate::statefile::FileLock, String> {
let path = store_dir(worker).join(STORE_LOCK);
crate::statefile::FileLock::acquire_path(&path).map_err(|e| format!("store lock: {e}"))
}
fn snapshot_locked(
worker: &str,
store: &Path,
run_id: Option<&str>,
) -> Result<Option<SnapshotOutcome>, String> {
let root = snapshots_dir(worker);
std::fs::create_dir_all(&root).map_err(|e| e.to_string())?;
let prev = list_snapshots(worker).into_iter().next();
let now = time::OffsetDateTime::now_utc();
let name = format!(
"{:04}{:02}{:02}-{:02}{:02}{:02}.{:03}",
now.year(),
u8::from(now.month()),
now.day(),
now.hour(),
now.minute(),
now.second(),
now.millisecond()
);
let final_dir = root.join(&name);
if final_dir.exists() {
return Ok(None);
}
let tmp = root.join(format!("{name}.tmp"));
let _ = std::fs::remove_dir_all(&tmp);
let mut cmd = std::process::Command::new("rsync");
cmd.arg("-a")
.arg("--delete")
.arg("--itemize-changes")
.arg("--exclude=/.locks")
.arg("--exclude=/.channel-stores");
if let Some(p) = &prev {
cmd.arg(format!("--link-dest={}", root.join(p).display()));
}
let out = cmd
.arg(format!("{}/", store.display()))
.arg(&tmp)
.output()
.map_err(|e| format!("rsync: {e} (is rsync installed?)"))?;
if !out.status.success() {
let _ = std::fs::remove_dir_all(&tmp);
return Err(format!(
"rsync failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
let channel_stores =
crate::paths::worker_channel_stores_dir(crate::paths::short_worker(worker));
let mut channel_out: Vec<u8> = Vec::new();
if channel_stores.is_dir() {
let mut cmd = std::process::Command::new("rsync");
cmd.arg("-a").arg("--delete").arg("--itemize-changes");
if let Some(p) = &prev {
let prev_sub = root.join(p).join(".channel-stores");
if prev_sub.is_dir() {
cmd.arg(format!("--link-dest={}", prev_sub.display()));
}
}
let out = cmd
.arg(format!("{}/", channel_stores.display()))
.arg(tmp.join(".channel-stores"))
.output()
.map_err(|e| format!("rsync: {e} (is rsync installed?)"))?;
if !out.status.success() {
let _ = std::fs::remove_dir_all(&tmp);
return Err(format!(
"rsync failed on channel stores: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
channel_out = out.stdout;
}
let is_itemized = |l: &&str| {
let b = l.as_bytes();
l.starts_with("*deleting")
|| (b.len() > 11
&& matches!(b[0], b'<' | b'>' | b'c' | b'h' | b'.')
&& matches!(b[1], b'f' | b'd' | b'L' | b'D' | b'S')
&& !l.starts_with(".d"))
};
let channel_text = String::from_utf8_lossy(&channel_out).into_owned();
let itemized: Vec<&str> = std::str::from_utf8(&out.stdout)
.unwrap_or_default()
.lines()
.filter(is_itemized)
.chain(channel_text.lines().filter(is_itemized))
.collect();
if prev.is_some() && itemized.is_empty() {
let _ = std::fs::remove_dir_all(&tmp);
return Ok(None);
}
std::fs::rename(&tmp, &final_dir).map_err(|e| e.to_string())?;
if let Some(rid) = run_id {
let _ = std::fs::write(
crate::paths::run_dir(rid).join("store-changes.txt"),
itemized.join("\n") + "\n",
);
}
Ok(Some(SnapshotOutcome {
dir: final_dir,
changes: itemized.len(),
}))
}
fn prune(worker: &str, keep: usize) {
let root = snapshots_dir(worker);
for name in list_snapshots(worker).into_iter().skip(keep) {
let _ = std::fs::remove_dir_all(root.join(name));
}
}
pub fn restore(
worker: &str,
from: Option<&str>,
channel: Option<&str>,
) -> Result<(PathBuf, Option<PathBuf>), String> {
let snaps = list_snapshots(worker);
let pick = match from {
Some(name) => snaps
.iter()
.find(|n| n.as_str() == name)
.ok_or_else(|| {
format!(
"no snapshot \"{name}\" — have: {}",
if snaps.is_empty() {
"none".into()
} else {
snaps.join(", ")
}
)
})?
.clone(),
None => snaps
.first()
.ok_or("no snapshots yet — nothing to restore from")?
.clone(),
};
let store = provision(worker)?;
let _lock = store_lock(worker)?;
let undo = snapshot_locked(worker, &store, None)?.map(|o| o.dir);
let src = snapshots_dir(worker).join(&pick);
let rsync = |from: &Path, to: &Path| -> Result<(), String> {
let out = std::process::Command::new("rsync")
.arg("-a")
.arg("--delete")
.arg("--exclude=/.locks")
.arg(format!("{}/", from.display()))
.arg(to)
.output()
.map_err(|e| format!("rsync: {e} (is rsync installed?)"))?;
if !out.status.success() {
return Err(format!(
"rsync failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
Ok(())
};
if let Some(channel) = channel {
let sub = src.join(".channel-stores").join(channel);
if !sub.is_dir() {
return Err(format!(
"snapshot {pick} has no channel store for \"{channel}\""
));
}
let dest =
crate::paths::worker_channel_store_dir(crate::paths::short_worker(worker), channel);
std::fs::create_dir_all(&dest).map_err(|e| e.to_string())?;
rsync(&sub, &dest)?;
return Ok((sub, undo));
}
let channel_sub = src.join(".channel-stores");
if channel_sub.is_dir() {
let dest = crate::paths::worker_channel_stores_dir(crate::paths::short_worker(worker));
std::fs::create_dir_all(&dest).map_err(|e| e.to_string())?;
rsync(&channel_sub, &dest)?;
}
let out = std::process::Command::new("rsync")
.arg("-a")
.arg("--delete")
.arg("--exclude=/.locks")
.arg("--exclude=/.channel-stores")
.arg(format!("{}/", src.display()))
.arg(&store)
.output()
.map_err(|e| format!("rsync: {e} (is rsync installed?)"))?;
if !out.status.success() {
return Err(format!(
"rsync failed: {}",
String::from_utf8_lossy(&out.stderr).trim()
));
}
Ok((src, undo))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn provision_is_idempotent_and_creates_locks() {
let _guard = crate::statefile::TEST_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
std::env::set_var("ROSTER_ROOT", dir.path());
let p = provision("dobby").unwrap();
assert!(p.join(".locks").is_dir());
let p2 = provision("org/dobby").unwrap();
assert_eq!(p, p2, "org/ prefix resolves to the same store");
}
#[test]
fn snapshot_rotate_restore_lifecycle() {
if std::process::Command::new("rsync")
.arg("--version")
.output()
.is_err()
{
eprintln!("skipping — rsync not installed");
return;
}
let _guard = crate::statefile::TEST_ENV_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
std::env::set_var("ROSTER_ROOT", dir.path());
let store = provision("dobby").unwrap();
std::fs::write(store.join("notes.md"), "v1").unwrap();
let first = snapshot("dobby", None, 5).unwrap().expect("first snapshot");
assert!(first.changes >= 1);
assert_eq!(
std::fs::read_to_string(first.dir.join("notes.md")).unwrap(),
"v1"
);
assert!(!first.dir.join(".locks").exists(), "locks are not content");
assert!(snapshot("dobby", None, 5).unwrap().is_none());
std::fs::write(store.join("notes.md"), "v2 with more text").unwrap();
let second = snapshot("dobby", None, 5)
.unwrap()
.expect("second snapshot");
assert_eq!(second.changes, 1);
assert_eq!(list_snapshots("dobby").len(), 2);
let older = list_snapshots("dobby").pop().unwrap();
std::fs::write(store.join("notes.md"), "wrecked-by-a-bad-run").unwrap();
let (from, undo) = restore("dobby", Some(&older), None).unwrap();
assert!(from.ends_with(&older));
assert!(undo.is_some(), "wrecked state was preserved for undo");
assert_eq!(
std::fs::read_to_string(store.join("notes.md")).unwrap(),
"v1"
);
assert!(store.join(".locks").is_dir(), "restore keeps the lock dir");
std::fs::write(store.join("notes.md"), "v3").unwrap();
snapshot("dobby", None, 1).unwrap().expect("third snapshot");
assert_eq!(list_snapshots("dobby").len(), 1);
std::fs::write(store.join("notes.md"), "v4").unwrap();
assert!(snapshot("dobby", None, 0).unwrap().is_none());
let chan = crate::paths::worker_channel_store_dir("dobby", "manas");
std::fs::create_dir_all(&chan).unwrap();
std::fs::write(chan.join("context.md"), "room notes v1").unwrap();
let with_chan = snapshot("dobby", None, 5)
.unwrap()
.expect("channel-store change snapshots");
assert_eq!(
std::fs::read_to_string(
with_chan
.dir
.join(".channel-stores")
.join("manas")
.join("context.md")
)
.unwrap(),
"room notes v1"
);
std::fs::write(chan.join("context.md"), "wrecked").unwrap();
std::fs::write(store.join("notes.md"), "store-stays").unwrap();
let (from, _) = restore("dobby", None, Some("manas")).unwrap();
assert!(from.ends_with("manas"));
assert_eq!(
std::fs::read_to_string(chan.join("context.md")).unwrap(),
"room notes v1"
);
assert_eq!(
std::fs::read_to_string(store.join("notes.md")).unwrap(),
"store-stays",
"a channel restore leaves the global store alone"
);
}
}