use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, MutexGuard, Weak};
use std::time::Duration;
use notify_debouncer_full::notify::event::{ModifyKind, RenameMode};
use notify_debouncer_full::notify::{EventKind, RecommendedWatcher, RecursiveMode};
use notify_debouncer_full::{
DebounceEventResult, DebouncedEvent, Debouncer, RecommendedCache, new_debouncer,
};
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use crate::scanner::{IgnoreFilter, WorkspaceScanner};
use crate::service::EngineService;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum WatcherAction {
Upsert(PathBuf),
Delete(PathBuf),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct WatchedRoot {
pub project_id: u64,
pub root: PathBuf,
pub directories: usize,
}
pub struct WorkspaceWatcher {
inner: Arc<Inner>,
worker: JoinHandle<()>,
}
struct Inner {
engine: EngineService,
state: Mutex<State>,
}
struct State {
debouncer: Debouncer<RecommendedWatcher, RecommendedCache>,
projects: HashMap<u64, ProjectWatch>,
dir_refs: HashMap<PathBuf, usize>,
}
struct ProjectWatch {
root: PathBuf,
path_dependencies: Vec<PathBuf>,
scanner: WorkspaceScanner,
ignore: IgnoreFilter,
dirs: HashSet<PathBuf>,
}
struct FileOp {
project_id: u64,
root: PathBuf,
rel_path: PathBuf,
delete: bool,
}
impl WorkspaceWatcher {
pub fn start(engine: EngineService, debounce: Duration) -> anyhow::Result<Self> {
let (tx, mut rx) = mpsc::unbounded_channel::<DebounceEventResult>();
let debouncer = new_debouncer(debounce, None, move |res| {
let _ = tx.send(res);
})?;
let inner = Arc::new(Inner {
engine,
state: Mutex::new(State {
debouncer,
projects: HashMap::new(),
dir_refs: HashMap::new(),
}),
});
let weak: Weak<Inner> = Arc::downgrade(&inner);
let worker = tokio::spawn(async move {
while let Some(result) = rx.recv().await {
let Some(inner) = weak.upgrade() else { break };
match result {
Ok(events) => inner.handle_events(events).await,
Err(errors) => {
for err in errors {
tracing::debug!("Watcher error: {err}");
}
}
}
}
});
Ok(Self { inner, worker })
}
pub fn add_project(&self, project_id: u64, root: &Path) -> anyhow::Result<usize> {
let root = dunce::canonicalize(root)?;
{
let state = self.inner.lock();
if let Some(existing) = state.projects.get(&project_id)
&& existing.root == root
{
return Ok(existing.dirs.len());
}
}
self.remove_project(project_id);
let scanner = WorkspaceScanner::new(&root);
let dirs = scanner.scan_dirs();
let mut state = self.inner.lock();
state.projects.insert(
project_id,
ProjectWatch {
ignore: IgnoreFilter::new(&root),
root: root.clone(),
path_dependencies: Vec::new(),
scanner,
dirs: HashSet::new(),
},
);
let count = state.sync_dirs(project_id, dirs);
tracing::info!(
"Watching project {project_id} at {} ({count} directories)",
root.display()
);
Ok(count)
}
pub fn add_path_dependency(
&self,
project_id: u64,
path_dep_root: &Path,
) -> anyhow::Result<usize> {
let root = dunce::canonicalize(path_dep_root)?;
let scanner = WorkspaceScanner::new(&root);
let dirs = scanner.scan_dirs();
let mut state = self.inner.lock();
if let Some(project) = state.projects.get_mut(&project_id)
&& !project.path_dependencies.contains(&root)
{
project.path_dependencies.push(root);
}
let count = state.add_dirs(project_id, dirs);
Ok(count)
}
pub fn remove_project(&self, project_id: u64) {
let mut state = self.inner.lock();
if state.projects.contains_key(&project_id) {
state.sync_dirs(project_id, Vec::new());
state.projects.remove(&project_id);
tracing::info!("Stopped watching project {project_id}");
}
}
pub fn projects(&self) -> Vec<WatchedRoot> {
let state = self.inner.lock();
let mut roots: Vec<_> = state
.projects
.iter()
.map(|(id, p)| WatchedRoot {
project_id: *id,
root: p.root.clone(),
directories: p.dirs.len(),
})
.collect();
roots.sort_by_key(|r| r.project_id);
roots
}
pub async fn stop(self) {
drop(self.inner);
let _ = self.worker.await;
}
}
impl Inner {
fn lock(&self) -> MutexGuard<'_, State> {
self.state.lock().unwrap_or_else(|e| e.into_inner())
}
async fn handle_events(&self, events: Vec<DebouncedEvent>) {
let (ops, resync) = self.lock().route(&events);
for op in ops {
if resync.contains(&op.project_id) {
continue; }
let result = if op.delete {
self.engine.remove_file(op.project_id, &op.rel_path).await
} else {
self.engine
.index_file(op.project_id, &op.root, &op.rel_path)
.await
};
if let Err(e) = result {
tracing::warn!("Watcher update of {} failed: {e}", op.rel_path.display());
}
}
for project_id in resync {
self.resync(project_id).await;
}
}
async fn resync(&self, project_id: u64) {
let Some((root, scanner)) = self
.lock()
.projects
.get(&project_id)
.map(|p| (p.root.clone(), p.scanner.clone()))
else {
return;
};
let dirs = scanner.scan_dirs();
self.lock().sync_dirs(project_id, dirs);
if let Err(e) = self.engine.index_project(project_id, &root).await {
tracing::warn!("Re-index of {} failed: {e}", root.display());
}
}
}
impl State {
fn sync_dirs(&mut self, project_id: u64, dirs: Vec<PathBuf>) -> usize {
let Some(project) = self.projects.get_mut(&project_id) else {
return 0;
};
let wanted: HashSet<PathBuf> = dirs.into_iter().collect();
let removed: Vec<PathBuf> = project.dirs.difference(&wanted).cloned().collect();
let added: Vec<PathBuf> = wanted.difference(&project.dirs).cloned().collect();
for dir in removed {
project.dirs.remove(&dir);
if let Some(refs) = self.dir_refs.get_mut(&dir) {
*refs -= 1;
if *refs == 0 {
self.dir_refs.remove(&dir);
let _ = self.debouncer.unwatch(&dir);
}
}
}
for dir in added {
let refs = self.dir_refs.entry(dir.clone()).or_insert(0);
if *refs == 0
&& let Err(e) = self.debouncer.watch(&dir, RecursiveMode::NonRecursive)
{
tracing::warn!("Could not watch {}: {e}", dir.display());
self.dir_refs.remove(&dir);
continue;
}
*refs += 1;
project.dirs.insert(dir);
}
project.dirs.len()
}
fn add_dirs(&mut self, project_id: u64, dirs: Vec<PathBuf>) -> usize {
let Some(project) = self.projects.get_mut(&project_id) else {
return 0;
};
for dir in dirs {
let refs = self.dir_refs.entry(dir.clone()).or_insert(0);
if *refs == 0
&& let Err(e) = self.debouncer.watch(&dir, RecursiveMode::NonRecursive)
{
tracing::warn!("Could not watch {}: {e}", dir.display());
self.dir_refs.remove(&dir);
continue;
}
*refs += 1;
project.dirs.insert(dir);
}
project.dirs.len()
}
fn route(&mut self, events: &[DebouncedEvent]) -> (Vec<FileOp>, HashSet<u64>) {
let mut ops = Vec::new();
let mut resync = HashSet::new();
for debounced in events {
let event = &debounced.event;
if event.need_rescan() {
resync.extend(self.projects.keys().copied());
continue;
}
let actions: Vec<WatcherAction> = match (&event.kind, event.paths.as_slice()) {
(EventKind::Modify(ModifyKind::Name(RenameMode::Both)), [from, to, ..]) => vec![
WatcherAction::Delete(from.clone()),
WatcherAction::Upsert(to.clone()),
],
(
EventKind::Remove(_) | EventKind::Modify(ModifyKind::Name(RenameMode::From)),
paths,
) => paths.iter().cloned().map(WatcherAction::Delete).collect(),
(EventKind::Create(_) | EventKind::Modify(_), paths) => {
paths.iter().cloned().map(WatcherAction::Upsert).collect()
}
_ => Vec::new(),
};
for action in actions {
let (path, delete) = match action {
WatcherAction::Upsert(path) => (path, false),
WatcherAction::Delete(path) => (path, true),
};
for (id, project) in self.projects.iter_mut() {
let mut matched = path
.strip_prefix(&project.root)
.ok()
.map(|p| (project.root.clone(), p.to_path_buf()));
if matched.is_none() {
for p_dep in &project.path_dependencies {
if let Ok(rel) = path.strip_prefix(p_dep) {
matched = Some((p_dep.clone(), rel.to_path_buf()));
break;
}
}
}
let Some((base_root, rel_path)) = matched else {
continue;
};
if project.ignore.invalidate(&path) {
resync.insert(*id);
continue;
}
let structural = project.dirs.contains(&path)
|| (path.is_dir() && !project.ignore.is_ignored(&path));
if structural {
resync.insert(*id);
continue;
}
if !project.scanner.is_supported(&path) || project.ignore.is_ignored(&path) {
continue;
}
ops.push(FileOp {
project_id: *id,
root: base_root,
rel_path,
delete,
});
}
}
}
(ops, resync)
}
}