Skip to main content

faucet_core/
zip_columns.rs

1//! Inbuilt `zip_columns` transform (#551): turn a **columnar payload** —
2//! `{ columns: [{name}, …], rows: [[v0, v1, …], …] }` — into one object per row,
3//! keyed by column name.
4//!
5//! Analytics / report APIs (e.g. a query-language `tableData` payload) return results
6//! positionally: a list of column descriptors plus a list of value-arrays. This
7//! transform zips each row against the column names so downstream stages and
8//! sinks see ordinary `{col: value}` records. It is expressible today via the
9//! DuckDB `sql` transform, but a small declarative transform is cleaner and
10//! needs no embedded engine.
11//!
12//! The whole module is gated by `#[cfg(feature = "transform-zip-columns")]` at
13//! the `mod` site in `lib.rs`. It routes through
14//! [`TransformStage::PageFn`](crate::stage::TransformStage) (page-level, 1→0..N,
15//! fallible) so a row whose width doesn't match the column count fails loudly
16//! rather than silently dropping or misaligning fields.
17
18use crate::FaucetError;
19use crate::stage::TransformStage;
20use crate::util::extract_records;
21use schemars::JsonSchema;
22use serde::{Deserialize, Serialize};
23use serde_json::{Map, Value};
24use std::sync::Arc;
25
26/// User-facing `zip_columns` config.
27#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
28#[serde(deny_unknown_fields)]
29pub struct ZipColumnsSpec {
30    /// JSONPath to the column **names**. Point it at the name of each column
31    /// descriptor (`columns[*].name`) or at a plain array of strings
32    /// (`columns`). Every matched value must be a string. Set exactly one of
33    /// `columns_path` or `groups`.
34    #[serde(default, skip_serializing_if = "String::is_empty")]
35    pub columns_path: String,
36    /// JSONPath to the **rows**. With `columns_path`, an array of positional
37    /// value-arrays (`rows` or `rows[*]`), each exactly as wide as the column
38    /// list. With `groups`, the row objects (`rows[*]`).
39    pub rows_path: String,
40    /// Column groups (#746): each row object holds several positional cell
41    /// arrays, each named by its own header list (a `runReport`-style
42    /// `dimensionValues` / `metricValues` pair). Every group is zipped on its own and
43    /// the results are merged into one record; two groups naming the same
44    /// column is an error.
45    #[serde(default, skip_serializing_if = "Vec::is_empty")]
46    pub groups: Vec<ColumnGroupSpec>,
47}
48
49/// One positional column group of a [`ZipColumnsSpec`].
50#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
51#[serde(deny_unknown_fields)]
52pub struct ColumnGroupSpec {
53    /// Field (dot path) of each row holding this group's cell array
54    /// (`dimensionValues`). A row without it fails the page.
55    pub from: String,
56    /// JSONPath, evaluated against the **record**, to this group's header list
57    /// (`$.dimensionHeaders[*].name`, or `$.dimensionHeaders[*]` with
58    /// `header_label`).
59    pub header: String,
60    /// Field of each header object to use as the column name, when `header`
61    /// matches objects rather than strings.
62    #[serde(default, skip_serializing_if = "Option::is_none")]
63    pub header_label: Option<String>,
64    /// Field (dot path) of each cell holding the value (`value`). Omit when
65    /// the cells are the values themselves. A cell without it yields `null`.
66    #[serde(default, skip_serializing_if = "Option::is_none")]
67    pub value: Option<String>,
68}
69
70impl ZipColumnsSpec {
71    /// Validate the spec, returning a reusable [`CompiledZipColumns`].
72    pub fn compile(&self) -> Result<CompiledZipColumns, FaucetError> {
73        CompiledZipColumns::compile(self)
74    }
75
76    /// Compile and wrap as a [`TransformStage::PageFn`] (1→0..N per record,
77    /// fallible on a row/column width mismatch).
78    pub fn into_stage(&self) -> Result<TransformStage, FaucetError> {
79        let compiled = self.compile()?;
80        Ok(TransformStage::PageFn(Arc::new(move |page: Vec<Value>| {
81            let mut out = Vec::with_capacity(page.len());
82            for rec in page {
83                out.extend(compiled.apply(&rec)?);
84            }
85            Ok(out)
86        })))
87    }
88}
89
90/// Validated [`ZipColumnsSpec`] — apply per record with [`CompiledZipColumns::apply`].
91#[derive(Debug, Clone)]
92pub struct CompiledZipColumns {
93    columns_path: String,
94    rows_path: String,
95    groups: Vec<CompiledGroup>,
96}
97
98#[derive(Debug, Clone)]
99struct CompiledGroup {
100    from: String,
101    header: String,
102    header_label: Option<String>,
103    value: Option<String>,
104}
105
106impl CompiledZipColumns {
107    fn compile(spec: &ZipColumnsSpec) -> Result<Self, FaucetError> {
108        let has_columns = !spec.columns_path.trim().is_empty();
109        if has_columns != spec.groups.is_empty() {
110            return Err(FaucetError::Config(
111                "zip_columns: set exactly one of `columns_path` or `groups`".into(),
112            ));
113        }
114        if spec.rows_path.trim().is_empty() {
115            return Err(FaucetError::Config(
116                "zip_columns: `rows_path` must not be empty".into(),
117            ));
118        }
119        let blank = |s: &str| s.trim().is_empty();
120        let mut groups = Vec::with_capacity(spec.groups.len());
121        for (i, g) in spec.groups.iter().enumerate() {
122            if blank(&g.from) || blank(&g.header) {
123                return Err(FaucetError::Config(format!(
124                    "zip_columns: group {i} needs a non-empty `from` and `header`"
125                )));
126            }
127            if g.header_label.as_deref().is_some_and(blank) || g.value.as_deref().is_some_and(blank)
128            {
129                return Err(FaucetError::Config(format!(
130                    "zip_columns: group '{}' has an empty `header_label` or `value`",
131                    g.from
132                )));
133            }
134            groups.push(CompiledGroup {
135                from: g.from.trim().to_string(),
136                header: normalize_path(&g.header),
137                header_label: g.header_label.clone(),
138                value: g.value.clone(),
139            });
140        }
141        Ok(Self {
142            columns_path: if has_columns {
143                normalize_path(&spec.columns_path)
144            } else {
145                String::new()
146            },
147            rows_path: normalize_path(&spec.rows_path),
148            groups,
149        })
150    }
151
152    /// Zip one columnar record into one object per row. A record that carries no
153    /// rows produces zero output records; a row whose width differs from the
154    /// column count is a hard error (never silently misaligned).
155    pub fn apply(&self, rec: &Value) -> Result<Vec<Value>, FaucetError> {
156        if !self.groups.is_empty() {
157            return self.apply_groups(rec);
158        }
159        let columns = self.column_names(rec)?;
160        let rows = row_candidates(&extract_records(rec, Some(&self.rows_path))?);
161        let mut out = Vec::with_capacity(rows.len());
162        for (i, row) in rows.into_iter().enumerate() {
163            let Value::Array(values) = row else {
164                return Err(FaucetError::Transform(format!(
165                    "zip_columns: row {i} at `{}` is not an array",
166                    self.rows_path
167                )));
168            };
169            if values.len() != columns.len() {
170                return Err(FaucetError::Transform(format!(
171                    "zip_columns: row {i} has {} value(s) but there are {} column(s)",
172                    values.len(),
173                    columns.len()
174                )));
175            }
176            let obj: Map<String, Value> = columns.iter().cloned().zip(values).collect();
177            out.push(Value::Object(obj));
178        }
179        Ok(out)
180    }
181
182    /// Resolve the column names, requiring every matched value to be a string.
183    fn column_names(&self, rec: &Value) -> Result<Vec<String>, FaucetError> {
184        let matched = extract_records(rec, Some(&self.columns_path))?;
185        // A single match that is itself an array (`columns_path: columns` where
186        // columns is already a string array) is unwrapped to its elements.
187        let candidates = column_candidates(&matched);
188        let mut names = Vec::with_capacity(candidates.len());
189        for c in candidates {
190            match c {
191                Value::String(s) => names.push(s),
192                other => {
193                    return Err(FaucetError::Transform(format!(
194                        "zip_columns: column name at `{}` is not a string: {other}",
195                        self.columns_path
196                    )));
197                }
198            }
199        }
200        if names.is_empty() {
201            return Err(FaucetError::Transform(format!(
202                "zip_columns: `columns_path` `{}` matched no column names",
203                self.columns_path
204            )));
205        }
206        Ok(names)
207    }
208
209    /// The `groups` form: every group's cells are zipped against its own
210    /// header list and the groups are merged into one record per row.
211    fn apply_groups(&self, rec: &Value) -> Result<Vec<Value>, FaucetError> {
212        let mut headers: Vec<Vec<String>> = Vec::with_capacity(self.groups.len());
213        let mut owner: Map<String, Value> = Map::new();
214        for g in &self.groups {
215            let names = g.header_names(rec)?;
216            for n in &names {
217                if let Some(Value::String(prev)) =
218                    owner.insert(n.clone(), Value::String(g.from.clone()))
219                {
220                    let which = if prev == g.from {
221                        format!("group '{prev}' names it twice")
222                    } else {
223                        format!("groups '{prev}' and '{}' both name it", g.from)
224                    };
225                    return Err(FaucetError::Transform(format!(
226                        "zip_columns: duplicate column '{n}': {which}"
227                    )));
228                }
229            }
230            headers.push(names);
231        }
232        let matched = extract_records(rec, Some(&self.rows_path))?;
233        let rows = match matched.as_slice() {
234            [Value::Array(inner)] => inner.clone(),
235            _ => matched,
236        };
237        let mut out = Vec::with_capacity(rows.len());
238        for (i, row) in rows.iter().enumerate() {
239            if !row.is_object() {
240                return Err(FaucetError::Transform(format!(
241                    "zip_columns: row {i} at `{}` is not an object",
242                    self.rows_path
243                )));
244            }
245            let mut obj = Map::new();
246            for (g, names) in self.groups.iter().zip(&headers) {
247                let Some(Value::Array(cells)) = path_get(row, &g.from) else {
248                    return Err(FaucetError::Transform(format!(
249                        "zip_columns: row {i} has no `{}` array (group '{}')",
250                        g.from, g.from
251                    )));
252                };
253                if cells.len() != names.len() {
254                    return Err(FaucetError::Transform(format!(
255                        "zip_columns: row {i}, group '{}': {} cell(s) but {} header(s)",
256                        g.from,
257                        cells.len(),
258                        names.len()
259                    )));
260                }
261                for (name, cell) in names.iter().zip(cells) {
262                    let v = match &g.value {
263                        Some(field) => path_get(cell, field).cloned().unwrap_or(Value::Null),
264                        None => cell.clone(),
265                    };
266                    obj.insert(name.clone(), v);
267                }
268            }
269            out.push(Value::Object(obj));
270        }
271        Ok(out)
272    }
273}
274
275impl CompiledGroup {
276    fn header_names(&self, rec: &Value) -> Result<Vec<String>, FaucetError> {
277        let matched = extract_records(rec, Some(&self.header))?;
278        let mut names = Vec::new();
279        for h in column_candidates(&matched) {
280            let name = match (&self.header_label, &h) {
281                (Some(label), Value::Object(_)) => path_get(&h, label).cloned(),
282                (None, _) => Some(h.clone()),
283                (Some(_), _) => None,
284            };
285            match name {
286                Some(Value::String(s)) => names.push(s),
287                _ => {
288                    return Err(FaucetError::Transform(format!(
289                        "zip_columns: group '{}': header at `{}` is not a string: {h}",
290                        self.from, self.header
291                    )));
292                }
293            }
294        }
295        Ok(names)
296    }
297}
298
299/// Resolve a dot path (`a.b`) inside a value.
300fn path_get<'a>(root: &'a Value, path: &str) -> Option<&'a Value> {
301    path.split('.').try_fold(root, |cur, seg| cur.get(seg))
302}
303
304/// Accept a bare path (`rows`, `columns[*].name`) by rooting it at `$`, while
305/// leaving an already-`$`-rooted expression untouched.
306fn normalize_path(path: &str) -> String {
307    let p = path.trim();
308    if p.starts_with('$') {
309        p.to_string()
310    } else {
311        format!("$.{p}")
312    }
313}
314
315/// Column candidates: a single array match (`columns_path: columns` pointing at
316/// a string array) is unwrapped to its elements; a `columns[*].name`-style match
317/// already yields the names directly.
318fn column_candidates(matched: &[Value]) -> Vec<Value> {
319    match matched {
320        [Value::Array(inner)] => inner.clone(),
321        other => other.to_vec(),
322    }
323}
324
325/// Row candidates: `rows` matches the rows array (one match, an array *of
326/// arrays*) → unwrap to the rows; `rows[*]` yields each row directly. Unwrapping
327/// only when every element is itself an array disambiguates a single-row
328/// `rows[*]` (one array of scalars) from the whole rows container.
329fn row_candidates(matched: &[Value]) -> Vec<Value> {
330    if let [Value::Array(inner)] = matched
331        && inner.iter().all(|v| matches!(v, Value::Array(_)))
332    {
333        return inner.clone();
334    }
335    matched.to_vec()
336}
337
338#[cfg(test)]
339mod tests {
340    use super::*;
341    use serde_json::json;
342
343    fn spec() -> CompiledZipColumns {
344        ZipColumnsSpec {
345            columns_path: "columns[*].name".into(),
346            rows_path: "rows".into(),
347            groups: vec![],
348        }
349        .compile()
350        .unwrap()
351    }
352
353    #[test]
354    fn zips_columns_into_row_objects() {
355        let rec = json!({
356            "columns": [{"name": "day"}, {"name": "sessions"}],
357            "rows": [["2026-01-01", 12], ["2026-01-02", 7]],
358        });
359        let out = spec().apply(&rec).unwrap();
360        assert_eq!(out.len(), 2);
361        assert_eq!(out[0], json!({"day": "2026-01-01", "sessions": 12}));
362        assert_eq!(out[1], json!({"day": "2026-01-02", "sessions": 7}));
363    }
364
365    #[test]
366    fn direct_string_array_columns_and_rows_star() {
367        let compiled = ZipColumnsSpec {
368            columns_path: "columns".into(),
369            rows_path: "rows[*]".into(),
370            groups: vec![],
371        }
372        .compile()
373        .unwrap();
374        let rec = json!({"columns": ["a", "b"], "rows": [[1, 2]]});
375        let out = compiled.apply(&rec).unwrap();
376        assert_eq!(out, vec![json!({"a": 1, "b": 2})]);
377    }
378
379    #[test]
380    fn no_rows_yields_no_records() {
381        let rec = json!({"columns": [{"name": "a"}], "rows": []});
382        assert!(spec().apply(&rec).unwrap().is_empty());
383    }
384
385    #[test]
386    fn width_mismatch_errors_clearly() {
387        let rec = json!({"columns": [{"name": "a"}, {"name": "b"}], "rows": [[1]]});
388        let err = spec().apply(&rec).unwrap_err();
389        let msg = err.to_string();
390        assert!(msg.contains("1 value") && msg.contains("2 column"), "{msg}");
391    }
392
393    #[test]
394    fn non_string_column_name_errors() {
395        let rec = json!({"columns": [{"name": 7}], "rows": [[1]]});
396        assert!(spec().apply(&rec).is_err());
397    }
398
399    #[test]
400    fn empty_paths_rejected_at_compile() {
401        assert!(
402            ZipColumnsSpec {
403                columns_path: "".into(),
404                rows_path: "rows".into(),
405                groups: vec![],
406            }
407            .compile()
408            .is_err()
409        );
410        assert!(
411            ZipColumnsSpec {
412                columns_path: "columns".into(),
413                rows_path: " ".into(),
414                groups: vec![],
415            }
416            .compile()
417            .is_err()
418        );
419    }
420
421    #[test]
422    fn into_stage_is_pagefn_and_flat_maps() {
423        let stage = spec_spec().into_stage().unwrap();
424        match stage {
425            TransformStage::PageFn(f) => {
426                let page = vec![json!({
427                    "columns": [{"name": "a"}],
428                    "rows": [[1], [2]],
429                })];
430                let out = f(page).unwrap();
431                assert_eq!(out, vec![json!({"a": 1}), json!({"a": 2})]);
432            }
433            other => panic!("expected PageFn, got {other:?}"),
434        }
435    }
436
437    fn spec_spec() -> ZipColumnsSpec {
438        ZipColumnsSpec {
439            columns_path: "columns[*].name".into(),
440            rows_path: "rows".into(),
441            groups: vec![],
442        }
443    }
444
445    fn grouped_report() -> ZipColumnsSpec {
446        serde_json::from_value(json!({
447            "rows_path": "$.rows[*]",
448            "groups": [
449                {"from": "dimensionValues", "header": "$.dimensionHeaders[*].name", "value": "value"},
450                {"from": "metricValues", "header": "$.metricHeaders[*]", "header_label": "name", "value": "value"}
451            ]
452        }))
453        .unwrap()
454    }
455
456    fn report() -> Value {
457        json!({
458            "dimensionHeaders": [{"name": "date"}, {"name": "country"}],
459            "metricHeaders": [{"name": "sessions", "type": "TYPE_INTEGER"}, {"name": "bounceRate", "type": "TYPE_FLOAT"}],
460            "rows": [
461                {"dimensionValues": [{"value": "20260901"}, {"value": "DE"}],
462                 "metricValues": [{"value": "1204"}, {"value": "0.41"}]},
463                {"dimensionValues": [{"value": "20260902"}, {}],
464                 "metricValues": [{"value": "9"}, null]}
465            ]
466        })
467    }
468
469    #[test]
470    fn groups_zipped_report_rows() {
471        let out = grouped_report()
472            .compile()
473            .unwrap()
474            .apply(&report())
475            .unwrap();
476        assert_eq!(
477            out,
478            vec![
479                json!({"date": "20260901", "country": "DE", "sessions": "1204", "bounceRate": "0.41"}),
480                json!({"date": "20260902", "country": null, "sessions": "9", "bounceRate": null}),
481            ]
482        );
483    }
484
485    #[test]
486    fn groups_with_scalar_cells_and_a_rows_container() {
487        let spec: ZipColumnsSpec = serde_json::from_value(json!({
488            "rows_path": "rows",
489            "groups": [{"from": "a.cells", "header": "names"}]
490        }))
491        .unwrap();
492        let rec = json!({"names": ["x", "y"], "rows": [{"a": {"cells": [1, 2]}}]});
493        let out = spec.compile().unwrap().apply(&rec).unwrap();
494        assert_eq!(out, vec![json!({"x": 1, "y": 2})]);
495    }
496
497    #[test]
498    fn groups_edge_cases() {
499        let c = grouped_report().compile().unwrap();
500        // An empty report (the API omits `rows`) yields nothing.
501        assert!(
502            c.apply(&json!({"dimensionHeaders": [], "metricHeaders": []}))
503                .unwrap()
504                .is_empty()
505        );
506        // Empty headers with empty cells yield an empty record.
507        let empty = json!({"dimensionHeaders": [], "metricHeaders": [],
508            "rows": [{"dimensionValues": [], "metricValues": []}]});
509        assert_eq!(c.apply(&empty).unwrap(), vec![json!({})]);
510
511        let err = |rec: Value| c.apply(&rec).unwrap_err().to_string();
512        let mut r = report();
513        r["rows"][1]["metricValues"] = json!([{"value": "1"}]);
514        let e = err(r);
515        assert!(
516            e.contains("row 1, group 'metricValues'") && e.contains("1 cell(s) but 2 header(s)"),
517            "{e}"
518        );
519
520        let mut r = report();
521        r["rows"][0].as_object_mut().unwrap().remove("metricValues");
522        let e = err(r);
523        assert!(e.contains("row 0 has no `metricValues` array"), "{e}");
524
525        let mut r = report();
526        r["metricHeaders"][0]["name"] = json!("date");
527        let e = err(r);
528        assert!(
529            e.contains("duplicate column 'date'")
530                && e.contains("'dimensionValues' and 'metricValues'"),
531            "{e}"
532        );
533
534        let mut r = report();
535        r["dimensionHeaders"][1]["name"] = json!("date");
536        assert!(err(r).contains("group 'dimensionValues' names it twice"));
537
538        let mut r = report();
539        r["dimensionHeaders"][0]["name"] = json!(7);
540        assert!(err(r).contains("group 'dimensionValues': header"));
541
542        let mut r = report();
543        r["metricHeaders"][0] = json!("sessions");
544        assert!(err(r).contains("group 'metricValues': header"));
545
546        let mut r = report();
547        r["rows"][0] = json!([1]);
548        assert!(err(r).contains("row 0 at `$.rows[*]` is not an object"));
549    }
550
551    #[test]
552    fn groups_compile_validation() {
553        let bad = |v: Value| {
554            serde_json::from_value::<ZipColumnsSpec>(v)
555                .unwrap()
556                .compile()
557                .unwrap_err()
558                .to_string()
559        };
560        let g = json!({"from": "a", "header": "h"});
561        assert!(bad(json!({"rows_path": "r"})).contains("exactly one"));
562        assert!(
563            bad(json!({"rows_path": "r", "columns_path": "c", "groups": [g]}))
564                .contains("exactly one")
565        );
566        assert!(
567            bad(json!({"rows_path": "r", "groups": [{"from": " ", "header": "h"}]}))
568                .contains("group 0")
569        );
570        assert!(
571            bad(json!({"rows_path": "r", "groups": [{"from": "a", "header": "h", "value": ""}]}))
572                .contains("group 'a'")
573        );
574        assert!(
575            bad(json!({"rows_path": "r", "groups": [{"from": "a", "header": "h", "header_label": " "}]}))
576                .contains("group 'a'")
577        );
578        assert!(bad(json!({"rows_path": " ", "groups": [g]})).contains("rows_path"));
579        let s = grouped_report();
580        assert_eq!(
581            serde_json::from_value::<ZipColumnsSpec>(serde_json::to_value(&s).unwrap()).unwrap(),
582            s
583        );
584    }
585}