1use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
2
3use crate::buffer::BufReaderWithPosition;
4use crate::core::{CoreReader, ReadResult};
5use crate::error::{self, Error, ErrorKind};
6use crate::records::{ByteRecord, ByteRecordBuilder};
7use crate::select::{Selection, Selector};
8use crate::utils::{self, trim_bom};
9
10#[cfg(feature = "str")]
11use crate::records::StringRecord;
12
13pub struct ReaderBuilder {
15 delimiter: u8,
16 quote: u8,
17 buffer_capacity: usize,
18 flexible: bool,
19 has_headers: bool,
20}
21
22impl Default for ReaderBuilder {
23 fn default() -> Self {
24 Self {
25 delimiter: b',',
26 quote: b'"',
27 buffer_capacity: 8192,
28 flexible: false,
29 has_headers: true,
30 }
31 }
32}
33
34impl ReaderBuilder {
35 pub fn new() -> Self {
37 Self::default()
38 }
39
40 pub fn with_capacity(capacity: usize) -> Self {
42 let mut reader = Self::default();
43 reader.buffer_capacity(capacity);
44 reader
45 }
46
47 pub fn delimiter(&mut self, delimiter: u8) -> &mut Self {
53 self.delimiter = delimiter;
54 self
55 }
56
57 pub fn quote(&mut self, quote: u8) -> &mut Self {
63 self.quote = quote;
64 self
65 }
66
67 pub fn buffer_capacity(&mut self, capacity: usize) -> &mut Self {
69 self.buffer_capacity = capacity;
70 self
71 }
72
73 pub fn flexible(&mut self, yes: bool) -> &mut Self {
79 self.flexible = yes;
80 self
81 }
82
83 pub fn has_headers(&mut self, yes: bool) -> &mut Self {
87 self.has_headers = yes;
88 self
89 }
90
91 pub fn from_reader<R: Read>(&self, reader: R) -> Reader<R> {
94 Reader {
95 buffer: BufReaderWithPosition::with_capacity(self.buffer_capacity, reader),
96 inner: CoreReader::new(self.delimiter, self.quote),
97 flexible: self.flexible,
98 headers: ByteRecord::new(),
99 has_read: false,
100 must_reemit_headers: !self.has_headers,
101 has_headers: self.has_headers,
102 index: 0,
103 }
104 }
105
106 pub fn reverse_from_reader<R: Read + Seek>(
109 &self,
110 mut reader: R,
111 ) -> error::Result<ReverseReader<R>> {
112 let initial_pos = reader.stream_position()?;
113
114 let mut forward_reader = self.from_reader(reader);
115 let headers = forward_reader.byte_headers()?.clone();
116 let position_after_headers = forward_reader.position();
117
118 let mut reader = forward_reader.into_inner();
119
120 let file_len = reader.seek(SeekFrom::End(0))?;
121
122 let offset = if self.has_headers {
123 initial_pos + position_after_headers
124 } else {
125 initial_pos
126 };
127
128 let reverse_io_reader = utils::ReverseReader::new(reader, file_len, offset);
129
130 Ok(ReverseReader {
131 buffer: BufReader::with_capacity(self.buffer_capacity, reverse_io_reader),
132 inner: CoreReader::new(self.delimiter, self.quote),
133 flexible: self.flexible,
134 headers,
135 })
136 }
137}
138
139pub struct Reader<R> {
147 buffer: BufReaderWithPosition<R>,
148 inner: CoreReader,
149 flexible: bool,
150 headers: ByteRecord,
151 has_read: bool,
152 must_reemit_headers: bool,
153 has_headers: bool,
154 index: u64,
155}
156
157impl<R: Read> Reader<R> {
158 pub fn from_reader(reader: R) -> Self {
161 ReaderBuilder::new().from_reader(reader)
162 }
163
164 #[inline]
165 fn check_field_count(&mut self, byte: u64, written: usize) -> error::Result<()> {
166 if self.flexible {
167 return Ok(());
168 }
169
170 if self.has_read && written != self.headers.len() {
171 return Err(Error::new(ErrorKind::UnequalLengths {
172 expected_len: self.headers.len(),
173 len: written,
174 pos: Some((
175 byte,
176 self.index
177 .saturating_sub(if self.has_headers { 1 } else { 0 }),
178 )),
179 }));
180 }
181
182 Ok(())
183 }
184
185 fn read_byte_record_impl(&mut self, record: &mut ByteRecord) -> error::Result<bool> {
186 use ReadResult::*;
187
188 record.clear();
189
190 let mut record_builder = ByteRecordBuilder::wrap(record);
191 let byte = self.position();
192
193 loop {
194 let input = self.buffer.fill_buf()?;
195
196 let (result, pos) = self.inner.read_record(input, &mut record_builder);
197
198 self.buffer.consume(pos);
199
200 match result {
201 End => {
202 return Ok(false);
203 }
204 Cr | Lf | InputEmpty => {
205 continue;
206 }
207 Record => {
208 self.index += 1;
209 self.check_field_count(byte, record.len())?;
210 return Ok(true);
211 }
212 };
213 }
214 }
215
216 #[inline]
217 fn on_first_read(&mut self) -> error::Result<()> {
218 if self.has_read {
219 return Ok(());
220 }
221
222 let input = self.buffer.fill_buf()?;
224 let bom_len = trim_bom(input);
225 self.buffer.consume(bom_len);
226
227 let mut headers = ByteRecord::new();
229
230 let has_data = self.read_byte_record_impl(&mut headers)?;
231
232 if !has_data {
233 self.must_reemit_headers = false;
234 }
235
236 self.headers = headers;
237 self.has_read = true;
238
239 Ok(())
240 }
241
242 #[inline]
245 pub fn has_headers(&self) -> bool {
246 self.has_headers
247 }
248
249 #[inline]
251 pub fn byte_headers(&mut self) -> error::Result<&ByteRecord> {
252 self.on_first_read()?;
253
254 Ok(&self.headers)
255 }
256
257 pub fn select(&mut self, selector: &Selector) -> error::Result<Selection> {
259 let has_headers = self.has_headers;
260 let headers = self.byte_headers()?;
261 selector.select(headers, has_headers)
262 }
263
264 pub fn select_one(&mut self, selector: &Selector) -> error::Result<usize> {
266 let has_headers = self.has_headers;
267 let headers = self.byte_headers()?;
268 selector.select_one(headers, has_headers)
269 }
270
271 #[inline(always)]
276 pub fn read_byte_record(&mut self, record: &mut ByteRecord) -> error::Result<bool> {
277 self.on_first_read()?;
278
279 if self.must_reemit_headers {
280 self.headers.clone_into(record);
281 self.must_reemit_headers = false;
282 return Ok(true);
283 }
284
285 self.read_byte_record_impl(record)
286 }
287
288 #[cfg(feature = "str")]
289 pub fn read_record(&mut self, record: &mut StringRecord) -> error::Result<bool> {
290 if self.read_byte_record(record.as_inner_mut())? {
291 if !record.validate_utf8() {
292 Err(Error::new(ErrorKind::Utf8Error))
293 } else {
294 Ok(true)
295 }
296 } else {
297 Ok(false)
298 }
299 }
300
301 pub fn byte_records(&mut self) -> ByteRecordsIter<'_, R> {
303 ByteRecordsIter {
304 reader: self,
305 record: ByteRecord::new(),
306 }
307 }
308
309 pub fn into_byte_records(self) -> ByteRecordsIntoIter<R> {
311 ByteRecordsIntoIter {
312 reader: self,
313 record: ByteRecord::new(),
314 }
315 }
316
317 #[cfg(feature = "str")]
319 pub fn records(&mut self) -> StringRecordsIter<'_, R> {
320 StringRecordsIter {
321 reader: self,
322 record: StringRecord::new(),
323 }
324 }
325
326 #[cfg(feature = "str")]
328 pub fn into_records(self) -> StringRecordsIntoIter<R> {
329 StringRecordsIntoIter {
330 reader: self,
331 record: StringRecord::new(),
332 }
333 }
334
335 pub fn get_ref(&self) -> &R {
337 self.buffer.get_ref()
338 }
339
340 pub fn get_mut(&mut self) -> &mut R {
342 self.buffer.get_mut()
343 }
344
345 pub fn into_inner(self) -> R {
349 self.buffer.into_inner().into_inner()
350 }
351
352 #[inline(always)]
354 pub fn position(&self) -> u64 {
355 if self.must_reemit_headers {
356 0
357 } else {
358 self.buffer.position()
359 }
360 }
361}
362
363pub struct ByteRecordsIter<'r, R> {
364 reader: &'r mut Reader<R>,
365 record: ByteRecord,
366}
367
368impl<R: Read> Iterator for ByteRecordsIter<'_, R> {
369 type Item = error::Result<ByteRecord>;
370
371 #[inline]
372 fn next(&mut self) -> Option<Self::Item> {
373 match self.reader.read_byte_record(&mut self.record) {
376 Err(err) => Some(Err(err)),
377 Ok(true) => Some(Ok(self.record.clone())),
378 Ok(false) => None,
379 }
380 }
381}
382
383pub struct ByteRecordsIntoIter<R> {
384 reader: Reader<R>,
385 record: ByteRecord,
386}
387
388impl<R: Read> Iterator for ByteRecordsIntoIter<R> {
389 type Item = error::Result<ByteRecord>;
390
391 #[inline]
392 fn next(&mut self) -> Option<Self::Item> {
393 match self.reader.read_byte_record(&mut self.record) {
396 Err(err) => Some(Err(err)),
397 Ok(true) => Some(Ok(self.record.clone())),
398 Ok(false) => None,
399 }
400 }
401}
402
403#[cfg(feature = "str")]
404pub struct StringRecordsIter<'r, R> {
405 reader: &'r mut Reader<R>,
406 record: StringRecord,
407}
408
409#[cfg(feature = "str")]
410impl<R: Read> Iterator for StringRecordsIter<'_, R> {
411 type Item = error::Result<StringRecord>;
412
413 #[inline]
414 fn next(&mut self) -> Option<Self::Item> {
415 match self.reader.read_record(&mut self.record) {
418 Err(err) => Some(Err(err)),
419 Ok(true) => Some(Ok(self.record.clone())),
420 Ok(false) => None,
421 }
422 }
423}
424
425#[cfg(feature = "str")]
426pub struct StringRecordsIntoIter<R> {
427 reader: Reader<R>,
428 record: StringRecord,
429}
430
431#[cfg(feature = "str")]
432impl<R: Read> Iterator for StringRecordsIntoIter<R> {
433 type Item = error::Result<StringRecord>;
434
435 #[inline]
436 fn next(&mut self) -> Option<Self::Item> {
437 match self.reader.read_record(&mut self.record) {
440 Err(err) => Some(Err(err)),
441 Ok(true) => Some(Ok(self.record.clone())),
442 Ok(false) => None,
443 }
444 }
445}
446
447pub struct ReverseReader<R> {
455 inner: CoreReader,
456 buffer: BufReader<utils::ReverseReader<R>>,
457 flexible: bool,
458 headers: ByteRecord,
459}
460
461impl<R: Read + Seek> ReverseReader<R> {
462 pub fn from_reader(reader: R) -> error::Result<Self> {
466 ReaderBuilder::new().reverse_from_reader(reader)
467 }
468
469 pub fn byte_headers(&self) -> &ByteRecord {
471 &self.headers
472 }
473
474 #[inline]
475 fn check_field_count(&mut self, written: usize) -> error::Result<()> {
476 if self.flexible {
477 return Ok(());
478 }
479
480 if written != self.headers.len() {
481 return Err(Error::new(ErrorKind::UnequalLengths {
482 expected_len: self.headers.len(),
483 len: written,
484 pos: None,
485 }));
486 }
487
488 Ok(())
489 }
490
491 pub fn read_byte_record(&mut self, record: &mut ByteRecord) -> error::Result<bool> {
496 use ReadResult::*;
497
498 record.clear();
499
500 let mut record_builder = ByteRecordBuilder::wrap(record);
501
502 loop {
503 let input = self.buffer.fill_buf()?;
504
505 let (result, pos) = self.inner.read_record(input, &mut record_builder);
506
507 self.buffer.consume(pos);
508
509 match result {
510 End => {
511 return Ok(false);
512 }
513 Cr | Lf | InputEmpty => {
514 continue;
515 }
516 Record => {
517 self.check_field_count(record.len())?;
518 record.reverse();
519 return Ok(true);
520 }
521 };
522 }
523 }
524
525 pub fn byte_records(&mut self) -> ReverseByteRecordsIter<'_, R> {
527 ReverseByteRecordsIter {
528 reader: self,
529 record: ByteRecord::new(),
530 }
531 }
532
533 pub fn into_byte_records(self) -> ReverseByteRecordsIntoIter<R> {
535 ReverseByteRecordsIntoIter {
536 reader: self,
537 record: ByteRecord::new(),
538 }
539 }
540}
541
542pub struct ReverseByteRecordsIter<'r, R> {
543 reader: &'r mut ReverseReader<R>,
544 record: ByteRecord,
545}
546
547impl<R: Read + Seek> Iterator for ReverseByteRecordsIter<'_, R> {
548 type Item = error::Result<ByteRecord>;
549
550 #[inline]
551 fn next(&mut self) -> Option<Self::Item> {
552 match self.reader.read_byte_record(&mut self.record) {
555 Err(err) => Some(Err(err)),
556 Ok(true) => Some(Ok(self.record.clone())),
557 Ok(false) => None,
558 }
559 }
560}
561
562pub struct ReverseByteRecordsIntoIter<R> {
563 reader: ReverseReader<R>,
564 record: ByteRecord,
565}
566
567impl<R: Read + Seek> Iterator for ReverseByteRecordsIntoIter<R> {
568 type Item = error::Result<ByteRecord>;
569
570 #[inline]
571 fn next(&mut self) -> Option<Self::Item> {
572 match self.reader.read_byte_record(&mut self.record) {
575 Err(err) => Some(Err(err)),
576 Ok(true) => Some(Ok(self.record.clone())),
577 Ok(false) => None,
578 }
579 }
580}
581
582#[cfg(test)]
583mod tests {
584 use std::io::Cursor;
585
586 use super::*;
587
588 impl<R: Read> Reader<R> {
589 fn from_reader_no_headers(reader: R) -> Self {
590 ReaderBuilder::new().has_headers(false).from_reader(reader)
591 }
592 }
593
594 #[test]
595 fn test_read_byte_record() -> error::Result<()> {
596 let csv = "name,surname,age\n\"john\",\"landy, the \"\"everlasting\"\" bastard\",45\n\"\"\"ok\"\"\",whatever,dude\nlucy,rose,\"67\"\njermaine,jackson,\"89\"\n\nkarine,loucan,\"52\"\nrose,\"glib\",12\n\"guillaume\",\"plique\",\"42\"\r\n";
597
598 let expected = vec![
599 brec!["name", "surname", "age"],
600 brec!["john", "landy, the \"everlasting\" bastard", "45"],
601 brec!["\"ok\"", "whatever", "dude"],
602 brec!["lucy", "rose", "67"],
603 brec!["jermaine", "jackson", "89"],
604 brec!["karine", "loucan", "52"],
605 brec!["rose", "glib", "12"],
606 brec!["guillaume", "plique", "42"],
607 ];
608
609 for capacity in [32usize, 4, 3, 2, 1] {
610 let mut reader = ReaderBuilder::with_capacity(capacity)
611 .has_headers(false)
612 .from_reader(Cursor::new(csv));
613
614 assert_eq!(
615 reader.byte_records().collect::<Result<Vec<_>, _>>()?,
616 expected,
617 );
618 }
619
620 Ok(())
621 }
622
623 #[test]
624 #[cfg(feature = "str")]
625 fn test_read_record() -> error::Result<()> {
626 let csv =
627 "french,chinese\nReine-Mère de l'Ouest,西王母\nEmpereur du Pic de l'Est,东华帝君\r\n";
628
629 let expected = vec![
630 srec!["french", "chinese"],
631 srec!["Reine-Mère de l'Ouest", "西王母"],
632 srec!["Empereur du Pic de l'Est", "东华帝君"],
633 ];
634
635 for capacity in [32usize, 4, 3, 2, 1] {
636 let mut reader = ReaderBuilder::with_capacity(capacity)
637 .has_headers(false)
638 .from_reader(Cursor::new(csv));
639
640 assert_eq!(reader.records().collect::<Result<Vec<_>, _>>()?, expected,);
641 }
642
643 Ok(())
644 }
645
646 #[test]
647 fn test_strip_bom() -> error::Result<()> {
648 let mut reader = Reader::from_reader_no_headers(Cursor::new("name,surname,age"));
649
650 assert_eq!(
651 reader.byte_records().next().unwrap()?,
652 brec!["name", "surname", "age"]
653 );
654
655 let mut reader =
656 Reader::from_reader_no_headers(Cursor::new(b"\xef\xbb\xbfname,surname,age"));
657
658 assert_eq!(
659 reader.byte_records().next().unwrap()?,
660 brec!["name", "surname", "age"]
661 );
662
663 Ok(())
664 }
665
666 #[test]
667 fn test_empty_row() -> error::Result<()> {
668 let data = "name\n\"\"\nlucy\n\"\"";
669
670 let reader = Reader::from_reader_no_headers(Cursor::new(data));
672
673 let expected = vec![brec!["name"], brec![""], brec!["lucy"], brec![""]];
674
675 let records = reader.into_byte_records().collect::<Result<Vec<_>, _>>()?;
676
677 assert_eq!(records, expected);
678
679 Ok(())
680 }
681
682 #[test]
683 fn test_crlf() -> error::Result<()> {
684 let reader = Reader::from_reader_no_headers(Cursor::new(
685 "name,surname\r\nlucy,\"john\"\r\nevan,zhong\r\nbéatrice,glougou\r\n",
686 ));
687
688 let expected = vec![
689 brec!["name", "surname"],
690 brec!["lucy", "john"],
691 brec!["evan", "zhong"],
692 brec!["béatrice", "glougou"],
693 ];
694
695 let records = reader.into_byte_records().collect::<Result<Vec<_>, _>>()?;
696
697 assert_eq!(records, expected);
698
699 Ok(())
700 }
701
702 #[test]
703 fn test_quote_always() -> error::Result<()> {
704 let reader = Reader::from_reader_no_headers(Cursor::new(
705 "\"name\",\"surname\"\n\"lucy\",\"rose\"\n\"john\",\"mayhew\"",
706 ));
707
708 let expected = vec![
709 brec!["name", "surname"],
710 brec!["lucy", "rose"],
711 brec!["john", "mayhew"],
712 ];
713
714 let records = reader.into_byte_records().collect::<Result<Vec<_>, _>>()?;
715
716 assert_eq!(records, expected);
717
718 Ok(())
719 }
720
721 #[test]
722 fn test_byte_headers() -> error::Result<()> {
723 let data = b"name,surname\njohn,dandy";
724
725 let mut reader = Reader::from_reader(Cursor::new(data));
727 assert_eq!(reader.byte_headers()?, &brec!["name", "surname"]);
728 assert_eq!(
729 reader.byte_records().next().unwrap()?,
730 brec!["john", "dandy"]
731 );
732
733 let mut reader = Reader::from_reader(Cursor::new(data));
735 assert_eq!(
736 reader.byte_records().next().unwrap()?,
737 brec!["john", "dandy"]
738 );
739 assert_eq!(reader.byte_headers()?, &brec!["name", "surname"]);
740
741 let mut reader = Reader::from_reader_no_headers(Cursor::new(data));
743 assert_eq!(reader.byte_headers()?, &brec!["name", "surname"]);
744 assert_eq!(
745 reader.byte_records().next().unwrap()?,
746 brec!["name", "surname"]
747 );
748
749 let mut reader = Reader::from_reader_no_headers(Cursor::new(data));
751 assert_eq!(
752 reader.byte_records().next().unwrap()?,
753 brec!["name", "surname"]
754 );
755 assert_eq!(reader.byte_headers()?, &brec!["name", "surname"]);
756
757 let mut reader = Reader::from_reader(Cursor::new(b""));
759 assert_eq!(reader.byte_headers()?, &brec![]);
760 assert!(reader.byte_records().next().is_none());
761
762 let mut reader = Reader::from_reader_no_headers(Cursor::new(b""));
764 assert_eq!(reader.byte_headers()?, &brec![]);
765 assert!(reader.byte_records().next().is_none());
766
767 Ok(())
768 }
769
770 #[test]
771 fn test_weirdness() -> error::Result<()> {
772 let data =
774 b"name,surname\n\"test\" \"wat\", ok\ntest \"wat\",ok \ntest,\"whatever\" ok\n\"test\" there,\"ok\"\r\n";
775 let mut reader = Reader::from_reader_no_headers(Cursor::new(data));
776
777 let records = reader.byte_records().collect::<Result<Vec<_>, _>>()?;
778
779 let expected = vec![
780 brec!["name", "surname"],
781 brec!["test \"wat", " ok"],
782 brec!["test \"wat", "ok "],
783 brec!["test", "whatever ok"],
784 brec!["test there", "ok"],
785 ];
786
787 assert_eq!(records, expected);
788
789 let data = b"name,surname\n\r\rjohn,coucou";
796 let mut reader = Reader::from_reader_no_headers(Cursor::new(data));
797 let records = reader.byte_records().collect::<Result<Vec<_>, _>>()?;
798
799 assert_eq!(
800 records,
801 vec![brec!["name", "surname"], brec!["john", "coucou"]]
802 );
803
804 Ok(())
805 }
806
807 #[test]
808 fn test_position() -> error::Result<()> {
809 let data = b"name,surname\njohnny,landis crue\nbabka,bob caterpillar\n";
810
811 let mut reader = Reader::from_reader(&data[..]);
812 let mut record = ByteRecord::new();
813
814 let mut positions = vec![reader.position()];
815
816 reader.byte_headers()?;
817
818 positions.push(reader.position());
819
820 while reader.read_byte_record(&mut record)? {
821 positions.push(reader.position());
822 }
823
824 assert_eq!(positions, vec![0, 13, 32, 54]);
825
826 let mut reader = ReaderBuilder::new()
827 .has_headers(false)
828 .from_reader(&data[..]);
829
830 reader.byte_headers()?;
831
832 assert_eq!(reader.position(), 0);
833
834 Ok(())
835 }
836
837 #[test]
838 fn test_reverse_reader() -> error::Result<()> {
839 let data = b"name,surname\njohn,landis\nbeatrice,babka\nevan,michalak";
840 let mut reader = ReverseReader::from_reader(Cursor::new(data))?;
841
842 assert_eq!(
843 reader.byte_records().collect::<Result<Vec<_>, _>>()?,
844 vec![
845 brec!["evan", "michalak"],
846 brec!["beatrice", "babka"],
847 brec!["john", "landis"]
848 ]
849 );
850
851 assert_eq!(reader.byte_headers(), &brec!["name", "surname"]);
852
853 Ok(())
854 }
855
856 #[test]
857 fn test_reverse_reader_crlf() -> error::Result<()> {
858 let data = b"name,surname\r\njohn,landis\r\nbeatrice,babka\r\nevan,michalak";
859 let mut reader = ReverseReader::from_reader(Cursor::new(data))?;
860
861 assert_eq!(
862 reader.byte_records().collect::<Result<Vec<_>, _>>()?,
863 vec![
864 brec!["evan", "michalak"],
865 brec!["beatrice", "babka"],
866 brec!["john", "landis"]
867 ]
868 );
869
870 assert_eq!(reader.byte_headers(), &brec!["name", "surname"]);
871
872 Ok(())
873 }
874
875 #[test]
876 fn test_weird_sequence() -> error::Result<()> {
877 let data = b"\r\r`\"\",\n,`\"\r\",\n";
878 let mut record = ByteRecord::new();
879 let mut reader = ReaderBuilder::new()
880 .flexible(true)
881 .has_headers(false)
882 .from_reader(&data[..]);
883
884 reader.read_byte_record(&mut record)?;
885 assert_eq!(record, brec!["`\"", ""]);
886
887 reader.read_byte_record(&mut record)?;
888
889 assert_eq!(record, brec!["", "\"\r", ""]);
890
891 Ok(())
892 }
893
894 #[test]
895 fn test_quoted_final_cr() -> error::Result<()> {
896 let csv = b"name,surname\n\"test\",\"\r\"\njohn,landis";
897
898 let expected = vec![
899 brec!["name", "surname"],
900 brec!["test", "\r"],
901 brec!["john", "landis"],
902 ];
903
904 for capacity in [32usize, 4, 3, 2, 1] {
905 let mut reader = ReaderBuilder::with_capacity(capacity)
906 .has_headers(false)
907 .from_reader(Cursor::new(csv));
908
909 assert_eq!(
910 reader.byte_records().collect::<Result<Vec<_>, _>>()?,
911 expected,
912 );
913 }
914
915 Ok(())
916 }
917
918 }