pub mod error;
pub use error::RegistryError;
use crate::registry::error::Result;
use crate::repository::{Repository as BaseRepository, scan_journal_files};
use crate::{File, FileInfo, TimeRange};
use journal_common::Seconds;
use journal_common::collections::{HashMap, HashSet};
use notify::{
Event,
event::{EventKind, ModifyKind, RenameMode},
};
use parking_lot::RwLock;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use tracing::{error, info, trace, warn};
mod monitor;
pub use monitor::Monitor;
struct Repository {
base: BaseRepository,
file_metadata: HashMap<File, FileInfo>,
}
impl Repository {
fn new() -> Self {
Self {
base: BaseRepository::default(),
file_metadata: HashMap::default(),
}
}
fn insert(&mut self, file: File) -> Result<()> {
let file_info = FileInfo {
file: file.clone(),
time_range: TimeRange::Unknown,
};
self.base.insert(file.clone())?;
self.file_metadata.insert(file, file_info);
Ok(())
}
fn remove(&mut self, file: &File) -> Result<()> {
self.base.remove(file)?;
self.file_metadata.remove(file);
Ok(())
}
fn remove_directory(&mut self, path: &str) {
self.base.remove_directory(path);
self.file_metadata
.retain(|file, _| file.dir().ok().map(|dir| dir != path).unwrap_or(true));
}
fn find_files_in_range(&self, start: Seconds, end: Seconds) -> Vec<FileInfo> {
let files: Vec<File> = self.base.find_files_in_range(start, end);
files
.into_iter()
.filter_map(|file| {
let file_info =
self.file_metadata
.get(&file)
.cloned()
.unwrap_or_else(|| FileInfo {
file: file.clone(),
time_range: TimeRange::Unknown,
});
let include = match file_info.time_range {
TimeRange::Unknown => {
true
}
TimeRange::Active { end: _file_end, .. } => {
true
}
TimeRange::Bounded {
start: file_start,
end: file_end,
..
} => {
file_start.0 < end.0 && file_end.0 > start.0
}
};
if include { Some(file_info) } else { None }
})
.collect()
}
fn update_file_info(&mut self, file_info: FileInfo) {
let file = file_info.file.clone();
self.file_metadata.insert(file, file_info);
}
}
impl Default for Repository {
fn default() -> Self {
Self::new()
}
}
struct RegistryInner {
repository: Repository,
watched_directories: HashSet<String>,
monitor: Monitor,
}
#[derive(Clone)]
pub struct Registry {
inner: Arc<RwLock<RegistryInner>>,
}
impl Registry {
pub fn new(monitor: Monitor) -> Self {
let inner = RegistryInner {
repository: Repository::new(),
watched_directories: HashSet::default(),
monitor,
};
Self {
inner: Arc::new(RwLock::new(inner)),
}
}
pub fn watch_directory(&self, path: &str) -> Result<()> {
let mut inner = self.inner.write();
if inner.watched_directories.contains(path) {
warn!("Directory {} is already being watched", path);
return Ok(());
}
info!("scanning directory: {}", path);
let files = scan_journal_files(path)?;
info!("found {} journal files in {}", files.len(), path);
inner.monitor.watch_directory(path)?;
inner.watched_directories.insert(String::from(path));
for file in files {
trace!("adding file to repository: {:?}", file.path());
if let Err(e) = inner.repository.insert(file) {
error!("failed to insert file into repository: {}", e);
}
}
info!(
"now watching directory: {} (total directories: {})",
path,
inner.watched_directories.len()
);
Ok(())
}
pub fn unwatch_directory(&self, path: &str) -> Result<()> {
let mut inner = self.inner.write();
if !inner.watched_directories.contains(path) {
warn!("directory {} is not being watched", path);
return Ok(());
}
inner.monitor.unwatch_directory(path)?;
inner.repository.remove_directory(path); inner.watched_directories.remove(path);
info!("stopped watching directory: {}", path);
Ok(())
}
pub fn process_event(&self, event: Event) -> Result<()> {
let mut inner = self.inner.write();
match event.kind {
EventKind::Create(_) => {
for path in &event.paths {
Self::insert_event_path(&mut inner, path);
}
}
EventKind::Remove(_) => {
for path in &event.paths {
Self::remove_event_path(&mut inner, path);
}
}
EventKind::Modify(ModifyKind::Name(RenameMode::Both)) => {
Self::process_rename_event(&mut inner, &event.paths);
}
EventKind::Modify(ModifyKind::Name(rename_mode)) => {
info!(
"unhandled modify event: '{:?}', expecting rename event for newly archived file",
rename_mode
);
}
_ => {}
}
Ok(())
}
fn insert_event_path(inner: &mut RegistryInner, path: &Path) {
trace!("adding file to repository: {:?}", path);
let Some(file) = File::from_path(path) else {
warn!("path is not a valid journal file: {:?}", path);
return;
};
if let Err(e) = inner.repository.insert(file) {
error!("failed to insert file: {}", e);
}
}
fn remove_event_path(inner: &mut RegistryInner, path: &Path) {
trace!("removing file from repository: {:?}", path);
let Some(file) = File::from_path(path) else {
warn!("path is not a valid journal file: {:?}", path);
return;
};
if let Err(e) = inner.repository.remove(&file) {
error!("failed to remove file: {}", e);
}
}
fn process_rename_event(inner: &mut RegistryInner, paths: &[PathBuf]) {
let Some((old_path, new_path)) = paths.first().zip(paths.get(1)) else {
error!("rename event with unexpected path count: {:#?}", paths);
return;
};
info!("rename event: {:?} -> {:?}", old_path, new_path);
Self::remove_renamed_old_path(inner, old_path);
Self::insert_renamed_new_path(inner, new_path);
}
fn remove_renamed_old_path(inner: &mut RegistryInner, old_path: &Path) {
let Some(old_file) = File::from_path(old_path) else {
return;
};
info!("removing old file: {:?}", old_file.path());
if let Err(e) = inner.repository.remove(&old_file) {
error!("failed to remove old file: {}", e);
}
}
fn insert_renamed_new_path(inner: &mut RegistryInner, new_path: &Path) {
let Some(new_file) = File::from_path(new_path) else {
return;
};
info!("inserting new file: {:?}", new_file.path());
if let Err(e) = inner.repository.insert(new_file) {
error!("failed to insert new file: {}", e);
}
}
pub fn find_files_in_range(&self, start: Seconds, end: Seconds) -> Result<Vec<FileInfo>> {
let inner = self.inner.read();
Ok(inner.repository.find_files_in_range(start, end))
}
pub fn update_time_range(
&self,
file: &File,
start_time: Seconds,
end_time: Seconds,
indexed_at: Seconds,
online: bool,
) {
let mut inner = self.inner.write();
let time_range = if online {
TimeRange::Active {
start: start_time,
end: end_time,
indexed_at: indexed_at,
}
} else {
TimeRange::Bounded {
start: start_time,
end: end_time,
indexed_at: indexed_at,
}
};
let file_info = FileInfo {
file: file.clone(),
time_range,
};
inner.repository.update_file_info(file_info);
}
}