[READ-ONLY] Mirror of https://github.com/probablykasper/csv-pipeline.
0

Configure Feed

Select the types of activity you want to include in your feed.

Rename Pipeline to PipelineBuilder, Pipe to Chain

Kasper (Jan 1, 2023, 12:37 PM +0100) 213d57cf 5db82c38

+144 -87
+72
src/chain.rs
··· 1 + use crate::RowResult; 2 + pub type BoxedIterator = Box<dyn Iterator<Item = RowResult>>; 3 + 4 + /// A struct that wraps a RowResult iterator for convenience 5 + pub struct Chain { 6 + iterator: BoxedIterator, 7 + } 8 + impl Chain { 9 + pub fn new(iterator: BoxedIterator) -> Self { 10 + Self { 11 + iterator: Box::new(iterator), 12 + } 13 + } 14 + pub fn with_state<S>(self, state: S) -> StatefulChainBuilder<S> { 15 + StatefulChainBuilder::new(Box::new(self), state) 16 + } 17 + } 18 + impl Iterator for Chain { 19 + type Item = RowResult; 20 + 21 + fn next(&mut self) -> Option<Self::Item> { 22 + self.iterator.next() 23 + } 24 + } 25 + 26 + pub struct StatefulChainBuilder<S> { 27 + iterator: BoxedIterator, 28 + state: S, 29 + } 30 + impl<S> StatefulChainBuilder<S> { 31 + pub fn new(iterator: BoxedIterator, state: S) -> Self { 32 + Self { state, iterator } 33 + } 34 + pub fn map<F>(self, f: F) -> Chain 35 + where 36 + F: FnMut(RowResult, &mut S) -> RowResult, 37 + { 38 + let x = StatefulChain { 39 + iterator: self.iterator, 40 + state: self.state, 41 + f, 42 + }; 43 + x.into_chain() 44 + } 45 + } 46 + 47 + pub struct StatefulChain<S, F: FnMut(RowResult, &mut S) -> RowResult> { 48 + iterator: BoxedIterator, 49 + state: S, 50 + f: F, 51 + } 52 + impl<S, F> StatefulChain<S, F> 53 + where 54 + F: FnMut(RowResult, &mut S) -> RowResult, 55 + { 56 + pub fn into_chain(self) -> Chain { 57 + Chain::new(Box::new(self)) 58 + } 59 + } 60 + impl<S, F> Iterator for StatefulChain<S, F> 61 + where 62 + F: FnMut(RowResult, &mut S) -> RowResult, 63 + { 64 + type Item = RowResult; 65 + 66 + fn next(&mut self) -> Option<Self::Item> { 67 + match self.iterator.next() { 68 + Some(item) => Some((self.f)(item, &mut self.state)), 69 + None => None, 70 + } 71 + } 72 + }
+23 -5
src/headers.rs
··· 1 1 use crate::Row; 2 + use csv::StringRecordIter; 2 3 use std::collections::HashMap; 3 4 4 5 #[derive(Debug, Clone, PartialEq)] 5 6 pub struct Headers { 6 7 indexes: HashMap<String, usize>, 7 - names: Row, 8 + row: Row, 8 9 } 9 10 impl Headers { 10 - pub fn add(&mut self, name: &str) -> bool { 11 + pub fn push_field(&mut self, name: &str) -> bool { 11 12 if self.indexes.contains_key(name) { 12 13 return false; 13 14 } 14 15 15 - self.names.push_field(name); 16 - self.indexes.insert(name.to_string(), self.names.len() - 1); 16 + self.row.push_field(name); 17 + self.indexes.insert(name.to_string(), self.row.len() - 1); 17 18 18 19 true 19 20 } ··· 21 22 pub fn contains(&self, name: &str) -> bool { 22 23 self.indexes.contains_key(name) 23 24 } 25 + 26 + pub fn get_row(&self) -> &Row { 27 + &self.row 28 + } 24 29 } 25 30 impl From<Row> for Headers { 26 31 fn from(row: Row) -> Headers { ··· 30 35 .enumerate() 31 36 .map(|(index, entry)| (entry.to_string(), index)) 32 37 .collect(), 33 - names: row, 38 + row, 34 39 } 35 40 } 36 41 } 42 + impl<'a> IntoIterator for &'a Headers { 43 + type Item = &'a str; 44 + type IntoIter = StringRecordIter<'a>; 45 + 46 + fn into_iter(self) -> StringRecordIter<'a> { 47 + self.row.into_iter() 48 + } 49 + } 50 + impl From<Headers> for Row { 51 + fn from(headers: Headers) -> Row { 52 + headers.row 53 + } 54 + }
+4 -5
src/lib.rs
··· 1 - pub mod headers; 2 - pub mod pipe; 3 - pub mod pipeline; 1 + mod chain; 2 + mod headers; 3 + mod pipeline; 4 4 5 5 pub use headers::Headers; 6 - pub use pipe::Pipe; 7 - pub use pipeline::Pipeline; 6 + pub use pipeline::{Pipeline, PipelineBuilder}; 8 7 9 8 #[derive(Debug, Clone, PartialEq)] 10 9 pub enum Error {
-63
src/pipe.rs
··· 1 - use crate::RowResult; 2 - 3 - pub type PipeIterator = Box<dyn Iterator<Item = RowResult>>; 4 - 5 - pub struct Pipe { 6 - iterator: PipeIterator, 7 - } 8 - impl Pipe { 9 - pub fn new(iterator: PipeIterator) -> Self { 10 - Self { 11 - iterator: Box::new(iterator.into_iter()), 12 - } 13 - } 14 - pub fn with_state<S>(self, state: S) -> StatefulPipeBuilder<S> { 15 - StatefulPipeBuilder::new(self.iterator, state) 16 - } 17 - } 18 - impl Iterator for Pipe { 19 - type Item = RowResult; 20 - 21 - fn next(&mut self) -> Option<Self::Item> { 22 - self.iterator.next() 23 - } 24 - } 25 - 26 - pub struct StatefulPipeBuilder<S> { 27 - iterator: PipeIterator, 28 - state: S, 29 - } 30 - impl<S> StatefulPipeBuilder<S> { 31 - pub fn new(iterator: PipeIterator, state: S) -> Self { 32 - Self { state, iterator } 33 - } 34 - pub fn map<F>(self, f: F) -> StatefulPipe<S, F> 35 - where 36 - F: FnMut(RowResult, &mut S) -> RowResult, 37 - { 38 - StatefulPipe { 39 - iterator: self.iterator, 40 - state: self.state, 41 - f, 42 - } 43 - } 44 - } 45 - 46 - pub struct StatefulPipe<S, F: FnMut(RowResult, &mut S) -> RowResult> { 47 - pub(crate) iterator: PipeIterator, 48 - state: S, 49 - f: F, 50 - } 51 - impl<S, F> Iterator for StatefulPipe<S, F> 52 - where 53 - F: FnMut(RowResult, &mut S) -> RowResult, 54 - { 55 - type Item = RowResult; 56 - 57 - fn next(&mut self) -> Option<Self::Item> { 58 - match self.iterator.next() { 59 - Some(item) => Some((self.f)(item, &mut self.state)), 60 - None => None, 61 - } 62 - } 63 - }
+45 -14
src/pipeline.rs
··· 1 + use super::chain::{BoxedIterator, Chain}; 1 2 use super::headers::Headers; 2 - use super::pipe::{Pipe, PipeIterator}; 3 3 use crate::{Error, Row, RowResult}; 4 4 use csv::{Reader, ReaderBuilder}; 5 5 use std::fs::File; 6 6 use std::path::Path; 7 7 8 - pub struct Pipeline { 8 + pub struct PipelineBuilder { 9 9 pub headers: Headers, 10 - pipe: PipeIterator, 10 + chain: Chain, 11 11 } 12 12 13 - impl Pipeline { 13 + impl PipelineBuilder { 14 14 pub fn from_reader(mut reader: Reader<File>) -> Self { 15 15 let headers_row = reader.headers().unwrap().clone(); 16 16 let records = reader.into_records().map(|r| { ··· 22 22 }); 23 23 Self { 24 24 headers: Headers::from(headers_row), 25 - pipe: Box::new(records), 25 + chain: Chain::new(Box::new(records)), 26 26 } 27 27 } 28 28 ··· 46 46 /// ## Example 47 47 /// 48 48 /// ``` 49 - /// use csv_pipeline::Pipeline; 49 + /// use csv_pipeline::PipelineBuilder; 50 50 /// 51 - /// Pipeline::from_path("test/Countries.csv") 51 + /// PipelineBuilder::from_path("test/Countries.csv") 52 52 /// .add_col("Language", |headers, row| { 53 53 /// Ok("".to_string()) 54 54 /// }); ··· 57 57 where 58 58 F: FnMut(&Headers, &Row) -> Result<String, Error>, 59 59 { 60 - self.headers.add(name); 60 + self.headers.push_field(name); 61 61 62 62 struct State<F> { 63 63 get_value: F, 64 64 headers: Headers, 65 65 } 66 - let pipe = Pipe::new(self.pipe).with_state(State { 66 + let stateful_chain = self.chain.with_state(State { 67 67 get_value, 68 68 headers: self.headers.clone(), 69 69 }); 70 - let newpipe = pipe.map(|row_result, state| { 70 + let new_chain = stateful_chain.map(|row_result, state| { 71 + println!("ADDCOL-map"); 71 72 let mut row = row_result?; 72 73 let value = (state.get_value)(&state.headers, &row)?; 73 74 row.push_field(&value); 74 75 Ok(row) 75 76 }); 76 77 77 - self.pipe = Box::new(newpipe.iterator); 78 + self.chain = new_chain; 78 79 79 80 self 81 + } 82 + 83 + pub fn build(self) -> Pipeline { 84 + Pipeline { 85 + headers: self.headers, 86 + iterator: Box::new(self.chain), 87 + } 88 + } 89 + } 90 + 91 + pub struct Pipeline { 92 + pub headers: Headers, 93 + pub iterator: BoxedIterator, 94 + } 95 + impl Iterator for Pipeline { 96 + type Item = RowResult; 97 + 98 + fn next(&mut self) -> Option<Self::Item> { 99 + self.iterator.next() 80 100 } 81 101 } 82 102 83 103 #[cfg(test)] 84 104 mod tests { 85 - use crate::Pipeline; 105 + use crate::PipelineBuilder; 86 106 87 107 #[test] 88 108 fn add_col() { 89 - let mut pipeline = Pipeline::from_path("test/Countries.csv") 90 - .add_col("Language", |_headers, _row| Ok("".to_string())); 109 + let mut pipeline = PipelineBuilder::from_path("test/Countries.csv") 110 + .add_col("Language", |_headers, _row| Ok("".to_string())) 111 + .build(); 112 + 113 + let mut writer = csv::Writer::from_writer(vec![]); 114 + writer.write_record(&pipeline.headers).unwrap(); 115 + println!("{:?}", pipeline.headers.get_row()); 116 + while let Some(item) = pipeline.next() { 117 + println!("{:?}", item.clone().unwrap()); 118 + writer.write_record(&item.unwrap()).unwrap(); 119 + } 120 + let s = String::from_utf8(writer.into_inner().unwrap()).unwrap(); 121 + print!("{}", s); 91 122 } 92 123 }