#![cfg(all(feature = "watch", unix))]
use std::collections::VecDeque;
use std::fs;
use std::io::{BufRead, BufReader, Read};
use std::path::Path;
use std::process::{Child, Command, Stdio};
use std::sync::mpsc::{self, Receiver, RecvTimeoutError};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
use fdu_core::control::DEFAULT_CONTROL_LINE_LIMIT;
const DEADLINE: Duration = Duration::from_secs(30);
fn oversized_rule() -> Vec<u8> {
let mut source = vec![b'a'; DEFAULT_CONTROL_LINE_LIMIT + 1];
source.push(b'\n');
source
}
struct Watching {
child: Child,
lines: Receiver<String>,
stderr: Option<JoinHandle<String>>,
}
impl Watching {
fn spawn(tree: &Path, cache: &Path) -> Self {
Self::spawn_view(tree, cache, "files")
}
fn spawn_view(tree: &Path, cache: &Path, view: &str) -> Self {
Self::spawn_selecting(tree, cache, view, &[])
}
fn spawn_selecting(tree: &Path, cache: &Path, view: &str, selection: &[&str]) -> Self {
let mut child = Command::new(env!("CARGO_BIN_EXE_fdu"))
.args(["--watch", "--view", view, "--format", "jsonl", "--interval", "1s"])
.args(selection)
.arg(tree)
.env("XDG_CACHE_HOME", cache)
.env_remove("FDU_CACHE_DIR")
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn watching fdu");
let stdout = child.stdout.take().expect("piped stdout");
let (sender, lines) = mpsc::channel();
std::thread::spawn(move || {
for line in BufReader::new(stdout).lines() {
let Ok(line) = line else { break };
if sender.send(line).is_err() {
break;
}
}
});
let mut stderr_pipe = child.stderr.take().expect("piped stderr");
let stderr = std::thread::spawn(move || {
let mut text = String::new();
let _ = stderr_pipe.read_to_string(&mut text);
text
});
Self { child, lines, stderr: Some(stderr) }
}
fn wait_for(&mut self, description: &str, wanted: impl Fn(&str) -> bool) -> String {
let started = Instant::now();
let mut recent = VecDeque::new();
loop {
match self.lines.recv_timeout(DEADLINE.saturating_sub(started.elapsed())) {
Ok(line) if wanted(&line) => return line,
Ok(line) => {
if recent.len() == 12 {
recent.pop_front();
}
recent.push_back(line);
}
Err(RecvTimeoutError::Timeout) => {
self.fail(&format!("timed out after {DEADLINE:?} waiting for {description}; recent lines: {recent:?}"));
}
Err(RecvTimeoutError::Disconnected) => {
self.fail(&format!("the watch exited before {description}"));
}
}
}
}
fn initial_report(&mut self) -> String {
let envelope =
self.wait_for("the initial report envelope", |line| line.contains("\"fdu.report/"));
self.wait_for("the initial files section", |line| line.contains("\"view\": \"files\""));
envelope
}
fn wait_for_change(&mut self, path: &str) -> String {
let quoted = format!("\"path\": \"{path}\"");
self.wait_for(&format!("a change record for {path}"), |line| {
line.contains("\"record\": \"change\"") && line.contains("ed)
})
}
fn wait_for_op(&mut self, op: &str, path: &str) -> String {
let quoted = format!("\"path\": \"{path}\"");
let operation = format!("\"op\": \"{op}\"");
self.wait_for(&format!("a {op} record for {path}"), |line| {
line.contains("\"record\": \"change\"")
&& line.contains(&operation)
&& line.contains("ed)
})
}
fn fail(&mut self, reason: &str) -> ! {
let _ = self.child.kill();
let status = self.child.wait().ok();
let stderr = self.stderr.take().and_then(|reader| reader.join().ok()).unwrap_or_default();
panic!("{reason}; exit status {status:?}; stderr: {stderr}");
}
}
impl Drop for Watching {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn report(tree: &Path, cache: &Path, args: &[&str]) -> String {
let output = Command::new(env!("CARGO_BIN_EXE_fdu"))
.args(args)
.arg(tree)
.env("XDG_CACHE_HOME", cache)
.env_remove("FDU_CACHE_DIR")
.output()
.expect("run fdu");
assert!(
output.status.success(),
"fdu {args:?} failed: {}",
String::from_utf8_lossy(&output.stderr),
);
String::from_utf8_lossy(&output.stdout).into_owned()
}
fn tree_with(files: &[(&str, &[u8])]) -> (tempfile::TempDir, std::path::PathBuf) {
let root = tempfile::tempdir().expect("tempdir");
let tree = root.path().join("tree");
fs::create_dir(&tree).expect("create tree");
for (name, bytes) in files {
fs::write(tree.join(name), bytes).expect("write fixture file");
}
(root, tree)
}
#[test]
fn a_watch_serves_a_tree_whose_ignore_rule_exceeds_the_control_bound() {
let rule = oversized_rule();
let (_root, tree) = tree_with(&[(".gitignore", &rule), ("kept.txt", b"kept")]);
let cache = tempfile::tempdir().expect("cache tempdir");
let mut watch = Watching::spawn(&tree, cache.path());
let envelope = watch.initial_report();
assert!(
envelope.contains("\"complete\": true"),
"a watch over an oversized ignore rule answered only partially: {envelope}",
);
assert!(
envelope.contains("\"applied\": 0, \"rules\": 0, \"refused\": 1"),
"the watch read the control file and named its refusal: {envelope}",
);
}
#[test]
fn a_control_file_edit_repaints_the_ignored_share() {
let (_root, tree) =
tree_with(&[(".gitignore", b"*.log\n"), ("kept.txt", b"kept"), ("debug.log", b"debug")]);
let cache = tempfile::tempdir().expect("cache tempdir");
let mut watch = Watching::spawn_view(&tree, cache.path(), "summary");
watch.wait_for("the initial summary's ignored share", |line| {
line.contains("\"view\": \"summary\"")
&& line.contains("\"ignored\": {\"files\": 1, \"dirs\": 0, \"bytes\": 5,")
});
fs::write(tree.join(".gitignore"), b"# nothing ignored\n").expect("rewrite the control file");
watch.wait_for("a repaint with nothing ignored", |line| {
line.contains("\"view\": \"summary\"")
&& line.contains("\"ignored\": {\"files\": 0, \"dirs\": 0, \"bytes\": 0,")
});
}
#[test]
fn a_control_file_edited_past_the_bound_does_not_end_a_watch() {
let (_root, tree) = tree_with(&[(".gitignore", b"*.log\n"), ("kept.txt", b"kept")]);
let cache = tempfile::tempdir().expect("cache tempdir");
let mut watch = Watching::spawn(&tree, cache.path());
watch.initial_report();
fs::write(tree.join(".gitignore"), oversized_rule()).expect("rewrite the control file");
watch.wait_for_change(".gitignore");
fs::write(tree.join("after.txt"), b"after").expect("write a later file");
watch.wait_for_change("after.txt");
let started = Instant::now();
while !report(&tree, cache.path(), &["--view", "files", "--format", "jsonl", "--stale-ok"])
.contains("\"path\": \"after.txt\"")
{
assert!(started.elapsed() < DEADLINE, "the watch never persisted the later change");
std::thread::sleep(Duration::from_millis(100));
}
drop(watch);
let mut next = Watching::spawn(&tree, cache.path());
let envelope = next.initial_report();
assert!(
envelope.contains("\"source\": \"warm_revalidate\"")
&& envelope.contains("\"complete\": true"),
"the snapshot a watch saved after a control edit did not start the next watch: {envelope}",
);
}
#[test]
fn a_watch_and_a_one_shot_report_start_warm_from_each_others_snapshot() {
let (_root, tree) = tree_with(&[(".gitignore", b"*.log\n"), ("kept.txt", b"kept")]);
let cache = tempfile::tempdir().expect("cache tempdir");
{
let mut watch = Watching::spawn(&tree, cache.path());
let envelope = watch.initial_report();
assert!(
envelope.contains("\"source\": \"cold_scan\""),
"expected a cold start: {envelope}"
);
}
let analyzed =
report(&tree, cache.path(), &["--analyze", "lines", "--view", "files", "--format", "json"]);
assert!(
analyzed.contains("\"source\": \"warm_revalidate\""),
"a report after a watch rescanned instead of reusing the watch's snapshot: {analyzed}",
);
report(&tree, cache.path(), &["--view", "files", "--format", "json", "--cache", "on"]);
let mut watch = Watching::spawn(&tree, cache.path());
let envelope = watch.initial_report();
assert!(
envelope.contains("\"source\": \"warm_revalidate\""),
"a watch after a one-shot report rescanned instead of reusing its snapshot: {envelope}",
);
}
#[test]
fn a_rule_edit_moves_a_streamed_entry_out_of_and_back_into_excluded_ignored() {
let (_root, tree) = tree_with(&[
(".gitignore", b"# nothing ignored\n"),
("kept.txt", b"kept"),
("debug.log", b"debug"),
]);
let cache = tempfile::tempdir().expect("cache tempdir");
let mut watch = Watching::spawn_selecting(&tree, cache.path(), "files", &["--ignored=exclude"]);
let row = watch
.wait_for("the initial row for debug.log", |line| line.contains("\"path\": \"debug.log\""));
assert!(
row.contains("\"ignored\": false"),
"the initial row must state the classification the stream maintains: {row}",
);
fs::write(tree.join(".gitignore"), b"*.log\n").expect("rewrite the control file");
let left = watch.wait_for_op("remove", "debug.log");
assert!(
left.contains("\"ignored\": true"),
"the record that drops the row must say why it left: {left}",
);
fs::write(tree.join(".gitignore"), b"# nothing ignored\n").expect("rewrite the control file");
let returned = watch.wait_for_op("upsert", "debug.log");
assert!(
returned.contains("\"ignored\": false")
&& returned.contains("\"kind\": \"file\"")
&& returned.contains("\"bytes\": 5"),
"the record that restores the row must carry the facts to draw it: {returned}",
);
}
#[test]
fn a_rule_edit_moves_a_streamed_entry_into_and_out_of_only_ignored() {
let (_root, tree) = tree_with(&[
(".gitignore", b"# nothing ignored\n"),
("kept.txt", b"kept"),
("debug.log", b"debug"),
]);
let cache = tempfile::tempdir().expect("cache tempdir");
let mut watch = Watching::spawn_selecting(&tree, cache.path(), "files", &["--ignored=only"]);
let section =
watch.wait_for("the initial files section", |line| line.contains("\"view\": \"files\""));
assert!(
!section.contains("debug.log"),
"no rule ignores anything yet, so the listing is empty: {section}",
);
fs::write(tree.join(".gitignore"), b"*.log\n").expect("rewrite the control file");
let arrived = watch.wait_for_op("upsert", "debug.log");
assert!(
arrived.contains("\"ignored\": true") && arrived.contains("\"bytes\": 5"),
"an entry entering the selection arrives with its facts: {arrived}",
);
fs::write(tree.join(".gitignore"), b"# nothing ignored\n").expect("rewrite the control file");
let left = watch.wait_for_op("remove", "debug.log");
assert!(
left.contains("\"ignored\": false"),
"the record that drops the row must say why it left: {left}",
);
}