use std::fs::{self, File};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunRecord {
pub pid: u32,
#[serde(rename = "startedIso")]
pub started_iso: String,
pub backend: String,
#[serde(rename = "msbPath", skip_serializing_if = "Option::is_none", default)]
pub msb_path: Option<String>,
}
static WRITE_LOCK: Mutex<()> = Mutex::new(());
#[derive(Debug, Clone)]
pub struct Ledger {
runs_dir: PathBuf,
run_id: String,
}
impl Ledger {
pub fn new(cache_dir: &Path, run_id: &str) -> Self {
Ledger {
runs_dir: cache_dir.join("runs"),
run_id: run_id.to_string(),
}
}
pub fn record_path(&self) -> PathBuf {
self.runs_dir.join(format!("{}.json", self.run_id))
}
pub fn sandboxes_path(&self) -> PathBuf {
self.runs_dir.join(format!("{}.sandboxes", self.run_id))
}
pub fn networks_path(&self) -> PathBuf {
self.runs_dir.join(format!("{}.networks", self.run_id))
}
pub fn write_record(&self, record: &RunRecord) -> std::io::Result<()> {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
fs::create_dir_all(&self.runs_dir)?;
let json =
serde_json::to_vec_pretty(record).expect("RunRecord has no non-serializable fields");
let tmp = self.runs_dir.join(format!("{}.json.tmp", self.run_id));
{
let mut f = File::create(&tmp)?;
f.write_all(&json)?;
}
fs::rename(&tmp, self.record_path())?;
Ok(())
}
pub fn read_record(&self) -> Option<RunRecord> {
let raw = fs::read(self.record_path()).ok()?;
serde_json::from_slice(&raw).ok()
}
pub fn append_sandbox(&self, name: &str) -> std::io::Result<()> {
self.append_line(&self.sandboxes_path(), name)
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn remove_sandbox(&self, name: &str) -> std::io::Result<()> {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
self.remove_line_locked(&self.sandboxes_path(), name)
}
pub fn remove_sandbox_and_prune_if_empty(&self, name: &str) {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
let _ = self.remove_line_locked(&self.sandboxes_path(), name);
if self.is_empty_locked() {
self.delete_files_locked();
}
}
pub fn sandbox_names(&self) -> Vec<String> {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
read_lines(&self.sandboxes_path())
}
pub fn append_network(&self, id: &str) -> std::io::Result<()> {
self.append_line(&self.networks_path(), id)
}
pub fn remove_network_and_prune_if_empty(&self, id: &str) {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
let _ = self.remove_line_locked(&self.networks_path(), id);
if self.is_empty_locked() {
self.delete_files_locked();
}
}
pub fn network_ids(&self) -> Vec<String> {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
read_lines(&self.networks_path())
}
#[cfg_attr(not(test), allow(dead_code))]
pub fn is_empty(&self) -> bool {
self.sandbox_names().is_empty() && self.network_ids().is_empty()
}
pub fn delete_files(&self) {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
self.delete_files_locked();
}
fn is_empty_locked(&self) -> bool {
read_lines(&self.sandboxes_path()).is_empty()
&& read_lines(&self.networks_path()).is_empty()
}
fn delete_files_locked(&self) {
let _ = fs::remove_file(self.record_path());
let _ = fs::remove_file(self.sandboxes_path());
let _ = fs::remove_file(self.networks_path());
}
fn append_line(&self, path: &Path, line: &str) -> std::io::Result<()> {
let _guard = WRITE_LOCK.lock().expect("ledger write lock poisoned");
self.append_line_locked(path, line)
}
fn append_line_locked(&self, path: &Path, line: &str) -> std::io::Result<()> {
fs::create_dir_all(&self.runs_dir)?;
let mut existing = read_lines(path);
if existing.iter().any(|l| l == line) {
return Ok(());
}
existing.push(line.to_string());
write_lines(path, &existing)
}
fn remove_line_locked(&self, path: &Path, line: &str) -> std::io::Result<()> {
let mut existing = read_lines(path);
let before = existing.len();
existing.retain(|l| l != line);
if existing.len() == before {
return Ok(()); }
write_lines(path, &existing)
}
}
fn read_lines(path: &Path) -> Vec<String> {
match fs::read_to_string(path) {
Ok(text) => text
.lines()
.map(str::trim)
.filter(|l| !l.is_empty())
.map(str::to_string)
.collect(),
Err(_) => Vec::new(),
}
}
fn write_lines(path: &Path, lines: &[String]) -> std::io::Result<()> {
let mut content = lines.join("\n");
if !lines.is_empty() {
content.push('\n');
}
let tmp = tmp_sibling(path);
fs::write(&tmp, content)?;
fs::rename(&tmp, path)
}
fn tmp_sibling(path: &Path) -> PathBuf {
let mut tmp = path.as_os_str().to_os_string();
tmp.push(".tmp");
PathBuf::from(tmp)
}
pub(crate) fn candidate_run_ids(cache_dir: &Path) -> Vec<String> {
let runs_dir = cache_dir.join("runs");
let Ok(entries) = fs::read_dir(&runs_dir) else {
return Vec::new();
};
let mut ids = Vec::new();
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|e| e.to_str()) != Some("json") {
continue;
}
if let Some(stem) = path.file_stem().and_then(|s| s.to_str()) {
ids.push(stem.to_string());
}
}
ids
}
pub(crate) fn file_age(path: &Path) -> Option<std::time::Duration> {
fs::metadata(path).ok()?.modified().ok()?.elapsed().ok()
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
fn temp_cache_dir(label: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"rz-ledger-{label}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&dir).unwrap();
dir
}
fn sample_record() -> RunRecord {
RunRecord {
pid: 12345,
started_iso: "2025-01-01T00:00:00Z".to_string(),
backend: "docker".to_string(),
msb_path: None,
}
}
#[test]
fn write_record_then_read_record_round_trips() {
let cache = temp_cache_dir("write-read");
let ledger = Ledger::new(&cache, "run1");
ledger.write_record(&sample_record()).unwrap();
assert_eq!(ledger.read_record(), Some(sample_record()));
}
#[test]
fn record_json_uses_camel_case_field_names() {
let cache = temp_cache_dir("camel-case");
let ledger = Ledger::new(&cache, "run1");
let mut record = sample_record();
record.msb_path = Some("/opt/msb".to_string());
ledger.write_record(&record).unwrap();
let raw = fs::read_to_string(ledger.record_path()).unwrap();
assert!(raw.contains("\"startedIso\""), "{raw}");
assert!(raw.contains("\"msbPath\""), "{raw}");
assert!(!raw.contains("started_iso"), "{raw}");
}
#[test]
fn msb_path_is_omitted_from_json_when_none() {
let cache = temp_cache_dir("omit-msbpath");
let ledger = Ledger::new(&cache, "run1");
ledger.write_record(&sample_record()).unwrap();
let raw = fs::read_to_string(ledger.record_path()).unwrap();
assert!(!raw.contains("msbPath"), "{raw}");
}
#[test]
fn append_sandbox_is_before_create_and_dedupes() {
let cache = temp_cache_dir("append-dedup");
let ledger = Ledger::new(&cache, "run1");
ledger.append_sandbox("rz-run1-0").unwrap();
ledger.append_sandbox("rz-run1-1").unwrap();
ledger.append_sandbox("rz-run1-0").unwrap(); assert_eq!(
ledger.sandbox_names(),
vec!["rz-run1-0".to_string(), "rz-run1-1".to_string()]
);
}
#[test]
fn remove_sandbox_after_stop_leaves_only_the_survivors() {
let cache = temp_cache_dir("remove-after-stop");
let ledger = Ledger::new(&cache, "run1");
ledger.append_sandbox("a").unwrap();
ledger.append_sandbox("b").unwrap();
ledger.remove_sandbox("a").unwrap();
assert_eq!(ledger.sandbox_names(), vec!["b".to_string()]);
}
#[test]
fn removing_an_absent_sandbox_is_a_harmless_no_op() {
let cache = temp_cache_dir("remove-absent");
let ledger = Ledger::new(&cache, "run1");
ledger.remove_sandbox("never-there").unwrap();
assert!(ledger.sandbox_names().is_empty());
}
#[test]
fn is_empty_true_only_when_both_files_have_no_entries() {
let cache = temp_cache_dir("is-empty");
let ledger = Ledger::new(&cache, "run1");
assert!(ledger.is_empty(), "no files at all: empty");
ledger.append_sandbox("a").unwrap();
assert!(!ledger.is_empty());
ledger.remove_sandbox("a").unwrap();
assert!(ledger.is_empty());
ledger.append_network("net-1").unwrap();
assert!(!ledger.is_empty());
}
#[test]
fn delete_files_removes_all_three_and_is_idempotent() {
let cache = temp_cache_dir("delete-files");
let ledger = Ledger::new(&cache, "run1");
ledger.write_record(&sample_record()).unwrap();
ledger.append_sandbox("a").unwrap();
ledger.append_network("n").unwrap();
ledger.delete_files();
assert!(!ledger.record_path().exists());
assert!(!ledger.sandboxes_path().exists());
assert!(!ledger.networks_path().exists());
ledger.delete_files(); }
#[test]
fn candidate_run_ids_lists_every_json_stem_and_ignores_other_extensions() {
let cache = temp_cache_dir("candidates");
Ledger::new(&cache, "runA")
.write_record(&sample_record())
.unwrap();
Ledger::new(&cache, "runB")
.write_record(&sample_record())
.unwrap();
let mut ids = candidate_run_ids(&cache);
ids.sort();
assert_eq!(ids, vec!["runA".to_string(), "runB".to_string()]);
}
#[test]
fn candidate_run_ids_on_a_missing_runs_dir_is_empty_not_an_error() {
let cache = temp_cache_dir("no-runs-dir");
assert!(candidate_run_ids(&cache).is_empty());
}
#[test]
fn concurrent_append_never_lands_in_the_remove_then_prune_gap() {
let cache = temp_cache_dir("remove-prune-race");
let ledger = Arc::new(Ledger::new(&cache, "run1"));
const ITERATIONS: usize = 200;
for i in 0..ITERATIONS {
let fresh_name = format!("fresh-{i}");
ledger.append_sandbox("victim").unwrap();
assert!(
ledger.network_ids().is_empty(),
"sanity: no networks tracked"
);
let barrier = Arc::new(std::sync::Barrier::new(2));
let remover = {
let ledger = ledger.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
barrier.wait();
ledger.remove_sandbox_and_prune_if_empty("victim");
})
};
let appender = {
let ledger = ledger.clone();
let barrier = barrier.clone();
let fresh_name = fresh_name.clone();
std::thread::spawn(move || {
barrier.wait();
ledger.append_sandbox(&fresh_name).unwrap();
})
};
remover
.join()
.expect("remove_sandbox_and_prune_if_empty must not panic");
appender.join().expect("append_sandbox must not panic");
assert_eq!(
ledger.sandbox_names(),
vec![fresh_name.clone()],
"iteration {i}: a concurrent append_sandbox racing \
remove_sandbox_and_prune_if_empty must never be discarded, \
regardless of which of the two operations the lock let run first"
);
ledger.remove_sandbox_and_prune_if_empty(&fresh_name);
assert!(ledger.sandbox_names().is_empty());
}
}
#[test]
fn write_lines_truncate_then_write_gap_is_visible_to_an_unlocked_reader() {
use std::fs::OpenOptions;
use std::io::Write as _;
let cache = temp_cache_dir("torn-read-mechanism");
let path = cache.join("mechanism.sandboxes");
fs::write(&path, "existing-line\n").unwrap();
let barrier = Arc::new(std::sync::Barrier::new(2));
let saw_empty = Arc::new(std::sync::atomic::AtomicBool::new(false));
let writer = {
let path = path.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
let mut f = OpenOptions::new()
.write(true)
.truncate(true)
.open(&path)
.unwrap();
barrier.wait(); std::thread::sleep(std::time::Duration::from_millis(50));
f.write_all(b"existing-line\nfresh-line\n").unwrap();
})
};
let reader = {
let path = path.clone();
let barrier = barrier.clone();
let saw_empty = saw_empty.clone();
std::thread::spawn(move || {
barrier.wait(); let seen = read_lines(&path);
eprintln!("[probe] unlocked reader saw: {seen:?}");
if seen.is_empty() {
saw_empty.store(true, std::sync::atomic::Ordering::SeqCst);
}
})
};
writer.join().unwrap();
reader.join().unwrap();
assert!(
saw_empty.load(std::sync::atomic::Ordering::SeqCst),
"an unlocked read landing between a truncate-then-write's two steps \
must observe a transiently empty file — demonstrates the exact gap \
the OLD write_lines exposed to sandbox_names/network_ids/is_empty"
);
assert_eq!(
read_lines(&path),
vec!["existing-line".to_string(), "fresh-line".to_string()],
"sanity: the write does eventually land correctly, after the gap"
);
}
#[test]
fn concurrent_writes_and_unlocked_reads_never_expose_a_torn_file() {
let cache = temp_cache_dir("torn-read-fixed");
let path = cache.join("fixed.sandboxes");
const WRITES: usize = 40;
let big_line = "x".repeat(64 * 1024);
let payloads: Vec<Vec<String>> = (0..WRITES)
.map(|i| vec![format!("{big_line}-{i}"); 32])
.collect();
write_lines(&path, &payloads[0]).unwrap();
let done = Arc::new(std::sync::atomic::AtomicBool::new(false));
let bad_reads = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let writer = {
let path = path.clone();
let payloads = payloads.clone();
let done = done.clone();
std::thread::spawn(move || {
for payload in &payloads[1..] {
write_lines(&path, payload).unwrap();
}
done.store(true, std::sync::atomic::Ordering::SeqCst);
})
};
let reader = {
let path = path.clone();
let done = done.clone();
let bad_reads = bad_reads.clone();
std::thread::spawn(move || {
while !done.load(std::sync::atomic::Ordering::SeqCst) {
let seen = read_lines(&path);
if seen.is_empty() {
bad_reads.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
}
})
};
writer.join().unwrap();
reader.join().unwrap();
assert_eq!(
bad_reads.load(std::sync::atomic::Ordering::SeqCst),
0,
"an unlocked reader must never observe an empty/torn file while a \
concurrent writer holds non-empty content, now that write_lines is \
atomic (tmp file + rename)"
);
assert_eq!(
read_lines(&path),
payloads[WRITES - 1],
"the file must end up exactly at the last write's content"
);
}
}