Skip to main content

ragfs_index/
watcher.rs

1//! File system watcher for detecting changes.
2
3use 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
14/// File system watcher with debouncing.
15pub struct FileWatcher {
16    debouncer: Debouncer<RecommendedWatcher, RecommendedCache>,
17}
18
19impl FileWatcher {
20    /// Create a new file watcher.
21    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    /// Create a watcher that increments `pending` for each queued event.
29    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        // Spawn thread to convert events
37        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    /// Start watching a path.
55    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    /// Stop watching a path.
62    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                    // Use blocking send since we're in a std thread
82                    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: &notify_debouncer_full::DebouncedEvent) -> Option<FileEvent> {
102    use notify_debouncer_full::notify::EventKind;
103
104    let path = event.paths.first()?.clone();
105
106    // Skip hidden files and directories
107    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            // Handle rename as "other" event with two paths
120            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}