use std::path::Path;
use std::sync::mpsc as std_mpsc;
use std::time::Duration;
pub use notify_debouncer_mini::Debouncer;
use notify_debouncer_mini::new_debouncer;
use notify_debouncer_mini::notify::RecommendedWatcher;
use notify_debouncer_mini::notify::RecursiveMode;
use tokio::sync::mpsc as tokio_mpsc;
const DEBOUNCE_TIMEOUT: Duration = Duration::from_millis(300);
const CHANNEL_CAPACITY: usize = 8;
pub fn watch_bears_dir(
dir: &Path,
) -> Result<
(Debouncer<RecommendedWatcher>, tokio_mpsc::Receiver<()>),
notify_debouncer_mini::notify::Error,
> {
let (std_tx, std_rx) = std_mpsc::channel();
let mut debouncer = new_debouncer(DEBOUNCE_TIMEOUT, std_tx)?;
debouncer.watcher().watch(dir, RecursiveMode::Recursive)?;
let (tok_tx, tok_rx) = tokio_mpsc::channel(CHANNEL_CAPACITY);
tokio::task::spawn_blocking(move || {
for _result in &std_rx {
if tok_tx.blocking_send(()).is_err() {
break;
}
}
});
Ok((debouncer, tok_rx))
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tempfile::tempdir;
#[tokio::test]
#[ignore = "timing-sensitive filesystem watcher; run with `cargo test -- --ignored`"]
async fn watcher_emits_event_on_file_create() {
let tmp = tempdir().expect("tempdir");
let bears_dir = tmp.path().join(".bears");
std::fs::create_dir_all(&bears_dir).expect("create .bears dir");
let (_debouncer, mut rx) = watch_bears_dir(&bears_dir).expect("start watcher");
tokio::time::sleep(Duration::from_millis(100)).await;
std::fs::write(bears_dir.join("test-task.md"), "# hello\n").expect("write test file");
let received = tokio::time::timeout(Duration::from_secs(2), rx.recv()).await;
assert!(
received.is_ok(),
"timeout waiting for watcher event after file create"
);
assert_eq!(
received.unwrap(),
Some(()),
"expected Some(()) from watcher channel"
);
}
#[tokio::test]
#[ignore = "timing-sensitive filesystem watcher; run with `cargo test -- --ignored`"]
async fn watcher_emits_event_on_file_modify_and_delete() {
let tmp = tempdir().expect("tempdir");
let bears_dir = tmp.path().join(".bears");
std::fs::create_dir_all(&bears_dir).expect("create .bears dir");
let task_file = bears_dir.join("abc-existing-task.md");
std::fs::write(&task_file, "# initial\n").expect("write initial file");
let (_debouncer, mut rx) = watch_bears_dir(&bears_dir).expect("start watcher");
tokio::time::sleep(Duration::from_millis(100)).await;
std::fs::write(&task_file, "# modified\n").expect("modify file");
let sig = tokio::time::timeout(Duration::from_secs(2), rx.recv()).await;
assert!(
sig.is_ok(),
"timeout waiting for watcher event after file modify"
);
assert_eq!(sig.unwrap(), Some(()));
std::fs::remove_file(&task_file).expect("delete file");
let sig = tokio::time::timeout(Duration::from_secs(2), rx.recv()).await;
assert!(
sig.is_ok(),
"timeout waiting for watcher event after file delete"
);
assert_eq!(sig.unwrap(), Some(()));
}
#[tokio::test]
#[ignore = "timing-sensitive debounce test; run with `cargo test -- --ignored`"]
async fn watcher_debounce_coalesces_rapid_writes() {
const WRITE_COUNT: usize = 10;
let tmp = tempdir().expect("tempdir");
let bears_dir = tmp.path().join(".bears");
std::fs::create_dir_all(&bears_dir).expect("create .bears dir");
let (_debouncer, mut rx) = watch_bears_dir(&bears_dir).expect("start watcher");
tokio::time::sleep(Duration::from_millis(100)).await;
for i in 0..WRITE_COUNT {
std::fs::write(
bears_dir.join(format!("task-{i:02}.md")),
format!("# task {i}\n"),
)
.expect("write task file");
}
let drain_deadline = Duration::from_secs(1);
let mut signal_count: usize = 0;
let start = std::time::Instant::now();
while start.elapsed() < drain_deadline {
let remaining = drain_deadline.saturating_sub(start.elapsed());
match tokio::time::timeout(remaining, rx.recv()).await {
Ok(Some(())) => signal_count += 1,
Ok(None) | Err(_) => break,
}
}
assert!(
signal_count > 0,
"expected at least one signal after rapid writes"
);
assert!(
signal_count < WRITE_COUNT,
"expected debouncing to coalesce signals: got {signal_count} for {WRITE_COUNT} writes"
);
}
}