use std::collections::HashMap;
use std::fs::File;
use std::io::{self, BufRead, BufReader, Seek, SeekFrom};
#[cfg(unix)]
use std::os::unix::fs::MetadataExt;
use std::path::{Path, PathBuf};
use miette::WrapErr;
use crate::errors::TaleError;
#[derive(Debug, Clone)]
pub struct FileState {
pub path: PathBuf,
pub position: u64,
pub size: u64,
pub inode: Option<u64>,
pub available: bool,
}
impl FileState {
pub fn new(path: PathBuf) -> Self {
Self {
path,
position: 0,
size: 0,
inode: None,
available: false,
}
}
pub fn new_and_refresh(path: PathBuf) -> Result<Self, TaleError> {
let mut state = Self::new(path);
state
.refresh()
.map_err(|e| TaleError::from(std::io::Error::other(e.to_string())))?;
Ok(state)
}
pub fn set_position(&mut self, position: u64) {
self.position = position;
}
pub fn refresh(&mut self) -> miette::Result<bool> {
let metadata = match std::fs::metadata(&self.path) {
Ok(metadata) => metadata,
Err(err) if err.kind() == io::ErrorKind::NotFound => {
let was_available = self.available;
self.available = false;
self.size = 0;
self.inode = None;
return Ok(was_available); }
Err(err) => {
return Err(TaleError::from(err))
.with_context(|| format!("Failed to get metadata for {}", self.path.display()));
}
};
let new_size = metadata.len();
#[cfg(unix)]
let new_inode = Some(metadata.ino());
#[cfg(not(unix))]
let new_inode = None;
let was_available = self.available;
let old_inode = self.inode;
let old_size = self.size;
let file_rotated = if old_inode.is_some() && new_inode != old_inode {
if crate::config::sticky() {
self.inode = new_inode;
self.position = 0;
self.available = metadata.is_file();
self.size = new_size;
true
} else {
self.available = false;
self.size = 0;
self.inode = None;
true
}
} else {
self.available = metadata.is_file();
self.size = new_size;
self.inode = new_inode;
false
};
let changed = was_available != self.available || file_rotated || (self.available && new_size > old_size);
Ok(changed)
}
pub fn has_new_data(&self) -> bool {
self.available && self.size > self.position
}
pub fn update_position(&mut self, new_position: u64) {
self.position = new_position;
}
pub fn read_new_lines(&mut self) -> Result<Vec<String>, TaleError> {
if !self.available || !self.has_new_data() {
return Ok(Vec::new());
}
let mut file = File::open(&self.path).map_err(TaleError::from)?;
file.seek(SeekFrom::Start(self.position)).map_err(TaleError::from)?;
let mut reader = BufReader::new(file);
let mut lines = Vec::new();
loop {
let mut line = String::new();
match reader.read_line(&mut line).map_err(TaleError::from)? {
0 => break, bytes_read => {
if line.ends_with('\n') {
line.pop();
if line.ends_with('\r') {
line.pop();
}
}
lines.push(line);
self.position += bytes_read as u64;
}
}
}
Ok(lines)
}
}
#[derive(Debug)]
pub struct FileStateManager {
states: HashMap<PathBuf, FileState>,
}
impl FileStateManager {
pub fn new() -> Self {
Self { states: HashMap::new() }
}
pub fn add_file<P: AsRef<Path>>(&mut self, path: P) -> Result<(), TaleError> {
let path = path.as_ref().to_path_buf();
let state = FileState::new_and_refresh(path.clone())?;
self.states.insert(path, state);
Ok(())
}
pub fn add_file_for_tailing<P: AsRef<Path>>(&mut self, path: P) -> Result<(), TaleError> {
let path = path.as_ref().to_path_buf();
let mut state = FileState::new_and_refresh(path.clone())?;
state.set_position(state.size);
self.states.insert(path, state);
Ok(())
}
pub fn remove_file<P: AsRef<Path>>(&mut self, path: P) {
self.states.remove(path.as_ref());
}
pub fn get_state<P: AsRef<Path>>(&self, path: P) -> Option<&FileState> {
self.states.get(path.as_ref())
}
pub fn get_state_mut<P: AsRef<Path>>(&mut self, path: P) -> Option<&mut FileState> {
self.states.get_mut(path.as_ref())
}
pub fn refresh_all(&mut self) -> Result<Vec<PathBuf>, TaleError> {
let mut changed_files = Vec::new();
for (path, state) in &mut self.states {
if let Ok(changed) = state.refresh()
&& changed
{
changed_files.push(path.clone());
}
}
Ok(changed_files)
}
pub fn files_with_new_data(&self) -> Vec<&PathBuf> {
self.states
.iter()
.filter_map(|(path, state)| if state.has_new_data() { Some(path) } else { None })
.collect()
}
pub fn tracked_files(&self) -> Vec<&PathBuf> {
self.states.keys().collect()
}
pub fn read_new_lines(&mut self) -> Result<Vec<(PathBuf, Vec<String>)>, TaleError> {
let mut all_new_lines = Vec::new();
for (path, state) in &mut self.states {
if state.has_new_data() {
let lines = state.read_new_lines()?;
if !lines.is_empty() {
all_new_lines.push((path.clone(), lines));
}
}
}
Ok(all_new_lines)
}
pub fn read_all_lines(&mut self) -> Result<Vec<(PathBuf, Vec<String>)>, crate::errors::TaleError> {
let mut all_lines = Vec::new();
for (path, state) in &self.states {
if state.available {
let lines = Self::read_all_lines_from_file(state)?;
if !lines.is_empty() {
all_lines.push((path.clone(), lines));
}
}
}
Ok(all_lines)
}
fn read_all_lines_from_file(state: &FileState) -> Result<Vec<String>, TaleError> {
if !state.available {
return Ok(Vec::new());
}
let mut file = File::open(&state.path).map_err(TaleError::from)?;
file.seek(SeekFrom::Start(0)).map_err(TaleError::from)?;
let reader = BufReader::new(file);
let mut lines = Vec::new();
for line_result in reader.lines() {
let line = line_result.map_err(TaleError::from)?;
lines.push(line);
}
Ok(lines)
}
}
impl Default for FileStateManager {
fn default() -> Self {
Self::new()
}
}