1use notify_debouncer_full::notify::{RecommendedWatcher, RecursiveMode};
4use notify_debouncer_full::{DebounceEventResult, Debouncer, RecommendedCache, new_debouncer};
5use ragfs_core::FileEvent;
6use std::path::Path;
7use std::sync::Arc;
8use std::sync::atomic::{AtomicUsize, Ordering};
9use std::sync::mpsc;
10use std::time::Duration;
11use tokio::sync::mpsc as tokio_mpsc;
12use tracing::{debug, error, warn};
13
14pub struct FileWatcher {
16 debouncer: Debouncer<RecommendedWatcher, RecommendedCache>,
17}
18
19impl FileWatcher {
20 pub fn new(
22 event_tx: tokio_mpsc::Sender<FileEvent>,
23 debounce_duration: Duration,
24 ) -> Result<Self, notify::Error> {
25 Self::with_pending(event_tx, debounce_duration, None)
26 }
27
28 pub fn with_pending(
30 event_tx: tokio_mpsc::Sender<FileEvent>,
31 debounce_duration: Duration,
32 pending: Option<Arc<AtomicUsize>>,
33 ) -> Result<Self, notify::Error> {
34 let (tx, rx) = mpsc::channel();
35
36 let event_tx_clone = event_tx.clone();
38 std::thread::spawn(move || {
39 while let Ok(result) = rx.recv() {
40 if let Err(e) = handle_debounced_events(result, &event_tx_clone, pending.as_deref())
41 {
42 error!("Error handling file events: {e}");
43 }
44 }
45 });
46
47 let debouncer = new_debouncer(debounce_duration, None, move |result| {
48 let _ = tx.send(result);
49 })?;
50
51 Ok(Self { debouncer })
52 }
53
54 pub fn watch(&mut self, path: &Path) -> Result<(), notify_debouncer_full::notify::Error> {
56 debug!("Starting to watch: {:?}", path);
57 self.debouncer.watch(path, RecursiveMode::Recursive)?;
58 Ok(())
59 }
60
61 pub fn unwatch(&mut self, path: &Path) -> Result<(), notify_debouncer_full::notify::Error> {
63 debug!("Stopping watch: {:?}", path);
64 self.debouncer.unwatch(path)?;
65 Ok(())
66 }
67}
68
69fn handle_debounced_events(
70 result: DebounceEventResult,
71 event_tx: &tokio_mpsc::Sender<FileEvent>,
72 pending: Option<&AtomicUsize>,
73) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
74 match result {
75 Ok(events) => {
76 for event in events {
77 if let Some(file_event) = convert_event(&event) {
78 if let Some(pending) = pending {
79 pending.fetch_add(1, Ordering::SeqCst);
80 }
81 if event_tx.blocking_send(file_event).is_err() {
83 if let Some(pending) = pending {
84 pending.fetch_sub(1, Ordering::SeqCst);
85 }
86 warn!("Event channel closed");
87 break;
88 }
89 }
90 }
91 }
92 Err(errors) => {
93 for error in errors {
94 error!("Watch error: {error}");
95 }
96 }
97 }
98 Ok(())
99}
100
101fn convert_event(event: ¬ify_debouncer_full::DebouncedEvent) -> Option<FileEvent> {
102 use notify_debouncer_full::notify::EventKind;
103
104 let path = event.paths.first()?.clone();
105
106 if path
108 .file_name()
109 .is_some_and(|name| name.to_string_lossy().starts_with('.'))
110 {
111 return None;
112 }
113
114 match &event.kind {
115 EventKind::Create(_) => Some(FileEvent::Created(path)),
116 EventKind::Modify(_) => Some(FileEvent::Modified(path)),
117 EventKind::Remove(_) => Some(FileEvent::Deleted(path)),
118 EventKind::Other => {
119 if event.paths.len() >= 2 {
121 Some(FileEvent::Renamed {
122 from: event.paths[0].clone(),
123 to: event.paths[1].clone(),
124 })
125 } else {
126 None
127 }
128 }
129 _ => None,
130 }
131}
132
133#[cfg(test)]
134mod tests {
135 use super::*;
136 use notify_debouncer_full::DebouncedEvent;
137 use notify_debouncer_full::notify::EventKind;
138 use notify_debouncer_full::notify::event::{CreateKind, ModifyKind, RemoveKind};
139 use std::path::PathBuf;
140 use std::time::Instant;
141
142 fn make_event(kind: EventKind, paths: Vec<PathBuf>) -> DebouncedEvent {
143 DebouncedEvent {
144 event: notify_debouncer_full::notify::Event {
145 kind,
146 paths,
147 attrs: Default::default(),
148 },
149 time: Instant::now(),
150 }
151 }
152
153 #[test]
154 fn test_convert_event_create() {
155 let path = PathBuf::from("/tmp/test.txt");
156 let event = make_event(EventKind::Create(CreateKind::File), vec![path.clone()]);
157
158 let result = convert_event(&event);
159 assert!(matches!(result, Some(FileEvent::Created(p)) if p == path));
160 }
161
162 #[test]
163 fn test_convert_event_modify() {
164 use notify_debouncer_full::notify::event::DataChange;
165 let path = PathBuf::from("/tmp/test.txt");
166 let event = make_event(
167 EventKind::Modify(ModifyKind::Data(DataChange::Any)),
168 vec![path.clone()],
169 );
170
171 let result = convert_event(&event);
172 assert!(matches!(result, Some(FileEvent::Modified(p)) if p == path));
173 }
174
175 #[test]
176 fn test_convert_event_delete() {
177 let path = PathBuf::from("/tmp/test.txt");
178 let event = make_event(EventKind::Remove(RemoveKind::File), vec![path.clone()]);
179
180 let result = convert_event(&event);
181 assert!(matches!(result, Some(FileEvent::Deleted(p)) if p == path));
182 }
183
184 #[test]
185 fn test_hidden_files_skipped() {
186 let path = PathBuf::from("/tmp/.hidden");
187 let event = make_event(EventKind::Create(CreateKind::File), vec![path]);
188
189 let result = convert_event(&event);
190 assert!(result.is_none());
191 }
192}