Skip to main content

async_rev_buf/
buf_reader.rs

1use crate::DEFAULT_BUF_SIZE;
2use pin_project_lite::pin_project;
3use std::io::{Result as IoResult, SeekFrom};
4use std::pin::Pin;
5use std::task::{Context, Poll};
6use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeek, AsyncSeekExt, ReadBuf};
7
8pin_project! {
9    /// A high-performance buffered reader that reads lines in reverse order.
10    #[derive(Debug)]
11    pub struct RevBufReader<R> {
12        #[pin]
13        inner: R,
14        buf: Box<[u8]>,
15        pos: usize,        // Current position in buffer (reading from end to start)
16        cap: usize,        // Amount of valid data in buffer
17        file_pos: u64,     // Current position in file
18        file_size: u64,    // Total file size (cached)
19        initialized: bool, // Whether we've initialized file size
20    }
21}
22
23impl<R: AsyncRead> RevBufReader<R> {
24    /// Creates a new reverse buffered reader with default capacity.
25    pub fn new(inner: R) -> Self {
26        Self::with_capacity(DEFAULT_BUF_SIZE, inner)
27    }
28
29    /// Creates a new reverse buffered reader with the specified capacity.
30    pub fn with_capacity(capacity: usize, inner: R) -> Self {
31        Self {
32            inner,
33            buf: vec![0; capacity].into_boxed_slice(),
34            pos: 0,
35            cap: 0,
36            file_pos: 0,
37            file_size: 0,
38            initialized: false,
39        }
40    }
41
42    /// Gets a reference to the underlying reader.
43    pub fn get_ref(&self) -> &R {
44        &self.inner
45    }
46
47    /// Gets a mutable reference to the underlying reader.
48    pub fn get_mut(&mut self) -> &mut R {
49        &mut self.inner
50    }
51
52    /// Gets a pinned mutable reference to the underlying reader.
53    pub fn get_pin_mut(self: Pin<&mut Self>) -> Pin<&mut R> {
54        self.project().inner
55    }
56
57    /// Consumes this reader, returning the underlying reader.
58    pub fn into_inner(self) -> R {
59        self.inner
60    }
61
62    /// Returns a reference to the internally buffered data.
63    pub fn buffer(&self) -> &[u8] {
64        &self.buf[self.pos..self.cap]
65    }
66}
67
68impl<R: AsyncRead + AsyncSeek + Unpin> RevBufReader<R> {
69    /// Initialize file size if not already done - optimized to do this once
70    async fn ensure_initialized(&mut self) -> IoResult<()> {
71        if !self.initialized {
72            self.file_size = self.inner.seek(SeekFrom::End(0)).await?;
73            self.file_pos = self.file_size;
74            self.initialized = true;
75        }
76        Ok(())
77    }
78
79    /// Efficiently seek backward in the file
80    async fn seek_back(&mut self, length: usize) -> IoResult<usize> {
81        if self.file_pos == 0 {
82            return Ok(0);
83        }
84
85        // Calculate how much we can actually seek back
86        let seek_amount = std::cmp::min(length as u64, self.file_pos) as usize;
87        let new_pos = self.file_pos - seek_amount as u64;
88
89        // Seek to the new position
90        self.inner.seek(SeekFrom::Start(new_pos)).await?;
91        self.file_pos = new_pos;
92        self.cap = 0;
93        self.pos = 0;
94
95        Ok(seek_amount)
96    }
97
98    /// Fill the buffer with data from the current position
99    async fn fill_buffer(&mut self) -> IoResult<&[u8]> {
100        if self.pos == 0 {
101            let length = self.seek_back(self.buf.len()).await?;
102            if length == 0 {
103                return Ok(&[]);
104            }
105
106            // Read the data from current position
107            let mut total_read = 0;
108            while total_read < length {
109                match self.inner.read(&mut self.buf[total_read..length]).await? {
110                    0 => break,
111                    n => total_read += n,
112                }
113            }
114
115            self.cap = total_read;
116            self.pos = total_read; // Start from the end of the buffer
117        }
118        Ok(&self.buf[0..self.pos])
119    }
120
121    /// Consume bytes from the buffer (moving backward)
122    fn consume(&mut self, amt: usize) {
123        self.pos = self.pos.saturating_sub(amt);
124    }
125
126    /// Core optimized line reading implementation
127    async fn read_line_internal(&mut self, buf: &mut String) -> IoResult<usize> {
128        self.ensure_initialized().await?;
129
130        if self.file_size == 0 {
131            return Ok(0);
132        }
133
134        let mut line_buffer = Vec::new();
135
136        loop {
137            // Get buffer data efficiently
138            let (buffer_slice, current_pos) = {
139                let buffer_data = self.fill_buffer().await?;
140                if buffer_data.is_empty() {
141                    break;
142                }
143                (buffer_data.to_vec(), self.pos)
144            };
145
146            // Search for newline from the end
147            if let Some(newline_pos) = buffer_slice.iter().rposition(|&b| b == b'\n' || b == b'\r')
148            {
149                // Found a newline - extract the line after it
150                let line_start = newline_pos + 1;
151                let line_data = &buffer_slice[line_start..current_pos];
152
153                // Build the line (prepend since we're reading backward)
154                let mut new_line = line_data.to_vec();
155                new_line.extend_from_slice(&line_buffer);
156                line_buffer = new_line;
157
158                // Consume up to and including the newline
159                self.consume(current_pos - newline_pos);
160
161                // Convert to string
162                let line_str = String::from_utf8_lossy(&line_buffer);
163                let trimmed = line_str.trim_end_matches('\r');
164                if !trimmed.is_empty() {
165                    buf.push_str(trimmed);
166                    return Ok(trimmed.len());
167                }
168                // Empty line, continue to next
169                line_buffer.clear();
170            } else {
171                // No newline found - consume entire buffer
172                let mut new_line = buffer_slice;
173                new_line.extend_from_slice(&line_buffer);
174                line_buffer = new_line;
175
176                self.consume(current_pos);
177
178                if self.file_pos == 0 && self.pos == 0 {
179                    // Reached start of file
180                    if !line_buffer.is_empty() {
181                        let line_str = String::from_utf8_lossy(&line_buffer);
182                        let trimmed = line_str.trim_end_matches('\r');
183                        buf.push_str(trimmed);
184                        return Ok(trimmed.len());
185                    }
186                    break;
187                }
188            }
189        }
190
191        Ok(0)
192    }
193
194    /// Read the next line in reverse order
195    pub async fn next_line(&mut self) -> IoResult<Option<String>> {
196        let mut line = String::new();
197        match self.read_line_internal(&mut line).await? {
198            0 => Ok(None),
199            _ => Ok(Some(line)),
200        }
201    }
202
203    /// Returns a stream of lines read in reverse order
204    pub fn lines(self) -> crate::Lines<R>
205    where
206        R: AsyncRead + AsyncSeek + Unpin,
207    {
208        crate::Lines::new(self)
209    }
210}
211
212// AsyncRead implementation for completeness
213impl<R: AsyncRead + Unpin> AsyncRead for RevBufReader<R> {
214    fn poll_read(
215        self: Pin<&mut Self>,
216        _cx: &mut Context<'_>,
217        _buf: &mut ReadBuf<'_>,
218    ) -> Poll<IoResult<()>> {
219        Poll::Ready(Ok(()))
220    }
221}