use agentd::store::file::FileStore;
use agentd::store::{Envelope, PutOutcome, Store, StoreError};
use serde_json::{Value, json};
use std::collections::BTreeSet;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
fn env(kind: &str, id: &str, seq: u64, state: Value) -> Value {
Envelope::new(kind, id, seq, "inst", None, state).to_value()
}
fn key(kind: &str, id: &str) -> String {
agentd::store::key("agentd", "inst", kind, id)
}
fn tree(root: &Path) -> BTreeSet<String> {
fn walk(dir: &Path, root: &Path, out: &mut BTreeSet<String>) {
let Ok(rd) = fs::read_dir(dir) else { return };
for e in rd.flatten() {
let p = e.path();
if p.is_dir() {
walk(&p, root, out);
} else {
out.insert(
p.strip_prefix(root)
.unwrap_or(&p)
.to_string_lossy()
.replace('\\', "/"),
);
}
}
}
let mut out = BTreeSet::new();
walk(root, root, &mut out);
out
}
fn enc(seg: &str) -> String {
seg.bytes()
.map(|b| {
if b.is_ascii_alphanumeric() || b == b'-' || b == b'_' {
(b as char).to_string()
} else {
format!("%{b:02X}")
}
})
.collect()
}
fn expected_file(k: &str) -> String {
let segs: Vec<&str> = k.split('/').filter(|s| !s.is_empty()).collect();
let mut parts: Vec<String> = segs.iter().map(|s| enc(s)).collect();
let last = parts.len() - 1;
parts[last] = format!("{}.json", parts[last]);
parts.join("/")
}
fn tmp() -> tempfile::TempDir {
tempfile::tempdir().expect("tempdir")
}
#[test]
fn round_trip_put_get_list_delete_cas_and_tombstone() {
let td = tmp();
let root = td.path().join("state");
let s = FileStore::open(&root).expect("open");
assert_eq!(
s.kind(),
"file",
"the adapter identifies itself for metrics"
);
let k = key("run", "01M06");
assert_eq!(
s.get(&k, None).unwrap(),
None,
"an absent key reads as None"
);
assert_eq!(
s.put(&k, 1, &env("run", "01M06", 1, json!({"step": "a"})))
.unwrap(),
PutOutcome::Ok
);
let got = s.get(&k, None).unwrap().expect("stored");
let e = Envelope::from_value(got).expect("parses as a v2 envelope");
assert_eq!(
(e.seq, e.kind.as_str(), e.instance.as_str()),
(1, "run", "inst")
);
assert_eq!(e.state, json!({"step": "a"}));
for stale in [0u64, 1] {
assert_eq!(
s.put(
&k,
stale,
&env("run", "01M06", stale, json!({"step": "clobber"}))
)
.unwrap(),
PutOutcome::Conflict {
latest_seq: Some(1)
},
"seq {stale} must not overwrite the stored seq 1"
);
}
let e = Envelope::from_value(s.get(&k, None).unwrap().unwrap()).unwrap();
assert_eq!(e.state, json!({"step": "a"}), "a conflict writes nothing");
assert_eq!(
s.put(&k, 2, &env("run", "01M06", 2, json!({"step": "b"})))
.unwrap(),
PutOutcome::Ok
);
for pin in [None, Some(1), Some(2)] {
let e = Envelope::from_value(s.get(&k, pin).unwrap().unwrap()).unwrap();
assert_eq!((e.seq, &e.state), (2, &json!({"step": "b"})), "pin {pin:?}");
}
let k2 = key("memory", "notes");
s.put(&k2, 7, &env("memory", "notes", 7, json!(["one"])))
.unwrap();
let listed = s.list("agentd/inst").unwrap();
let pairs: Vec<(String, Option<u64>)> = listed.iter().map(|e| (e.key.clone(), e.seq)).collect();
assert_eq!(
pairs,
vec![(k2.clone(), Some(7)), (k.clone(), Some(2))],
"list returns decoded keys, sorted, with seqs"
);
assert_eq!(
s.list("agentd/inst/run").unwrap().len(),
1,
"the prefix selects one kind"
);
assert!(
s.list("agentd/nobody").unwrap().is_empty(),
"an unknown prefix is empty, not an error"
);
assert_eq!(
s.put(&k, 3, &env("run", "01M06", 3, Value::Null)).unwrap(),
PutOutcome::Ok
);
assert_eq!(
s.get(&k, None).unwrap(),
None,
"a tombstone reads as absent"
);
assert_eq!(
s.list("agentd/inst/run").unwrap()[0].seq,
Some(3),
"the tombstone still fences later writers at seq 3"
);
assert_eq!(
s.put(&k, 3, &env("run", "01M06", 3, json!({"step": "c"})))
.unwrap(),
PutOutcome::Conflict {
latest_seq: Some(3)
},
"a tombstone is a CAS floor like any other record"
);
s.delete(&k).unwrap();
assert_eq!(s.get(&k, None).unwrap(), None);
assert!(
s.list("agentd/inst/run").unwrap().is_empty(),
"the key is gone from list"
);
s.delete(&k).unwrap();
assert!(s.get(&k2, None).unwrap().is_some());
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let f = root.join(expected_file(&k2));
assert_eq!(
fs::metadata(&f).unwrap().permissions().mode() & 0o777,
0o600,
"state files are 0600: they hold conversation content"
);
assert_eq!(
fs::metadata(f.parent().unwrap())
.unwrap()
.permissions()
.mode()
& 0o777,
0o700,
"state directories are 0700"
);
}
assert_eq!(
tree(&root),
BTreeSet::from([".lock".to_string(), expected_file(&k2)]),
"one file per live key, plus the instance lock"
);
}
#[test]
fn a_reader_racing_a_writer_never_sees_a_partial_file() {
let td = tmp();
let root = td.path().join("state");
let s = Arc::new(FileStore::open(&root).expect("open"));
let k = key("run", "hot");
let path = root.join(expected_file(&k));
let stop = Arc::new(AtomicBool::new(false));
let reads = Arc::new(AtomicU64::new(0));
let small = Arc::new(AtomicU64::new(0));
let large = Arc::new(AtomicU64::new(0));
let reader = {
let (stop, reads, small, large, path) = (
stop.clone(),
reads.clone(),
small.clone(),
large.clone(),
path.clone(),
);
std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
let Ok(bytes) = fs::read(&path) else { continue }; let v: Value = serde_json::from_slice(&bytes).unwrap_or_else(|e| {
panic!(
"torn read: {} bytes at {} did not parse ({e}) — the write was not atomic",
bytes.len(),
path.display()
)
});
let declared = v["state"]["len"].as_u64().expect("state.len") as usize;
let actual = v["state"]["filler"].as_str().expect("state.filler").len();
assert_eq!(
declared, actual,
"a published envelope is internally inconsistent: declared {declared}, \
got {actual} — a partial file became visible"
);
reads.fetch_add(1, Ordering::Relaxed);
if actual > 100_000 { &large } else { &small }.fetch_add(1, Ordering::Relaxed);
}
})
};
for seq in 1..=240u64 {
let n = if seq % 2 == 0 { 10 } else { 128 * 1024 };
let filler = "x".repeat(n);
let st = json!({"len": n, "filler": filler});
assert_eq!(
s.put(&k, seq, &env("run", "hot", seq, st)).unwrap(),
PutOutcome::Ok
);
}
stop.store(true, Ordering::Relaxed);
reader.join().expect("the reader saw only whole envelopes");
let (r, sm, lg) = (
reads.load(Ordering::Relaxed),
small.load(Ordering::Relaxed),
large.load(Ordering::Relaxed),
);
assert!(
r > 50,
"only {r} reads landed — the race did not happen, the test proved nothing"
);
assert!(
sm > 0 && lg > 0,
"saw {sm} small / {lg} large envelopes — sizes did not interleave"
);
let leftovers: Vec<String> = tree(&root)
.into_iter()
.filter(|p| p.contains(".tmp."))
.collect();
assert!(
leftovers.is_empty(),
"temp files survived completed writes: {leftovers:?}"
);
}
#[test]
fn a_crash_between_tmp_and_rename_leaves_the_previous_envelope_intact() {
let _no_forks = FORK_FREE.lock().unwrap_or_else(|e| e.into_inner());
let td = tmp();
let root = td.path().join("state");
let s = FileStore::open(&root).expect("open");
let k = key("run", "crashy");
s.put(
&k,
4,
&env("run", "crashy", 4, json!({"step": "committed"})),
)
.unwrap();
let file = root.join(expected_file(&k));
let dir = file.parent().unwrap();
let residue = dir.join(format!(
".{}.tmp.{}",
file.file_name().unwrap().to_string_lossy(),
424242
));
let partial =
serde_json::to_vec(&env("run", "crashy", 5, json!({"step": "interrupted"}))).unwrap();
fs::write(&residue, &partial[..partial.len() / 2]).unwrap();
drop(s);
let s = FileStore::open(&root).expect("reopen after the crash");
let e =
Envelope::from_value(s.get(&k, None).unwrap().expect("the committed envelope")).unwrap();
assert_eq!(
(e.seq, &e.state),
(4, &json!({"step": "committed"})),
"seq 4 survived intact"
);
assert_eq!(
s.list("agentd/inst").unwrap().len(),
1,
"the truncated temp file is not a key — list must not surface it"
);
assert_eq!(
s.put(&k, 5, &env("run", "crashy", 5, json!({"step": "redone"})))
.unwrap(),
PutOutcome::Ok,
"the interrupted write can simply be redone: the CAS floor is still 4"
);
let e = Envelope::from_value(s.get(&k, None).unwrap().unwrap()).unwrap();
assert_eq!(e.state, json!({"step": "redone"}));
}
#[test]
fn hostile_ids_never_write_outside_the_root() {
let td = tmp();
let base = td.path().to_path_buf();
let root = base.join("a/b/c/state");
let s = FileStore::open(&root).expect("open");
let before: BTreeSet<String> = tree(&base);
assert_eq!(before, BTreeSet::from(["a/b/c/state/.lock".to_string()]));
let contained = [
"../../../../../../etc/passwd", "..", ".", "/etc/passwd", "/", "a/b", "..%2f..%2fpwned", "....//....//pwned", "sub/../../../pwned", "-", "ünïcøde", " leading and trailing ", "CON", ];
let mut expected: BTreeSet<String> = BTreeSet::from([".lock".to_string()]);
for id in contained {
let k = key("run", id);
assert_eq!(
s.put(&k, 1, &env("run", id, 1, json!({"id": id}))).unwrap(),
PutOutcome::Ok,
"id {id:?} is a legal opaque id"
);
expected.insert(expected_file(&k));
let e = Envelope::from_value(s.get(&k, None).unwrap().expect("readable")).unwrap();
assert_eq!(e.state, json!({"id": id}), "id {id:?} round-trips");
}
let k = key("run", "nul\0byte");
match s.put(&k, 1, &env("run", "x", 1, json!(1))) {
Err(StoreError::Mapping(m)) => assert!(m.contains("NUL"), "unhelpful message: {m}"),
other => panic!("a NUL id must be refused, got {other:?}"),
}
match s.get(&k, None) {
Err(StoreError::Mapping(_)) => {}
other => panic!("a NUL id must be refused on read too, got {other:?}"),
}
assert!(matches!(
s.put("", 1, &json!({})),
Err(StoreError::Mapping(_))
));
assert!(matches!(
s.put("///", 1, &json!({})),
Err(StoreError::Mapping(_))
));
let after = tree(&root);
assert_eq!(
after, expected,
"the root holds exactly the keys that were put"
);
let outside: Vec<String> = tree(&base)
.into_iter()
.filter(|p| !p.starts_with("a/b/c/state/"))
.collect();
assert!(
outside.is_empty(),
"files appeared outside the root: {outside:?}"
);
for probe in [
"a/b/c/passwd",
"a/b/pwned",
"a/pwned",
"pwned",
"passwd",
"etc",
] {
assert!(!base.join(probe).exists(), "an escape landed at {probe}");
}
let croot = root.canonicalize().unwrap();
for rel in &after {
let p = root.join(rel).canonicalize().unwrap();
assert!(
p.starts_with(&croot),
"{} is not under {}",
p.display(),
croot.display()
);
assert!(
!p.components()
.any(|c| c.as_os_str() == ".." || c.as_os_str() == "."),
"{} still carries a relative component",
p.display()
);
}
assert_eq!(
after.len(),
contained.len() + 1,
"one file per hostile id (plus .lock): ids collided on disk"
);
let listed = s.list("agentd/inst").unwrap();
let keys: BTreeSet<String> = listed.into_iter().map(|e| e.key).collect();
for id in contained {
let want: String = key("run", id)
.split('/')
.filter(|s| !s.is_empty())
.collect::<Vec<_>>()
.join("/");
assert!(
keys.contains(&want),
"list lost {id:?} (wanted {want:?}) — have {keys:?}"
);
}
}
static FORK_FREE: Mutex<()> = Mutex::new(());
#[test]
#[cfg(unix)]
fn a_second_store_on_the_same_root_is_refused_and_names_the_pid() {
let _no_forks = FORK_FREE.lock().unwrap_or_else(|e| e.into_inner());
let td = tmp();
let root = td.path().join("state");
let first = FileStore::open(&root).expect("the first open takes the lock");
let err = match FileStore::open(&root) {
Err(StoreError::Io(m)) => m,
Err(other) => panic!("the refusal must be an Io error, got {other:?}"),
Ok(_) => panic!("a second open on the same root must fail"),
};
let me = std::process::id();
assert!(
err.contains(&format!("pid {me}")),
"the message must name the holder: {err}"
);
assert!(
err.contains(&root.display().to_string()),
"…and the directory being fought over: {err}"
);
assert!(
err.contains("agent.name") || err.contains("store.file.path"),
"…and what to do about it: {err}"
);
let other_root = td.path().join("other");
let second = FileStore::open(&other_root).expect("a different root is a different lock");
assert_eq!(second.root(), other_root.as_path());
drop(first);
let deadline = Instant::now() + Duration::from_secs(2);
let reopened = loop {
match FileStore::open(&root) {
Ok(s) => break s,
Err(e) if Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(20));
let _ = e;
}
Err(e) => panic!("the lock was not released on drop: {e}"),
}
};
assert_eq!(reopened.root(), root.as_path());
assert!(
FileStore::open(&root).is_err(),
"…and the reopened store now holds it"
);
}
#[test]
#[cfg(unix)]
fn the_lock_is_held_across_processes() {
let _no_forks = FORK_FREE.lock().unwrap_or_else(|e| e.into_inner());
let td = tmp();
let root = td.path().join("state");
let holder = FileStore::open(&root).expect("open");
let out = std::process::Command::new("sh")
.arg("-c")
.arg(format!(
"flock -n {} -c true; echo $?",
shell_quote(&root.join(".lock"))
))
.output();
match out {
Ok(o) if o.status.success() => {
let code = String::from_utf8_lossy(&o.stdout).trim().to_string();
if code == "127" {
eprintln!("skipped: flock(1) not installed");
} else {
assert_ne!(
code, "0",
"another process acquired the lock while we hold it"
);
}
}
_ => eprintln!("skipped: flock(1) unavailable"),
}
drop(holder);
}
fn shell_quote(p: &Path) -> String {
format!("'{}'", p.display().to_string().replace('\'', r"'\''"))
}
#[test]
#[cfg(unix)]
fn a_stale_lock_file_from_a_dead_process_does_not_block_startup() {
let td = tmp();
let root = td.path().join("state");
fs::create_dir_all(&root).unwrap();
fs::write(root.join(".lock"), "999999").unwrap();
let s = FileStore::open(&root).expect("a stale lock FILE is not a held lock");
let k = key("run", "after-restart");
assert_eq!(
s.put(&k, 1, &env("run", "after-restart", 1, json!(true)))
.unwrap(),
PutOutcome::Ok
);
let pid = fs::read_to_string(root.join(".lock")).unwrap();
assert_eq!(
pid.trim(),
std::process::id().to_string(),
"the holder rewrites the lock file with its own pid"
);
}
#[test]
#[cfg(unix)]
fn a_contended_open_fails_fast_rather_than_blocking() {
let td = tmp();
let root = td.path().join("state");
let _held = FileStore::open(&root).expect("open");
let t0 = Instant::now();
assert!(FileStore::open(&root).is_err());
assert!(
t0.elapsed() < Duration::from_secs(2),
"the contended open took {:?} — LOCK_NB is what makes this a startup error",
t0.elapsed()
);
}
#[test]
fn a_restart_finds_its_state_under_the_same_instance_segment() {
let _no_forks = FORK_FREE.lock().unwrap_or_else(|e| e.into_inner());
let td = tmp();
let root = td.path().join("state");
let mut want: Vec<(String, u64)> = Vec::new();
{
let s = FileStore::open(&root).expect("open");
for (kind, id, seq) in [("run", "r1", 3), ("timer", "t1", 1), ("memory", "m1", 9)] {
let k = key(kind, id);
s.put(&k, seq, &env(kind, id, seq, json!({"k": kind})))
.unwrap();
want.push((k, seq));
}
s.put(
&agentd::store::key("agentd", "other", "run", "r1"),
1,
&env("run", "r1", 1, json!({"k": "other"})),
)
.unwrap();
}
let s = FileStore::open(&root).expect("reopen");
want.sort();
let got: Vec<(String, u64)> = s
.list("agentd/inst")
.unwrap()
.into_iter()
.map(|e| (e.key, e.seq.unwrap_or(0)))
.collect();
assert_eq!(got, want, "every key survived the restart, with its seq");
assert_eq!(
s.list("agentd").unwrap().len(),
want.len() + 1,
"the other instance's key is there, under its own segment"
);
}
#[test]
fn concurrent_readers_of_distinct_keys_are_consistent() {
let td = tmp();
let root = td.path().join("state");
let s = Arc::new(FileStore::open(&root).expect("open"));
for i in 0..20u64 {
let k = key("run", &format!("r{i}"));
s.put(
&k,
i + 1,
&env("run", &format!("r{i}"), i + 1, json!({"i": i})),
)
.unwrap();
}
let handles: Vec<_> = (0..4)
.map(|_| {
let s = s.clone();
std::thread::spawn(move || {
for i in 0..20u64 {
let e = Envelope::from_value(
s.get(&key("run", &format!("r{i}")), None).unwrap().unwrap(),
)
.unwrap();
assert_eq!(e.state, json!({"i": i}));
}
s.list("agentd/inst").unwrap().len()
})
})
.collect();
for h in handles {
assert_eq!(h.join().unwrap(), 20);
}
}
#[test]
fn an_unparseable_record_is_reported_with_its_path() {
let td = tmp();
let root = td.path().join("state");
let s = FileStore::open(&root).expect("open");
let k = key("run", "bad");
s.put(&k, 1, &env("run", "bad", 1, json!({"ok": true})))
.unwrap();
let path: PathBuf = root.join(expected_file(&k));
fs::write(&path, b"{not json").unwrap();
match s.get(&k, None) {
Err(StoreError::Corrupt(m)) => assert!(
m.contains(&path.display().to_string()),
"the diagnostic must name the file: {m}"
),
other => panic!("a corrupt record must surface as Corrupt, got {other:?}"),
}
}