use std::path::{Path, PathBuf};
use std::sync::mpsc as std_mpsc;
use notify::{Event, EventKind, RecommendedWatcher, RecursiveMode, Watcher};
use tokio::sync::mpsc;
#[non_exhaustive]
#[derive(Debug, Clone)]
pub enum ChangeEvent {
Touched(Vec<PathBuf>),
Removed(Vec<PathBuf>),
}
#[non_exhaustive]
#[derive(Debug, thiserror::Error)]
pub enum WatcherError {
#[error("watcher setup failed: {0}")]
Setup(String),
#[error("watch failed: {0}")]
Watch(String),
}
pub struct FileWatcher {
_inner: RecommendedWatcher,
rx: mpsc::UnboundedReceiver<ChangeEvent>,
}
impl FileWatcher {
pub fn new(root: &Path) -> Result<Self, WatcherError> {
let (event_tx, event_rx) = std_mpsc::channel::<notify::Result<Event>>();
let mut watcher = notify::recommended_watcher(event_tx)
.map_err(|e| WatcherError::Setup(e.to_string()))?;
watcher
.watch(root, RecursiveMode::Recursive)
.map_err(|e| WatcherError::Watch(e.to_string()))?;
let (out_tx, out_rx) = mpsc::unbounded_channel();
std::thread::spawn(move || {
for result in event_rx {
let Ok(event) = result else { continue };
let mapped = map_event(event);
if let Some(change) = mapped {
if out_tx.send(change).is_err() {
break;
}
}
}
});
Ok(Self {
_inner: watcher,
rx: out_rx,
})
}
pub async fn next(&mut self) -> Option<ChangeEvent> {
self.rx.recv().await
}
}
fn map_event(event: Event) -> Option<ChangeEvent> {
match event.kind {
EventKind::Create(_) | EventKind::Modify(_) => Some(ChangeEvent::Touched(event.paths)),
EventKind::Remove(_) => Some(ChangeEvent::Removed(event.paths)),
_ => None,
}
}