use super::*;
use std::fs;
use std::time::{Duration, Instant};
use tempfile::{TempDir, tempdir};
const WAIT_TIMEOUT: Duration = Duration::from_secs(5);
const POLL_INTERVAL: Duration = Duration::from_millis(50);
fn poll_until<T>(
mut probe: impl FnMut() -> Option<T>,
timeout: Duration,
interval: Duration,
) -> Option<T> {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if let Some(found) = probe() {
return Some(found);
}
std::thread::sleep(interval);
}
probe()
}
fn workspace() -> (TempDir, PathBuf) {
let dir = tempdir().unwrap();
let root = dir.path().canonicalize().unwrap();
(dir, root)
}
fn wait_quiet(watcher: &Watcher) {
std::thread::sleep(Duration::from_millis(200));
let _ = watcher.tick();
}
fn wait_for(watcher: &Watcher, pred: impl Fn(&Change) -> bool) -> Vec<Change> {
let probe = || {
let changes = watcher.tick();
let found = changes.iter().any(&pred);
found.then_some(changes)
};
poll_until(probe, WAIT_TIMEOUT, POLL_INTERVAL).unwrap_or_default()
}
fn expect_touch(root: &Path, target: &Path) {
fs::create_dir_all(target.parent().unwrap()).unwrap();
let watcher = Watcher::new(root).unwrap();
wait_quiet(&watcher);
fs::write(target, b"x").unwrap();
let changes = wait_for(&watcher, |c| c.path == *target);
let hit = changes.iter().find(|c| c.path == *target).expect("event");
assert_eq!(hit.kind, ChangeKind::Touched);
}
#[test]
fn poll_until_returns_the_first_some() {
let mut calls = 0;
let got = poll_until(
|| {
calls += 1;
(calls == 2).then_some(calls)
},
Duration::from_secs(5),
Duration::from_millis(1),
);
assert_eq!(got, Some(2));
}
#[test]
fn poll_until_times_out_to_none() {
let got = poll_until(
|| None::<()>,
Duration::from_millis(30),
Duration::from_millis(5),
);
assert_eq!(got, None);
}
#[test]
fn new_errors_on_missing_repo() {
let root = tempdir().unwrap();
let err = Watcher::new(&root.path().join("nope")).err().unwrap();
assert!(err.to_string().contains("filesystem watcher"));
}
#[test]
fn tick_is_empty_when_nothing_changed() {
let (_dir, root) = workspace();
let watcher = Watcher::new(&root).unwrap();
wait_quiet(&watcher);
assert!(watcher.tick().is_empty());
}
#[test]
fn detects_step_request_creation_at_conv_repo_root() {
let (_dir, root) = workspace();
expect_touch(&root, &root.join("steps/abc-1/001/request.json"));
}
#[test]
fn detects_subagent_step_record_at_conv_repo_root() {
let (_dir, root) = workspace();
expect_touch(&root, &root.join("steps/aa-bb/001/request.json"));
}
#[test]
fn detects_inbox_deposits_at_workspace_root() {
let (_dir, root) = workspace();
expect_touch(&root, &root.join("inbox/aa-bb/user-001.md"));
}
#[test]
fn detects_goal_md_update_in_an_agent_worktree() {
let (_dir, root) = workspace();
expect_touch(&root, &root.join("agents/aa-bb/goal.md"));
}
#[test]
fn detects_removal_under_summary() {
let (_dir, root) = workspace();
let target = root.join("agents/aa-bb/summary/001.md");
fs::create_dir_all(target.parent().unwrap()).unwrap();
fs::write(&target, b"hi").unwrap();
let watcher = Watcher::new(&root).unwrap();
wait_quiet(&watcher);
fs::remove_file(&target).unwrap();
let changes = wait_for(&watcher, |c| c.path == target);
let hit = changes.iter().find(|c| c.path == target).expect("event");
assert_eq!(hit.kind, ChangeKind::Removed);
}
#[test]
fn ignores_paths_outside_allowlist() {
let (_dir, root) = workspace();
let watcher = Watcher::new(&root).unwrap();
wait_quiet(&watcher);
fs::write(root.join("README.md"), b"x").unwrap();
fs::create_dir_all(root.join("random")).unwrap();
fs::write(root.join("random/x.txt"), b"x").unwrap();
fs::create_dir_all(root.join("agents/aa-bb/random")).unwrap();
fs::write(root.join("agents/aa-bb/random/x.txt"), b"x").unwrap();
std::thread::sleep(Duration::from_millis(200));
assert!(watcher.tick().is_empty());
}
#[test]
fn coalesces_rapid_writes_to_one_event() {
let (_dir, root) = workspace();
fs::create_dir_all(root.join("steps")).unwrap();
let target = root.join("steps/out.log");
let watcher = Watcher::new(&root).unwrap();
wait_quiet(&watcher);
for i in 0..5 {
fs::write(&target, format!("line {i}")).unwrap();
}
let changes = wait_for(&watcher, |c| c.path == target);
let hits: Vec<_> = changes.iter().filter(|e| e.path == target).collect();
assert_eq!(hits.len(), 1, "got {changes:?}");
assert_eq!(hits[0].kind, ChangeKind::Touched);
}
#[test]
fn coalesces_atomic_rename_to_destination() {
let (_dir, root) = workspace();
fs::create_dir_all(root.join("steps/abc/001")).unwrap();
let tmp = root.join("steps/abc/001/request.json.tmp");
let final_path = root.join("steps/abc/001/request.json");
let watcher = Watcher::new(&root).unwrap();
wait_quiet(&watcher);
fs::write(&tmp, b"{}").unwrap();
fs::rename(&tmp, &final_path).unwrap();
let changes = wait_for(&watcher, |c| c.path == final_path);
let finals: Vec<_> = changes.iter().filter(|e| e.path == final_path).collect();
let tmps: Vec<_> = changes.iter().filter(|e| e.path == tmp).collect();
assert_eq!(finals.len(), 1, "one event for destination: {changes:?}");
assert_eq!(finals[0].kind, ChangeKind::Touched);
assert!(tmps.is_empty(), "rename source leaked: {tmps:?}");
}
#[test]
fn coalesce_drops_rename_source_with_trailing_modify() {
use notify::event::{CreateKind, DataChange};
let (_dir, root) = workspace();
fs::create_dir_all(root.join("steps/abc/001")).unwrap();
let tmp = root.join("steps/abc/001/request.json.tmp");
let final_path = root.join("steps/abc/001/request.json");
fs::write(&final_path, b"{}").unwrap(); let name = EventKind::Modify(ModifyKind::Name(RenameMode::Any));
let modify = EventKind::Modify(ModifyKind::Data(DataChange::Content));
let raw = vec![
(tmp.clone(), EventKind::Create(CreateKind::File)),
(tmp.clone(), name),
(tmp.clone(), modify),
(final_path.clone(), EventKind::Create(CreateKind::File)),
];
let changes = coalesce(&root, RootKind::Workspace, raw);
assert_eq!(changes.len(), 1);
assert_eq!(changes[0].path, final_path);
assert_eq!(changes[0].kind, ChangeKind::Touched);
}
#[test]
fn coalesce_surfaces_genuine_deletion_with_non_remove_last_event() {
use notify::event::{CreateKind, DataChange, RemoveKind};
let (_dir, root) = workspace();
fs::create_dir_all(root.join("agents/aa-bb/summary")).unwrap();
let gone = root.join("agents/aa-bb/summary/001.md");
let modify = EventKind::Modify(ModifyKind::Data(DataChange::Content));
let raw = vec![
(gone.clone(), EventKind::Create(CreateKind::File)),
(gone.clone(), EventKind::Remove(RemoveKind::File)),
(gone.clone(), modify),
];
let changes = coalesce(&root, RootKind::Workspace, raw);
assert_eq!(changes.len(), 1);
assert_eq!(changes[0].path, gone);
assert_eq!(changes[0].kind, ChangeKind::Removed);
}
#[test]
fn ingest_splits_name_both_into_from_and_to() {
let mut raw = Vec::new();
ingest(
Event {
kind: EventKind::Modify(ModifyKind::Name(RenameMode::Both)),
paths: vec![PathBuf::from("/a"), PathBuf::from("/b")],
attrs: notify::event::EventAttributes::default(),
},
&mut raw,
);
assert_eq!(raw.len(), 2);
assert!(matches!(
raw[0].1,
EventKind::Modify(ModifyKind::Name(RenameMode::From))
));
assert!(matches!(
raw[1].1,
EventKind::Modify(ModifyKind::Name(RenameMode::To))
));
}
#[test]
fn coalesce_drops_prior_events_when_rename_from_arrives() {
let repo = Path::new("/r");
let p = PathBuf::from("/r/steps/abc/001/request.json");
let raw = vec![
(
p.clone(),
EventKind::Create(notify::event::CreateKind::File),
),
(
p.clone(),
EventKind::Modify(ModifyKind::Name(RenameMode::From)),
),
];
assert!(coalesce(repo, RootKind::Workspace, raw).is_empty());
}
#[test]
fn classify_over_existence_and_rename() {
let root = tempdir().unwrap();
let present = root.path().join("x");
fs::write(&present, b"").unwrap();
assert_eq!(classify(false, &present), Some(ChangeKind::Touched));
assert_eq!(classify(true, &present), Some(ChangeKind::Touched));
let absent = Path::new("/no/such/path");
assert_eq!(classify(true, absent), None);
assert_eq!(classify(false, absent), Some(ChangeKind::Removed));
}