1use crate::row::Row;
6use fv_value::{compile, ExprError, Value};
7
8#[derive(Debug, Clone)]
10pub enum Step {
11 Select { columns: Vec<String> },
12 Rename { mapping: Vec<(String, String)> },
13 Drop { columns: Vec<String> },
14 Filter { expression: String },
15 ApplyExpression { column: String, expression: String },
16}
17
18#[derive(Debug)]
19pub enum StepError {
20 Expr(ExprError),
21 Message(String),
22}
23
24impl std::fmt::Display for StepError {
25 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
26 match self {
27 StepError::Expr(e) => write!(f, "{e}"),
28 StepError::Message(m) => write!(f, "{m}"),
29 }
30 }
31}
32
33impl From<ExprError> for StepError {
34 fn from(e: ExprError) -> Self {
35 StepError::Expr(e)
36 }
37}
38
39pub fn apply_step(step: &Step, rows: &[Row]) -> Result<Vec<Row>, StepError> {
41 match step {
42 Step::Select { columns } => Ok(rows
44 .iter()
45 .map(|r| Row(columns.iter().map(|c| (c.clone(), r.get(c))).collect()))
46 .collect()),
47
48 Step::Rename { mapping } => Ok(rows
50 .iter()
51 .map(|r| {
52 Row(r
53 .0
54 .iter()
55 .map(|(k, v)| {
56 let nk = mapping
57 .iter()
58 .find(|(from, _)| from == k)
59 .map(|(_, to)| to.clone())
60 .unwrap_or_else(|| k.clone());
61 (nk, v.clone())
62 })
63 .collect())
64 })
65 .collect()),
66
67 Step::Drop { columns } => Ok(rows
69 .iter()
70 .map(|r| Row(r.0.iter().filter(|(k, _)| !columns.contains(k)).cloned().collect()))
71 .collect()),
72
73 Step::Filter { expression } => {
75 let ast = compile(expression)?;
76 let mut out = Vec::new();
77 for r in rows {
78 match ast.eval(r.lookup())? {
79 Value::Bool(true) => out.push(r.clone()),
80 Value::Bool(false) => {}
81 other => {
82 return Err(StepError::Message(format!(
83 "filter: predicate must yield a boolean, got {other:?}"
84 )))
85 }
86 }
87 }
88 Ok(out)
89 }
90
91 Step::ApplyExpression { column, expression } => {
93 let ast = compile(expression)?;
94 let mut out = Vec::with_capacity(rows.len());
95 for r in rows {
96 let v = ast.eval(r.lookup())?;
97 let mut nr = r.clone();
98 nr.set(column, v);
99 out.push(nr);
100 }
101 Ok(out)
102 }
103 }
104}
105
106pub fn apply_steps(steps: &[Step], rows: &[Row]) -> Result<Vec<Row>, StepError> {
108 let mut current = rows.to_vec();
109 for step in steps {
110 current = apply_step(step, ¤t)?;
111 }
112 Ok(current)
113}
114
115pub fn apply_steps_isolating(steps: &[Step], rows: &[Row]) -> (Vec<Row>, usize) {
124 match apply_steps(steps, rows) {
125 Ok(out) => (out, 0),
126 Err(_) => {
127 let mut survived = Vec::new();
128 let mut dropped = 0usize;
129 for row in rows {
130 match apply_steps(steps, std::slice::from_ref(row)) {
131 Ok(mut o) => survived.append(&mut o),
132 Err(_) => dropped += 1,
133 }
134 }
135 (survived, dropped)
136 }
137 }
138}
139
140pub fn validate_steps(steps: &[Step]) -> Result<(), StepError> {
144 apply_steps(steps, &[]).map(|_| ())
145}
146
147pub const INLINE_OPS: &[&str] = &["select", "rename", "drop", "filter", "applyExpression"];
150
151pub fn parse(step: &serde_json::Value) -> Result<Step, String> {
153 let strs = |k: &str| {
154 step[k]
155 .as_array()
156 .map(|a| a.iter().filter_map(|v| v.as_str().map(String::from)).collect())
157 .unwrap_or_default()
158 };
159 Ok(match step["op"].as_str().unwrap_or("") {
160 "select" => Step::Select {
161 columns: strs("columns"),
162 },
163 "drop" => Step::Drop {
164 columns: strs("columns"),
165 },
166 "rename" => Step::Rename {
167 mapping: step["mapping"]
168 .as_object()
169 .map(|m| {
170 m.iter()
171 .map(|(k, v)| (k.clone(), v.as_str().unwrap_or("").to_string()))
172 .collect()
173 })
174 .unwrap_or_default(),
175 },
176 "filter" => Step::Filter {
177 expression: step["expression"].as_str().unwrap_or("").to_string(),
178 },
179 "applyExpression" => Step::ApplyExpression {
180 column: step["column"].as_str().unwrap_or("").to_string(),
181 expression: step["expression"].as_str().unwrap_or("").to_string(),
182 },
183 other => return Err(format!("unknown inline step op '{other}'")),
184 })
185}
186
187#[cfg(test)]
188mod isolation_tests {
189 use super::*;
190 use crate::row::Row;
191 use fv_value::Value;
192
193 fn row(pairs: &[(&str, Value)]) -> Row {
194 Row(pairs.iter().map(|(k, v)| (k.to_string(), v.clone())).collect())
195 }
196
197 #[test]
198 fn isolating_drops_only_the_poison_row() {
199 let steps = vec![Step::Filter {
202 expression: "keep".into(),
203 }];
204 let rows = vec![
205 row(&[("keep", Value::Bool(true)), ("id", Value::Num(1.0))]),
206 row(&[("keep", Value::Num(9.0)), ("id", Value::Num(2.0))]), row(&[("keep", Value::Bool(true)), ("id", Value::Num(3.0))]),
208 ];
209 assert!(apply_steps(&steps, &rows).is_err(), "whole-batch apply is fail-closed");
210 let (out, dropped) = apply_steps_isolating(&steps, &rows);
211 assert_eq!(dropped, 1);
212 assert_eq!(out.len(), 2);
213 assert_eq!(out[0].get("id"), Value::Num(1.0));
214 assert_eq!(out[1].get("id"), Value::Num(3.0));
215 }
216
217 #[test]
218 fn isolating_fast_path_when_all_rows_ok() {
219 let steps = vec![Step::Filter {
220 expression: "keep".into(),
221 }];
222 let rows = vec![
223 row(&[("keep", Value::Bool(true)), ("id", Value::Num(1.0))]),
224 row(&[("keep", Value::Bool(false)), ("id", Value::Num(2.0))]),
225 ];
226 let (out, dropped) = apply_steps_isolating(&steps, &rows);
227 assert_eq!(dropped, 0);
228 assert_eq!(out.len(), 1, "filter keeps only the true row");
229 assert_eq!(out[0].get("id"), Value::Num(1.0));
230 }
231
232 #[test]
233 fn validate_rejects_bad_expression_accepts_good() {
234 assert!(validate_steps(&[Step::ApplyExpression {
235 column: "x".into(),
236 expression: "1 +".into()
237 }])
238 .is_err());
239 assert!(validate_steps(&[Step::ApplyExpression {
240 column: "x".into(),
241 expression: "1 + 2".into()
242 }])
243 .is_ok());
244 }
245}