use crate::DEFAULT_BUF_SIZE;
use pin_project_lite::pin_project;
use std::io::{Result as IoResult, SeekFrom};
use std::pin::Pin;
use std::task::{Context, Poll};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek, AsyncSeekExt, ReadBuf};
pin_project! {
#[derive(Debug)]
pub struct RevBufReader<R> {
#[pin]
inner: R,
buf: Box<[u8]>,
pos: usize, cap: usize, file_pos: u64, file_size: u64, initialized: bool, }
}
impl<R: AsyncRead> RevBufReader<R> {
pub fn new(inner: R) -> Self {
Self::with_capacity(DEFAULT_BUF_SIZE, inner)
}
pub fn with_capacity(capacity: usize, inner: R) -> Self {
Self {
inner,
buf: vec![0; capacity].into_boxed_slice(),
pos: 0,
cap: 0,
file_pos: 0,
file_size: 0,
initialized: false,
}
}
pub fn get_ref(&self) -> &R {
&self.inner
}
pub fn get_mut(&mut self) -> &mut R {
&mut self.inner
}
pub fn get_pin_mut(self: Pin<&mut Self>) -> Pin<&mut R> {
self.project().inner
}
pub fn into_inner(self) -> R {
self.inner
}
pub fn buffer(&self) -> &[u8] {
&self.buf[self.pos..self.cap]
}
}
impl<R: AsyncRead + AsyncSeek + Unpin> RevBufReader<R> {
async fn ensure_initialized(&mut self) -> IoResult<()> {
if !self.initialized {
self.file_size = self.inner.seek(SeekFrom::End(0)).await?;
self.file_pos = self.file_size;
self.initialized = true;
}
Ok(())
}
async fn seek_back(&mut self, length: usize) -> IoResult<usize> {
if self.file_pos == 0 {
return Ok(0);
}
let seek_amount = std::cmp::min(length as u64, self.file_pos) as usize;
let new_pos = self.file_pos - seek_amount as u64;
self.inner.seek(SeekFrom::Start(new_pos)).await?;
self.file_pos = new_pos;
self.cap = 0;
self.pos = 0;
Ok(seek_amount)
}
async fn fill_buffer(&mut self) -> IoResult<&[u8]> {
if self.pos == 0 {
let length = self.seek_back(self.buf.len()).await?;
if length == 0 {
return Ok(&[]);
}
let mut total_read = 0;
while total_read < length {
match self.inner.read(&mut self.buf[total_read..length]).await? {
0 => break,
n => total_read += n,
}
}
self.cap = total_read;
self.pos = total_read; }
Ok(&self.buf[0..self.pos])
}
fn consume(&mut self, amt: usize) {
self.pos = self.pos.saturating_sub(amt);
}
async fn read_line_internal(&mut self, buf: &mut String) -> IoResult<usize> {
self.ensure_initialized().await?;
if self.file_size == 0 {
return Ok(0);
}
let mut line_buffer = Vec::new();
loop {
let (buffer_slice, current_pos) = {
let buffer_data = self.fill_buffer().await?;
if buffer_data.is_empty() {
break;
}
(buffer_data.to_vec(), self.pos)
};
if let Some(newline_pos) = buffer_slice.iter().rposition(|&b| b == b'\n' || b == b'\r')
{
let line_start = newline_pos + 1;
let line_data = &buffer_slice[line_start..current_pos];
let mut new_line = line_data.to_vec();
new_line.extend_from_slice(&line_buffer);
line_buffer = new_line;
self.consume(current_pos - newline_pos);
let line_str = String::from_utf8_lossy(&line_buffer);
let trimmed = line_str.trim_end_matches('\r');
if !trimmed.is_empty() {
buf.push_str(trimmed);
return Ok(trimmed.len());
}
line_buffer.clear();
} else {
let mut new_line = buffer_slice;
new_line.extend_from_slice(&line_buffer);
line_buffer = new_line;
self.consume(current_pos);
if self.file_pos == 0 && self.pos == 0 {
if !line_buffer.is_empty() {
let line_str = String::from_utf8_lossy(&line_buffer);
let trimmed = line_str.trim_end_matches('\r');
buf.push_str(trimmed);
return Ok(trimmed.len());
}
break;
}
}
}
Ok(0)
}
pub async fn next_line(&mut self) -> IoResult<Option<String>> {
let mut line = String::new();
match self.read_line_internal(&mut line).await? {
0 => Ok(None),
_ => Ok(Some(line)),
}
}
pub fn lines(self) -> crate::Lines<R>
where
R: AsyncRead + AsyncSeek + Unpin,
{
crate::Lines::new(self)
}
}
impl<R: AsyncRead + Unpin> AsyncRead for RevBufReader<R> {
fn poll_read(
self: Pin<&mut Self>,
_cx: &mut Context<'_>,
_buf: &mut ReadBuf<'_>,
) -> Poll<IoResult<()>> {
Poll::Ready(Ok(()))
}
}