use std::path::{Path, PathBuf};
use std::time::Duration;
use miette::Result;
use notify::{Event, EventKind, RecommendedWatcher, RecursiveMode, Watcher};
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use super::file_state::FileStateManager;
use crate::errors::*;
#[derive(Debug, Clone)]
pub enum WatchEvent {
FileModified(PathBuf),
FileCreated(PathBuf),
FileDeleted(PathBuf),
FileMoved { from: PathBuf, to: PathBuf },
Error(String),
}
#[derive(Debug, Clone)]
pub struct WatcherConfig {
pub poll_interval: Duration,
pub event_buffer_size: usize,
pub use_polling: bool,
}
impl Default for WatcherConfig {
fn default() -> Self {
Self {
poll_interval: Duration::from_millis(100),
event_buffer_size: 1000,
use_polling: false,
}
}
}
pub struct MultiFileWatcher {
_config: WatcherConfig,
file_manager: FileStateManager,
event_sender: Option<mpsc::UnboundedSender<WatchEvent>>,
_watcher: Option<RecommendedWatcher>,
_task_handle: Option<JoinHandle<()>>,
}
impl MultiFileWatcher {
pub fn new(_config: WatcherConfig) -> Self {
Self {
_config,
file_manager: FileStateManager::new(),
event_sender: None,
_watcher: None,
_task_handle: None,
}
}
pub async fn add_files<I, P>(&mut self, paths: I) -> Result<(), TaleError>
where
I: IntoIterator<Item = P>,
P: AsRef<Path>,
{
for path in paths {
self.file_manager.add_file_for_tailing(path)?;
}
Ok(())
}
pub async fn watch(&mut self) -> Result<mpsc::UnboundedReceiver<WatchEvent>, TaleError> {
let (event_sender, event_receiver) = mpsc::unbounded_channel();
let (notify_sender, notify_receiver) = std::sync::mpsc::channel();
let mut watcher = notify::recommended_watcher(move |result: notify::Result<Event>| {
if let Err(e) = notify_sender.send(result) {
eprintln!("Failed to send notify event: {e}");
}
})?;
for file_path in self.file_manager.tracked_files() {
watcher
.watch(file_path, RecursiveMode::NonRecursive)
.map_err(TaleError::NotifyError)?;
}
let event_sender_clone = event_sender.clone();
let task_handle = tokio::spawn(async move {
while let Ok(result) = notify_receiver.recv() {
match result {
Ok(event) => {
if let Some(watch_event) = Self::convert_notify_event(event)
&& event_sender_clone.send(watch_event).is_err()
{
break; }
}
Err(e) => {
let error_event = WatchEvent::Error(format!("Notify error: {e}"));
if event_sender_clone.send(error_event).is_err() {
break; }
}
}
}
});
self.event_sender = Some(event_sender);
self._watcher = Some(watcher);
self._task_handle = Some(task_handle);
Ok(event_receiver)
}
pub async fn stop(&mut self) -> Result<(), TaleError> {
self._watcher = None;
self.event_sender = None;
if let Some(handle) = self._task_handle.take() {
let _ = handle.await;
}
Ok(())
}
fn convert_notify_event(event: Event) -> Option<WatchEvent> {
match event.kind {
EventKind::Modify(_) => {
event.paths.first().map(|path| WatchEvent::FileModified(path.clone()))
}
EventKind::Create(_) => {
event.paths.first().map(|path| WatchEvent::FileCreated(path.clone()))
}
EventKind::Remove(_) => {
event.paths.first().map(|path| WatchEvent::FileDeleted(path.clone()))
}
_ => {
None
}
}
}
pub fn file_manager(&self) -> &FileStateManager {
&self.file_manager
}
pub fn file_manager_mut(&mut self) -> &mut FileStateManager {
&mut self.file_manager
}
}
pub fn create_watcher() -> MultiFileWatcher {
MultiFileWatcher::new(WatcherConfig::default())
}
pub fn create_watcher_with_config(config: WatcherConfig) -> MultiFileWatcher {
MultiFileWatcher::new(config)
}