use std::fs::File;
use std::io::{self, BufRead, BufReader, Read, Seek, Write};
use std::path::PathBuf;
use std::time::{Duration, Instant};
use bytes::BytesMut;
use miette::IntoDiagnostic;
use super::FileProcessor;
use crate::defaults::io::*;
use crate::defaults::processing::*;
use crate::errors::TaleError;
use crate::{config, process_line, strip_line_ending};
pub struct BackSeekingProcessor<'a> {
fpath: PathBuf,
initial_file_size: u64,
file: Option<File>,
outlock: io::StdoutLock<'a>,
buffer: BytesMut,
count: u16,
}
impl<'a> FileProcessor for BackSeekingProcessor<'a> {
fn process_lines<F>(&mut self, _line_processor: F) -> Result<(), TaleError>
where
F: FnMut(&str) -> Result<(), TaleError>,
{
let _temp_buffer = BytesMut::with_capacity(OUTPUT_BUFFER_CAPACITY);
let _temp_outlock = io::stdout().lock();
todo!("Refactor existing tail() method to support callback-based processing")
}
fn skip_lines(&mut self, _count: u64) -> Result<(), TaleError> {
todo!("Implement using existing offset logic")
}
fn file_size(&self) -> u64 {
self.initial_file_size
}
fn seek(&mut self, pos: io::SeekFrom) -> Result<u64, TaleError> {
if let Some(mut f) = self.file.as_ref() {
Ok(f.seek(pos).map_err(TaleError::from)?)
} else {
let mut file = File::open(&self.fpath).map_err(TaleError::from)?;
let actual = file.seek(pos).map_err(TaleError::from)?;
self.file = Some(file);
Ok(actual)
}
}
fn position(&self) -> u64 {
if let Some(mut f) = self.file.as_ref() {
f.stream_position().unwrap_or_default()
} else {
0
}
}
}
impl<'a> BackSeekingProcessor<'a> {
pub fn new(fpath: PathBuf) -> Self {
let file_size = if let Ok(mut file) = File::open(&fpath) {
file.seek(io::SeekFrom::End(0)).unwrap_or_default()
} else {
0
};
Self {
fpath,
initial_file_size: file_size,
file: None,
outlock: io::stdout().lock(),
buffer: BytesMut::with_capacity(OUTPUT_BUFFER_CAPACITY),
count: 0,
}
}
pub fn move_to_position(
&mut self,
offset: i64,
units: config::OffsetUnit,
tailing: bool,
) -> Result<File, TaleError> {
let mut file = File::open(&self.fpath)?;
let file_size = file.seek(io::SeekFrom::End(0))?;
if file_size == 0 {
return Ok(file);
}
file.seek(io::SeekFrom::Start(0))?;
match units {
config::OffsetUnit::Lines => {
if offset > 0 {
file.seek(io::SeekFrom::Start(0))?;
} else if offset < 0 {
let start = self.move_n_lines_back(&mut file, (-offset) as u64)?;
file.seek(io::SeekFrom::Start(start))?;
} else if tailing {
file.seek(io::SeekFrom::End(0))?;
}
}
config::OffsetUnit::Bytes => {
if offset > 0 {
file.seek(io::SeekFrom::Start(offset as u64))?;
} else if offset < 0 {
file.seek(io::SeekFrom::End(offset))?;
} else if tailing {
file.seek(io::SeekFrom::End(0))?;
}
}
config::OffsetUnit::Blocks => {
if offset > 0 {
let byte_offset = (offset as u64) * BLOCK_SIZE;
file.seek(io::SeekFrom::Start(byte_offset))?;
} else if offset < 0 {
let byte_offset = offset * (BLOCK_SIZE as i64);
file.seek(io::SeekFrom::End(byte_offset))?;
} else if tailing {
file.seek(io::SeekFrom::End(0))?;
}
}
}
Ok(file)
}
fn move_n_lines_back(&mut self, file: &mut File, line_count: u64) -> Result<u64, TaleError> {
let file_size = file.seek(io::SeekFrom::End(0))?;
if file_size == 0 {
return Ok(0);
}
const BUFFER_SIZE: usize = 8192;
let mut buffer = vec![0u8; BUFFER_SIZE];
let mut lines_found = 0u64;
file.seek(io::SeekFrom::End(-1))?;
let mut last_byte = [0u8; 1];
file.read_exact(&mut last_byte)?;
let ends_with_newline = last_byte[0] == b'\n';
let target_newlines = if ends_with_newline { line_count } else { line_count - 1 };
let mut pos = file_size;
loop {
let chunk_size = std::cmp::min(BUFFER_SIZE as u64, pos) as usize;
if chunk_size == 0 {
return Ok(0);
}
pos -= chunk_size as u64;
file.seek(io::SeekFrom::Start(pos))?;
file.read_exact(&mut buffer[..chunk_size])?;
for (i, &byte) in buffer[..chunk_size].iter().enumerate().rev() {
if byte == b'\n' {
lines_found += 1;
if lines_found > target_newlines {
return Ok(pos + i as u64 + 1);
}
}
}
if pos == 0 {
return Ok(0);
}
}
}
pub fn process_line(&mut self, line: &str) -> Result<(), TaleError> {
process_line(line, &mut self.buffer, &mut self.outlock)
.map_err(|e| TaleError::from(std::io::Error::other(e.to_string())))?;
self.count += 1;
self.flush_if_needed()
}
pub fn flush_if_needed(&mut self) -> Result<(), TaleError> {
if self.count >= FLUSH_LINE_COUNT {
self.outlock.flush()?;
self.count = 0;
}
Ok(())
}
pub fn flush(&mut self) -> Result<(), TaleError> {
self.outlock.flush()?;
self.count = 0;
Ok(())
}
pub fn tail(&mut self) -> miette::Result<()> {
let tailing = config::tailing();
let offset_unit = config::offset_unit();
let offset = config::offset();
let file = self
.move_to_position(offset, offset_unit, tailing)
.map_err(miette::Report::from)?;
let mut reader = BufReader::new(file);
if offset > 0 && matches!(offset_unit, config::OffsetUnit::Lines) {
let consume_me = (&mut reader).lines().take(offset as usize);
let _count = consume_me.count();
};
let mut line = String::with_capacity(LINE_CAPACITY);
while reader.read_line(&mut line).into_diagnostic()? != 0 {
strip_line_ending(&mut line);
self.process_line(line.as_str()).map_err(miette::Report::from)?;
line.clear();
}
self.flush().map_err(miette::Report::from)?;
if !tailing {
return Ok(());
}
let mut last_flush = Instant::now();
let mut file = reader.into_inner();
let mut file_position = file.stream_position().into_diagnostic()?;
loop {
std::thread::sleep(Duration::from_millis(100));
let current_size = file.seek(io::SeekFrom::End(0)).into_diagnostic()?;
if current_size > file_position {
file.seek(io::SeekFrom::Start(file_position)).into_diagnostic()?;
let mut tail_reader = BufReader::new(&file);
match tail_reader.read_line(&mut line).into_diagnostic()? {
0 => {
continue;
}
_ => {
strip_line_ending(&mut line);
process_line(&line, &mut self.buffer, &mut self.outlock)?;
if last_flush.elapsed() >= TAIL_FLUSH_INTERVAL {
self.outlock.flush().into_diagnostic()?;
last_flush = Instant::now();
}
line.clear();
self.buffer.clear();
}
}
file_position = file.stream_position().into_diagnostic()?;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn seeking_backwards() {
use std::io::{Read, Seek, Write};
use tempfile::NamedTempFile;
let mut temp_file = NamedTempFile::new().expect("Failed to create temp file");
let content = "line1\nline2\nline3\nline4\nline5\n";
temp_file
.write_all(content.as_bytes())
.expect("Failed to write to temp file");
let pathbuf = PathBuf::from(temp_file.path());
let mut processor = BackSeekingProcessor::new(pathbuf);
let mut file = File::open(temp_file.path()).expect("Failed to open temp file");
let pos = processor
.move_n_lines_back(&mut file, 2)
.expect("Failed to find position");
file.seek(io::SeekFrom::Start(pos)).expect("Failed to seek");
let mut remaining = String::new();
file.read_to_string(&mut remaining).expect("Failed to read remaining");
assert_eq!(remaining, "line4\nline5\n");
let pos = processor
.move_n_lines_back(&mut file, 1)
.expect("Failed to find position");
file.seek(io::SeekFrom::Start(pos)).expect("Failed to seek");
let mut remaining = String::new();
file.read_to_string(&mut remaining).expect("Failed to read remaining");
assert_eq!(remaining, "line5\n");
let pos = processor
.move_n_lines_back(&mut file, 10)
.expect("Failed to find position");
assert_eq!(pos, 0);
}
#[test]
fn seeking_in_empty() {
use tempfile::NamedTempFile;
let temp_file = NamedTempFile::new().expect("Failed to create temp file");
let pathbuf = PathBuf::from(temp_file.path());
let mut processor = BackSeekingProcessor::new(pathbuf);
let mut file = File::open(temp_file.path()).expect("Failed to open temp file");
let pos = processor
.move_n_lines_back(&mut file, 5)
.expect("Failed to find position");
assert_eq!(pos, 0);
}
#[test]
fn good_circular_buffer_byte_logic() {
let input_data = b"0123456789abcdefghij";
let buffer_size = 10;
let mut circular_buffer = vec![0u8; buffer_size];
let mut pos = 0usize;
for &byte in input_data {
circular_buffer[pos % buffer_size] = byte;
pos += 1;
}
let _total_read = input_data.len() as u64;
let _bytes_to_show = buffer_size as u64;
let start_pos = pos % buffer_size;
let mut result = Vec::with_capacity(buffer_size);
for i in 0..buffer_size {
result.push(circular_buffer[(start_pos + i) % buffer_size]);
}
assert_eq!(result, b"abcdefghij");
}
#[test]
fn circular_buffer_partial_fill_works() {
let input_data = b"hello";
let buffer_size = 10;
let mut circular_buffer = vec![0u8; buffer_size];
for (i, &byte) in input_data.iter().enumerate() {
circular_buffer[i] = byte;
}
let bytes_to_output = input_data.len();
let result = &circular_buffer[..bytes_to_output];
assert_eq!(result, b"hello");
}
#[test]
fn circular_line_logic() {
use std::collections::VecDeque;
let lines_to_keep = 3;
let mut line_buffer: VecDeque<String> = VecDeque::with_capacity(lines_to_keep);
let input_lines = vec!["line1", "line2", "line3", "line4", "line5"];
for line in input_lines {
if line_buffer.len() >= lines_to_keep {
line_buffer.pop_front();
}
line_buffer.push_back(line.to_string());
}
let result: Vec<String> = line_buffer.into_iter().collect();
assert_eq!(result, vec!["line3", "line4", "line5"]);
}
#[test]
fn handles_overshoots_correctly() {
let overshoot = b"partial line\ncomplete line\nanother";
let overshoot_str = String::from_utf8_lossy(overshoot);
let mut remaining = overshoot_str.as_ref();
let mut complete_lines = Vec::new();
while let Some(pos) = remaining.find('\n') {
complete_lines.push(&remaining[..pos]);
remaining = &remaining[pos + 1..];
}
assert_eq!(complete_lines, vec!["partial line", "complete line"]);
assert_eq!(remaining, "another");
}
}