mod common;
use std::io::{BufRead, BufReader};
use std::process::{Child, Command, Stdio};
use std::sync::mpsc;
use std::time::Duration;
use common::write_file;
const EVENT_BUDGET: Duration = Duration::from_secs(30);
struct Watcher {
child: Child,
rx: mpsc::Receiver<String>,
}
impl Watcher {
fn spawn(store: &std::path::Path, extra: &[&str]) -> Self {
let mut child = Command::new(assert_cmd::cargo::cargo_bin("dbmd"))
.args(["--json", "watch", "--interval", "1", "--dir"])
.arg(store)
.args(extra)
.stdout(Stdio::piped())
.stderr(Stdio::null())
.spawn()
.expect("spawn dbmd watch");
let stdout = child.stdout.take().expect("piped stdout");
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
for line in BufReader::new(stdout).lines() {
match line {
Ok(l) => {
if tx.send(l).is_err() {
break;
}
}
Err(_) => break,
}
}
});
Self { child, rx }
}
fn next(&self) -> serde_json::Value {
let line = self
.rx
.recv_timeout(EVENT_BUDGET)
.expect("watch event within budget");
serde_json::from_str(&line).expect("valid NDJSON event line")
}
fn wait_for(&self, event: &str) -> serde_json::Value {
loop {
let v = self.next();
if v["event"] == serde_json::json!(event) {
return v;
}
}
}
}
impl Drop for Watcher {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn scratch_store() -> (tempfile::TempDir, std::path::PathBuf) {
let dir = tempfile::TempDir::new().unwrap();
let store = dir.path().to_path_buf();
write_file(
&store,
"DB.md",
"---\ntype: db-md\nscope: test\nowner: T\n---\n# T\n",
);
write_file(
&store,
"records/widgets/a.md",
"---\ntype: widget\nsummary: A\n---\n",
);
(dir, store)
}
#[test]
fn watch_streams_create_modify_remove() {
let (_tmp, store) = scratch_store();
let watcher = Watcher::spawn(&store, &[]);
let baseline = watcher.next();
assert_eq!(baseline["event"], serde_json::json!("baseline"));
assert_eq!(baseline["files"], serde_json::json!(2));
write_file(
&store,
"records/widgets/b.md",
"---\ntype: widget\nsummary: B\n---\n",
);
let created = watcher.wait_for("created");
assert_eq!(created["path"], serde_json::json!("records/widgets/b.md"));
assert!(created["at"].is_string());
write_file(
&store,
"records/widgets/b.md",
"---\ntype: widget\nsummary: B\n---\nnow with a body\n",
);
let modified = watcher.wait_for("modified");
assert_eq!(modified["path"], serde_json::json!("records/widgets/b.md"));
std::fs::remove_file(store.join("records/widgets/b.md")).unwrap();
let removed = watcher.wait_for("removed");
assert_eq!(removed["path"], serde_json::json!("records/widgets/b.md"));
}
#[test]
fn watch_path_scopes_events() {
let (_tmp, store) = scratch_store();
write_file(
&store,
"records/notes/n.md",
"---\ntype: note\nsummary: N\n---\n",
);
let watcher = Watcher::spawn(&store, &["--path", "records/widgets"]);
let baseline = watcher.next();
assert_eq!(baseline["files"], serde_json::json!(1)); assert_eq!(baseline["path"], serde_json::json!("records/widgets"));
write_file(
&store,
"records/notes/n2.md",
"---\ntype: note\nsummary: N2\n---\n",
);
write_file(
&store,
"records/widgets/w2.md",
"---\ntype: widget\nsummary: W2\n---\n",
);
let created = watcher.wait_for("created");
assert_eq!(created["path"], serde_json::json!("records/widgets/w2.md"));
}
#[test]
fn watch_zero_interval_refused() {
let (_tmp, store) = scratch_store();
common::dbmd()
.args(["watch", "--interval", "0", "--dir"])
.arg(&store)
.assert()
.failure()
.code(1); }