diff --git a/src/pipeline.rs b/src/pipeline.rs index 318f87c..19b785b 100644 --- a/src/pipeline.rs +++ b/src/pipeline.rs @@ -1,5 +1,5 @@ use super::headers::Headers; -use crate::pipeline_iterators::{AddCol, Flush, MapCol, MapRow, TransformInto}; +use crate::pipeline_iterators::{AddCol, Flush, MapCol, MapRow, TransformInto, Validate}; use crate::target::Target; use crate::transform::Transform; use crate::{Error, Row, RowResult, StringTarget}; @@ -225,6 +225,18 @@ impl<'a> Pipeline<'a> { } } + pub fn validate(mut self, get_row: F) -> Self + where + F: FnMut(&Headers, &Row) -> Result<(), Error> + 'a, + { + self.iterator = Box::new(Validate { + iterator: self.iterator, + f: get_row, + headers: self.headers.clone(), + }); + self + } + /// Write to the specified [`Target`]. /// /// ## Example diff --git a/src/pipeline_iterators.rs b/src/pipeline_iterators.rs index 7a18057..3354aa2 100644 --- a/src/pipeline_iterators.rs +++ b/src/pipeline_iterators.rs @@ -152,6 +152,30 @@ where } } +pub struct Validate { + pub iterator: I, + pub f: F, + pub headers: Headers, +} +impl Iterator for Validate +where + I: Iterator, + F: FnMut(&Headers, &Row) -> Result<(), Error>, +{ + type Item = RowResult; + + fn next(&mut self) -> Option { + let row = match self.iterator.next()? { + Ok(row) => row, + Err(e) => return Some(Err(e)), + }; + match (self.f)(&self.headers, &row) { + Ok(()) => Some(Ok(row)), + Err(e) => Some(Err(e)), + } + } +} + pub struct Flush { pub iterator: I, pub target: T,