1use 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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
28#[serde(deny_unknown_fields)]
29pub struct ZipColumnsSpec {
30 #[serde(default, skip_serializing_if = "String::is_empty")]
35 pub columns_path: String,
36 pub rows_path: String,
40 #[serde(default, skip_serializing_if = "Vec::is_empty")]
46 pub groups: Vec<ColumnGroupSpec>,
47}
48
49#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
51#[serde(deny_unknown_fields)]
52pub struct ColumnGroupSpec {
53 pub from: String,
56 pub header: String,
60 #[serde(default, skip_serializing_if = "Option::is_none")]
63 pub header_label: Option<String>,
64 #[serde(default, skip_serializing_if = "Option::is_none")]
67 pub value: Option<String>,
68}
69
70impl ZipColumnsSpec {
71 pub fn compile(&self) -> Result<CompiledZipColumns, FaucetError> {
73 CompiledZipColumns::compile(self)
74 }
75
76 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#[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 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 fn column_names(&self, rec: &Value) -> Result<Vec<String>, FaucetError> {
184 let matched = extract_records(rec, Some(&self.columns_path))?;
185 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 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
299fn path_get<'a>(root: &'a Value, path: &str) -> Option<&'a Value> {
301 path.split('.').try_fold(root, |cur, seg| cur.get(seg))
302}
303
304fn 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
315fn column_candidates(matched: &[Value]) -> Vec<Value> {
319 match matched {
320 [Value::Array(inner)] => inner.clone(),
321 other => other.to_vec(),
322 }
323}
324
325fn 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 assert!(
502 c.apply(&json!({"dimensionHeaders": [], "metricHeaders": []}))
503 .unwrap()
504 .is_empty()
505 );
506 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}