async_rev_buf/
buf_reader.rs1use 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 #[derive(Debug)]
11 pub struct RevBufReader<R> {
12 #[pin]
13 inner: R,
14 buf: Box<[u8]>,
15 pos: usize, cap: usize, file_pos: u64, file_size: u64, initialized: bool, }
21}
22
23impl<R: AsyncRead> RevBufReader<R> {
24 pub fn new(inner: R) -> Self {
26 Self::with_capacity(DEFAULT_BUF_SIZE, inner)
27 }
28
29 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 pub fn get_ref(&self) -> &R {
44 &self.inner
45 }
46
47 pub fn get_mut(&mut self) -> &mut R {
49 &mut self.inner
50 }
51
52 pub fn get_pin_mut(self: Pin<&mut Self>) -> Pin<&mut R> {
54 self.project().inner
55 }
56
57 pub fn into_inner(self) -> R {
59 self.inner
60 }
61
62 pub fn buffer(&self) -> &[u8] {
64 &self.buf[self.pos..self.cap]
65 }
66}
67
68impl<R: AsyncRead + AsyncSeek + Unpin> RevBufReader<R> {
69 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 async fn seek_back(&mut self, length: usize) -> IoResult<usize> {
81 if self.file_pos == 0 {
82 return Ok(0);
83 }
84
85 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 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 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 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; }
118 Ok(&self.buf[0..self.pos])
119 }
120
121 fn consume(&mut self, amt: usize) {
123 self.pos = self.pos.saturating_sub(amt);
124 }
125
126 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 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 if let Some(newline_pos) = buffer_slice.iter().rposition(|&b| b == b'\n' || b == b'\r')
148 {
149 let line_start = newline_pos + 1;
151 let line_data = &buffer_slice[line_start..current_pos];
152
153 let mut new_line = line_data.to_vec();
155 new_line.extend_from_slice(&line_buffer);
156 line_buffer = new_line;
157
158 self.consume(current_pos - newline_pos);
160
161 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 line_buffer.clear();
170 } else {
171 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 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 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 pub fn lines(self) -> crate::Lines<R>
205 where
206 R: AsyncRead + AsyncSeek + Unpin,
207 {
208 crate::Lines::new(self)
209 }
210}
211
212impl<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}