Skip to main content

simd_csv/
line_reader.rs

1use memchr::memchr;
2
3use std::io::{self, BufReader, Read};
4
5use crate::buffer::ScratchBuffer;
6use crate::utils::trim_trailing_crlf;
7
8/// A zero-copy & optimized line reader.
9///
10/// This reader recognizes both `LF` & `CRLF` line terminators, but not single
11/// `CR`.
12pub struct LineReader<R> {
13    inner: ScratchBuffer<R>,
14}
15
16impl<R: Read> LineReader<R> {
17    /// Create a new reader with default using the provided reader implementing
18    /// [`std::io::Read`].
19    ///
20    /// Avoid providing a buffered reader because buffering will be handled for
21    /// you by the [`LineReader`].
22    pub fn from_reader(inner: R) -> Self {
23        Self {
24            inner: ScratchBuffer::new(inner),
25        }
26    }
27
28    /// Create a new reader with provided buffer capacity and using the provided
29    /// reader implementing [`std::io::Read`].
30    ///
31    /// Avoid providing a buffered reader because buffering will be handled for
32    /// you by the [`LineReader`].
33    pub fn with_capacity(capacity: usize, inner: R) -> Self {
34        Self {
35            inner: ScratchBuffer::with_capacity(capacity, inner),
36        }
37    }
38
39    /// Consume the reader to count the number of lines as fast as possible.
40    pub fn count_lines(&mut self) -> io::Result<u64> {
41        let mut count: u64 = 0;
42        let mut current_is_empty = true;
43
44        loop {
45            let input = self.inner.fill_buf()?;
46            let len = input.len();
47
48            if len == 0 {
49                if !current_is_empty {
50                    count += 1;
51                }
52
53                return Ok(count);
54            }
55
56            match memchr(b'\n', input) {
57                None => {
58                    self.inner.consume(len);
59                    current_is_empty = false;
60                }
61                Some(pos) => {
62                    count += 1;
63                    self.inner.consume(pos + 1);
64                    current_is_empty = true;
65                }
66            };
67        }
68    }
69
70    /// Attempt to read the next line from underlying reader.
71    ///
72    /// Will return `Ok(None)` if end of stream was reached.
73    pub fn read_line(&mut self) -> io::Result<Option<&[u8]>> {
74        self.inner.reset();
75
76        loop {
77            let input = self.inner.fill_buf()?;
78            let len = input.len();
79
80            if len == 0 {
81                if self.inner.has_something_saved() {
82                    return Ok(Some(trim_trailing_crlf(self.inner.saved())));
83                }
84
85                return Ok(None);
86            }
87
88            match memchr(b'\n', input) {
89                None => {
90                    self.inner.save();
91                }
92                Some(pos) => {
93                    let bytes = self.inner.flush(pos + 1);
94                    return Ok(Some(trim_trailing_crlf(bytes)));
95                }
96            };
97        }
98    }
99
100    /// Attempt to read the next line from underlying reader but copy the line
101    /// to given buffer instead of using [`Self::read_line`] internal buffering
102    /// logic for zero-copy reading.
103    ///
104    /// Will return `Ok(false)` if end of stream was reached.
105    pub fn read_line_into(&mut self, buffer: &mut Vec<u8>) -> io::Result<bool> {
106        self.inner.reset();
107        buffer.clear();
108
109        loop {
110            let input = self.inner.fill_buf()?;
111            let len = input.len();
112
113            if len == 0 {
114                if buffer.is_empty() {
115                    return Ok(false);
116                } else {
117                    let trimmed_len = trim_trailing_crlf(buffer).len();
118
119                    buffer.truncate(trimmed_len);
120
121                    return Ok(true);
122                }
123            }
124
125            match memchr(b'\n', input) {
126                None => {
127                    buffer.extend_from_slice(input);
128                    self.inner.consume(len);
129                }
130                Some(pos) => {
131                    buffer.extend_from_slice(trim_trailing_crlf(&input[..pos + 1]));
132                    self.inner.consume(pos + 1);
133
134                    return Ok(true);
135                }
136            };
137        }
138    }
139
140    /// Return an iterator over the underlying reader's lines.
141    #[inline]
142    pub fn lines(&mut self) -> Lines<'_, R> {
143        Lines {
144            reader: self,
145            buffer: Vec::new(),
146        }
147    }
148
149    /// Return the current byte offset of the reader.
150    #[inline(always)]
151    pub fn position(&self) -> u64 {
152        self.inner.position()
153    }
154
155    /// Return the underlying [`BufReader`].
156    #[inline(always)]
157    pub fn into_bufreader(self) -> BufReader<R> {
158        self.inner.into_bufreader()
159    }
160
161    /// Return the underlying reader.
162    ///
163    /// **BEWARE**: Already buffered data will be lost!
164    #[inline(always)]
165    pub fn into_inner(self) -> R {
166        self.inner.into_bufreader().into_inner()
167    }
168}
169
170pub struct Lines<'r, R> {
171    reader: &'r mut LineReader<R>,
172    buffer: Vec<u8>,
173}
174
175impl<R: Read> Iterator for Lines<'_, R> {
176    type Item = io::Result<Vec<u8>>;
177
178    #[inline]
179    fn next(&mut self) -> Option<Self::Item> {
180        match self.reader.read_line_into(&mut self.buffer) {
181            Err(err) => Some(Err(err)),
182            Ok(true) => Some(Ok(self.buffer.clone())),
183            Ok(false) => None,
184        }
185    }
186}
187
188#[cfg(test)]
189mod tests {
190    use std::io::Cursor;
191
192    use super::*;
193
194    #[test]
195    fn test_read_line() -> io::Result<()> {
196        let tests: &[(&[u8], Vec<&[u8]>)] = &[
197            (b"", vec![]),
198            (b"test", vec![b"test"]),
199            (
200                b"hello\nwhatever\r\nbye!",
201                vec![b"hello", b"whatever", b"bye!"],
202            ),
203            (
204                b"hello\nwhatever\nbye!\n",
205                vec![b"hello", b"whatever", b"bye!"],
206            ),
207            (
208                b"hello\nwhatever\r\nbye!\n\n\r\n\n",
209                vec![b"hello", b"whatever", b"bye!", b"", b"", b""],
210            ),
211        ];
212
213        for (data, expected) in tests {
214            // Count
215            let mut reader = LineReader::from_reader(Cursor::new(data));
216
217            assert_eq!(reader.count_lines()?, expected.len() as u64);
218
219            // Zero-copy
220            let mut reader = LineReader::from_reader(Cursor::new(data));
221
222            let mut lines = Vec::new();
223
224            while let Some(line) = reader.read_line()? {
225                lines.push(line.to_vec());
226            }
227
228            assert_eq!(lines, *expected);
229
230            // Into buffer
231            let mut reader = LineReader::from_reader(Cursor::new(data));
232
233            lines.clear();
234            let mut buffer = vec![];
235
236            while reader.read_line_into(&mut buffer)? {
237                lines.push(buffer.clone());
238            }
239
240            assert_eq!(lines, *expected);
241
242            // Iterator
243            let mut reader = LineReader::from_reader(Cursor::new(data));
244
245            assert_eq!(
246                reader.lines().collect::<Result<Vec<_>, _>>().unwrap(),
247                *expected
248            );
249        }
250
251        Ok(())
252    }
253}