1use crate::config::TransformSpec;
7use crate::error::{CliError, CliResult};
8#[cfg(feature = "transforms")]
9use faucet_core::{CastOnError, CastType, KeyCaseMode, ValueCaseMode};
10#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
11use faucet_core::{JsonSchema, schema_for};
12use faucet_core::{RecordTransform, TransformStage};
13#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
14use serde::Deserialize;
15use serde_json::Value;
16#[cfg(feature = "transforms")]
17use std::collections::HashMap;
18
19#[cfg(feature = "transforms")]
21#[derive(Debug, Deserialize, JsonSchema)]
22struct FlattenConfig {
23 #[serde(default = "default_separator")]
25 separator: String,
26}
27
28#[cfg(feature = "transforms")]
29fn default_separator() -> String {
30 "__".to_owned()
31}
32
33#[cfg(feature = "transforms")]
35#[derive(Debug, Deserialize, JsonSchema)]
36struct RenameKeysConfig {
37 pattern: String,
39 replacement: String,
41}
42
43#[cfg(feature = "transforms")]
44#[derive(Debug, Deserialize, JsonSchema)]
45struct FieldsConfig {
46 fields: Vec<String>,
48}
49
50#[cfg(feature = "transforms")]
51#[derive(Debug, Deserialize, JsonSchema)]
52struct SetConfig {
53 values: serde_json::Map<String, Value>,
55}
56
57#[cfg(feature = "transforms")]
58#[derive(Debug, Deserialize, JsonSchema)]
59struct RenameFieldConfig {
60 fields: HashMap<String, String>,
62}
63
64#[cfg(feature = "transforms")]
65#[derive(Debug, Deserialize, JsonSchema)]
66struct CastConfig {
67 fields: HashMap<String, CastType>,
69 #[serde(default)]
71 on_error: CastOnError,
72}
73
74#[cfg(feature = "transforms")]
75#[derive(Debug, Deserialize, JsonSchema)]
76struct RedactConfig {
77 fields: Vec<String>,
79 #[serde(default = "default_mask")]
81 mask: Value,
82}
83
84#[cfg(feature = "transforms")]
85fn default_mask() -> Value {
86 Value::String("***".to_owned())
87}
88
89#[cfg(feature = "transforms")]
90#[derive(Debug, Deserialize, JsonSchema)]
91struct ValueCaseConfig {
92 fields: Vec<String>,
94 mode: ValueCaseMode,
96}
97
98#[cfg(feature = "transforms")]
99#[derive(Debug, Deserialize, JsonSchema)]
100struct SpellSymbolsConfig {
101 #[serde(default)]
103 extra: HashMap<String, String>,
104 #[serde(default = "default_spell_separator")]
106 separator: String,
107}
108
109#[cfg(feature = "transforms")]
110fn default_spell_separator() -> String {
111 " ".to_owned()
112}
113
114#[cfg(feature = "transforms")]
115#[derive(Debug, Deserialize, JsonSchema)]
116struct KeysCaseConfig {
117 mode: KeyCaseMode,
119}
120
121#[cfg(feature = "transform-filter")]
122#[derive(Debug, Deserialize, JsonSchema)]
123struct FilterConfig {
124 path: String,
126 op: faucet_core::FilterOp,
128 #[serde(default, skip_serializing_if = "Option::is_none")]
130 value: Option<Value>,
131}
132
133#[cfg(feature = "transform-explode")]
134#[derive(Debug, Deserialize, JsonSchema)]
135struct ExplodeConfig {
136 path: String,
138 #[serde(default, skip_serializing_if = "Option::is_none")]
141 prefix: Option<String>,
142 #[serde(default = "default_explode_separator_cli")]
144 separator: String,
145 #[serde(default)]
148 on_missing: faucet_core::OnMissing,
149}
150
151#[cfg(feature = "transform-explode")]
152fn default_explode_separator_cli() -> String {
153 "_".to_owned()
154}
155
156#[cfg(feature = "transform-cdc-unwrap")]
157#[derive(Debug, Deserialize, JsonSchema)]
158struct CdcUnwrapConfig {
159 #[serde(default = "cdc_op_field")]
161 op_field: String,
162 #[serde(default = "cdc_after_field")]
164 after_field: String,
165 #[serde(default = "cdc_before_field")]
167 before_field: String,
168 #[serde(default = "cdc_key_field")]
170 key_field: String,
171 #[serde(default = "cdc_marker_field")]
173 marker_field: String,
174 #[serde(default = "cdc_delete_ops")]
176 delete_ops: Vec<String>,
177 #[serde(default = "cdc_drop_ops")]
179 drop_ops: Vec<String>,
180}
181
182#[cfg(feature = "transform-cdc-unwrap")]
183fn cdc_op_field() -> String {
184 "op".into()
185}
186#[cfg(feature = "transform-cdc-unwrap")]
187fn cdc_after_field() -> String {
188 "after".into()
189}
190#[cfg(feature = "transform-cdc-unwrap")]
191fn cdc_before_field() -> String {
192 "before".into()
193}
194#[cfg(feature = "transform-cdc-unwrap")]
195fn cdc_key_field() -> String {
196 "document_key".into()
197}
198#[cfg(feature = "transform-cdc-unwrap")]
199fn cdc_marker_field() -> String {
200 "__op".into()
201}
202#[cfg(feature = "transform-cdc-unwrap")]
203fn cdc_delete_ops() -> Vec<String> {
204 vec!["d".into(), "delete".into()]
205}
206#[cfg(feature = "transform-cdc-unwrap")]
207fn cdc_drop_ops() -> Vec<String> {
208 vec!["ddl".into(), "truncate".into()]
209}
210
211struct TransformDef {
218 kind: &'static str,
219 description: &'static str,
220 schema_fn: fn() -> Value,
221 compile_fn: fn(&str, Value) -> CliResult<TransformStage>,
222}
223
224fn registry() -> Vec<TransformDef> {
230 #[allow(unused_mut)]
231 let mut defs: Vec<TransformDef> = Vec::new();
232 #[cfg(feature = "transforms")]
233 {
234 defs.extend(vec![
235 TransformDef {
236 kind: "flatten",
237 description: "Flatten nested objects into a single level (configurable separator).",
238 schema_fn: || schema::<FlattenConfig>(),
239 compile_fn: |kind, config| {
240 let cfg = decode::<FlattenConfig>(kind, config)?;
241 Ok(TransformStage::Map(RecordTransform::Flatten {
242 separator: cfg.separator,
243 }))
244 },
245 },
246 TransformDef {
247 kind: "rename_keys",
248 description: "Rewrite every key via a regex pattern + replacement.",
249 schema_fn: || schema::<RenameKeysConfig>(),
250 compile_fn: |kind, config| {
251 let cfg = decode::<RenameKeysConfig>(kind, config)?;
252 Ok(TransformStage::Map(RecordTransform::RenameKeys {
253 pattern: cfg.pattern,
254 replacement: cfg.replacement,
255 }))
256 },
257 },
258 TransformDef {
259 kind: "keys_case",
260 description: "Re-case every key (snake / camel / pascal / kebab / screaming_snake).",
261 schema_fn: || schema::<KeysCaseConfig>(),
262 compile_fn: |kind, config| {
263 let cfg = decode::<KeysCaseConfig>(kind, config)?;
264 Ok(TransformStage::Map(RecordTransform::KeysCase {
265 mode: cfg.mode,
266 }))
267 },
268 },
269 TransformDef {
270 kind: "select",
271 description: "Keep only the listed top-level fields; drop the rest.",
272 schema_fn: || schema::<FieldsConfig>(),
273 compile_fn: |kind, config| {
274 let cfg = decode::<FieldsConfig>(kind, config)?;
275 Ok(TransformStage::Map(RecordTransform::Select {
276 fields: cfg.fields,
277 }))
278 },
279 },
280 TransformDef {
281 kind: "drop",
282 description: "Remove the listed top-level fields.",
283 schema_fn: || schema::<FieldsConfig>(),
284 compile_fn: |kind, config| {
285 let cfg = decode::<FieldsConfig>(kind, config)?;
286 Ok(TransformStage::Map(RecordTransform::Drop {
287 fields: cfg.fields,
288 }))
289 },
290 },
291 TransformDef {
292 kind: "set",
293 description: "Set named fields to constant values on every record.",
294 schema_fn: || schema::<SetConfig>(),
295 compile_fn: |kind, config| {
296 let cfg = decode::<SetConfig>(kind, config)?;
297 Ok(TransformStage::Map(RecordTransform::Set {
298 values: cfg.values,
299 }))
300 },
301 },
302 TransformDef {
303 kind: "rename_field",
304 description: "Rename specific top-level fields by name.",
305 schema_fn: || schema::<RenameFieldConfig>(),
306 compile_fn: |kind, config| {
307 let cfg = decode::<RenameFieldConfig>(kind, config)?;
308 Ok(TransformStage::Map(RecordTransform::RenameField {
309 fields: cfg.fields,
310 }))
311 },
312 },
313 TransformDef {
314 kind: "cast",
315 description: "Coerce named fields to int / float / bool / string / timestamp.",
316 schema_fn: || schema::<CastConfig>(),
317 compile_fn: |kind, config| {
318 let cfg = decode::<CastConfig>(kind, config)?;
319 Ok(TransformStage::Map(RecordTransform::Cast {
320 fields: cfg.fields,
321 on_error: cfg.on_error,
322 }))
323 },
324 },
325 TransformDef {
326 kind: "redact",
327 description: "Overwrite the listed fields with a mask value (default `***`).",
328 schema_fn: || schema::<RedactConfig>(),
329 compile_fn: |kind, config| {
330 let cfg = decode::<RedactConfig>(kind, config)?;
331 Ok(TransformStage::Map(RecordTransform::Redact {
332 fields: cfg.fields,
333 mask: cfg.mask,
334 }))
335 },
336 },
337 TransformDef {
338 kind: "value_case",
339 description: "Lowercase, uppercase, or trim the value of named string fields.",
340 schema_fn: || schema::<ValueCaseConfig>(),
341 compile_fn: |kind, config| {
342 let cfg = decode::<ValueCaseConfig>(kind, config)?;
343 Ok(TransformStage::Map(RecordTransform::ValueCase {
344 fields: cfg.fields,
345 mode: cfg.mode,
346 }))
347 },
348 },
349 TransformDef {
350 kind: "spell_symbols",
351 description: "Replace punctuation/symbols in string values with their spelled-out words.",
352 schema_fn: || schema::<SpellSymbolsConfig>(),
353 compile_fn: |kind, config| {
354 let cfg = decode::<SpellSymbolsConfig>(kind, config)?;
355 Ok(TransformStage::Map(RecordTransform::SpellSymbols {
356 extra: cfg.extra,
357 separator: cfg.separator,
358 }))
359 },
360 },
361 #[cfg(feature = "transform-filter")]
362 TransformDef {
363 kind: "filter",
364 description: "Keep records where a JSONPath predicate is true.",
365 schema_fn: || schema::<FilterConfig>(),
366 compile_fn: |kind, config| {
367 let cfg = decode::<FilterConfig>(kind, config)?;
368 let stage = TransformStage::Filter(faucet_core::FilterSpec {
370 path: cfg.path,
371 op: cfg.op,
372 value: cfg.value,
373 });
374 faucet_core::compile_stage(&stage).map_err(|e| match e {
375 faucet_core::FaucetError::Transform(msg) => CliError::InvalidTransform {
376 name: kind.to_owned(),
377 message: msg,
378 },
379 other => CliError::InvalidTransform {
380 name: kind.to_owned(),
381 message: format!("{other}"),
382 },
383 })?;
384 Ok(stage)
385 },
386 },
387 #[cfg(feature = "transform-explode")]
388 TransformDef {
389 kind: "explode",
390 description: "Expand an array field into one record per element.",
391 schema_fn: || schema::<ExplodeConfig>(),
392 compile_fn: |kind, config| {
393 let cfg = decode::<ExplodeConfig>(kind, config)?;
394 let stage = TransformStage::Explode(faucet_core::ExplodeSpec {
395 path: cfg.path,
396 prefix: cfg.prefix,
397 separator: cfg.separator,
398 on_missing: cfg.on_missing,
399 });
400 faucet_core::compile_stage(&stage).map_err(|e| match e {
401 faucet_core::FaucetError::Transform(msg) => CliError::InvalidTransform {
402 name: kind.to_owned(),
403 message: msg,
404 },
405 other => CliError::InvalidTransform {
406 name: kind.to_owned(),
407 message: format!("{other}"),
408 },
409 })?;
410 Ok(stage)
411 },
412 },
413 ]);
414 }
415 #[cfg(feature = "transform-sql")]
416 {
417 defs.push(TransformDef {
418 kind: "sql",
419 description: "Run DuckDB SQL over the whole page; records are the `batch` relation.",
420 schema_fn: || schema_sql(),
421 compile_fn: |kind, config| {
422 let cfg: faucet_transform_sql::SqlTransformConfig = decode_sql(kind, config)?;
423 let transform = faucet_transform_sql::SqlTransform::compile(&cfg).map_err(|e| {
424 let message = match &e {
425 faucet_core::FaucetError::Transform(m)
426 | faucet_core::FaucetError::Config(m) => m.clone(),
427 other => format!("{other}"),
428 };
429 CliError::InvalidTransform {
430 name: kind.to_owned(),
431 message,
432 }
433 })?;
434 Ok(transform.into_page_stage())
435 },
436 });
437 }
438 #[cfg(feature = "transform-cdc-unwrap")]
439 {
440 defs.push(TransformDef {
441 kind: "cdc_unwrap",
442 description: "Normalize a CDC envelope into a flat row + delete marker (for upsert sinks).",
443 schema_fn: || schema::<CdcUnwrapConfig>(),
444 compile_fn: |kind, config| {
445 let cfg = decode::<CdcUnwrapConfig>(kind, config)?;
446 Ok(faucet_core::TransformStage::CdcUnwrap(faucet_core::CdcUnwrapSpec {
447 op_field: cfg.op_field,
448 after_field: cfg.after_field,
449 before_field: cfg.before_field,
450 key_field: cfg.key_field,
451 marker_field: cfg.marker_field,
452 delete_ops: cfg.delete_ops,
453 drop_ops: cfg.drop_ops,
454 }))
455 },
456 });
457 }
458 defs
459}
460
461pub fn compile_transforms(specs: &[TransformSpec]) -> CliResult<Vec<TransformStage>> {
466 let mut out = Vec::with_capacity(specs.len());
467 for s in specs {
468 out.push(compile_one(s)?);
469 }
470 Ok(out)
471}
472
473fn compile_one(spec: &TransformSpec) -> CliResult<TransformStage> {
474 match registry().into_iter().find(|t| t.kind == spec.kind) {
475 Some(def) => (def.compile_fn)(&spec.kind, spec.config.clone()),
476 None => Err(unknown_transform(&spec.kind)),
477 }
478}
479
480pub fn transform_descriptions() -> Vec<(&'static str, &'static str)> {
483 registry()
484 .into_iter()
485 .map(|t| (t.kind, t.description))
486 .collect()
487}
488
489pub fn available_transforms() -> Vec<&'static str> {
491 registry().into_iter().map(|t| t.kind).collect()
492}
493
494#[cfg(feature = "quality")]
499pub fn quality_descriptions() -> Vec<(&'static str, &'static str)> {
500 let mut checks = vec![
501 ("not_null", "field present and non-null"),
502 ("not_empty", "string non-empty after trim"),
503 ("regex_match", "string matches a regex"),
504 ("value_in_set", "value is in an allowed set"),
505 ("not_in_set", "value is not in a forbidden set"),
506 ("compare", "numeric/scalar comparison (gt/gte/lt/lte/eq/ne)"),
507 ("type_is", "value is of an expected JSON type"),
508 ("string_length", "string length within [min,max]"),
509 ];
510 #[cfg(feature = "quality-jsonschema")]
511 checks.push((
512 "json_schema",
513 "record validates against a JSON Schema (feature-gated)",
514 ));
515 checks.extend([
516 ("row_count", "batch row count within [min,max]"),
517 ("null_rate", "batch null rate of a field <= max"),
518 ("unique", "composite key unique within the batch"),
519 (
520 "distinct_count",
521 "distinct values of a field within [min,max]",
522 ),
523 ]);
524 checks
525}
526
527pub fn transform_schema(name: &str) -> CliResult<Value> {
531 registry()
532 .into_iter()
533 .find(|t| t.kind == name)
534 .map(|t| (t.schema_fn)())
535 .ok_or_else(|| unknown_transform(name))
536}
537
538fn unknown_transform(name: &str) -> CliError {
539 let available = available_transforms();
540 CliError::UnknownTransform {
541 name: name.to_owned(),
542 available: if available.is_empty() {
543 "(none — rebuild faucet-cli with the `transforms` feature enabled)".to_owned()
544 } else {
545 available.join(", ")
546 },
547 }
548}
549
550#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
551fn schema<T: JsonSchema>() -> Value {
552 serde_json::to_value(schema_for!(T)).unwrap_or_else(|_| serde_json::json!({"type": "object"}))
553}
554
555#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
556fn decode<T: serde::de::DeserializeOwned>(name: &str, config: Value) -> CliResult<T> {
557 serde_json::from_value(config).map_err(|e| CliError::InvalidTransform {
558 name: name.to_owned(),
559 message: e.to_string(),
560 })
561}
562
563#[cfg(feature = "transform-sql")]
564fn schema_sql() -> Value {
565 serde_json::to_value(faucet_core::schema_for!(
566 faucet_transform_sql::SqlTransformConfig
567 ))
568 .unwrap_or(Value::Null)
569}
570
571#[cfg(feature = "transform-sql")]
572fn decode_sql(kind: &str, config: Value) -> CliResult<faucet_transform_sql::SqlTransformConfig> {
573 serde_json::from_value(config).map_err(|e| CliError::InvalidTransform {
574 name: kind.to_owned(),
575 message: e.to_string(),
576 })
577}
578
579#[cfg(test)]
580mod tests {
581 use super::*;
582 use serde_json::json;
583
584 #[test]
585 fn empty_list_compiles_to_empty() {
586 let out = compile_transforms(&[]).unwrap();
587 assert!(out.is_empty());
588 }
589
590 #[cfg(feature = "transforms")]
591 #[test]
592 fn compiles_keys_case_and_flatten() {
593 let specs = vec![
594 TransformSpec {
595 kind: "keys_case".into(),
596 config: json!({"mode": "snake"}),
597 },
598 TransformSpec {
599 kind: "flatten".into(),
600 config: json!({"separator": "."}),
601 },
602 ];
603 let out = compile_transforms(&specs).unwrap();
604 assert_eq!(out.len(), 2);
605 }
606
607 #[cfg(feature = "transforms")]
608 #[test]
609 fn keys_case_rejects_unknown_mode() {
610 let specs = vec![TransformSpec {
611 kind: "keys_case".into(),
612 config: json!({"mode": "spongebob"}),
613 }];
614 let err = compile_transforms(&specs).unwrap_err();
615 match err {
616 CliError::InvalidTransform { name, .. } => assert_eq!(name, "keys_case"),
617 other => panic!("expected InvalidTransform, got {other:?}"),
618 }
619 }
620
621 #[cfg(feature = "transforms")]
622 #[test]
623 fn keys_case_requires_mode() {
624 let specs = vec![TransformSpec {
625 kind: "keys_case".into(),
626 config: json!({}),
627 }];
628 let err = compile_transforms(&specs).unwrap_err();
629 match err {
630 CliError::InvalidTransform { name, .. } => assert_eq!(name, "keys_case"),
631 other => panic!("expected InvalidTransform, got {other:?}"),
632 }
633 }
634
635 #[cfg(feature = "transforms")]
636 #[test]
637 fn snake_case_kind_is_no_longer_recognized() {
638 let specs = vec![TransformSpec {
640 kind: "snake_case".into(),
641 config: json!({}),
642 }];
643 let err = compile_transforms(&specs).unwrap_err();
644 match err {
645 CliError::UnknownTransform { name, .. } => assert_eq!(name, "snake_case"),
646 other => panic!("expected UnknownTransform, got {other:?}"),
647 }
648 }
649
650 #[cfg(feature = "transforms")]
651 #[test]
652 fn rename_keys_requires_pattern_and_replacement() {
653 let specs = vec![TransformSpec {
654 kind: "rename_keys".into(),
655 config: json!({"pattern": "^_"}),
656 }];
657 let err = compile_transforms(&specs).unwrap_err();
658 match err {
659 CliError::InvalidTransform { name, .. } => assert_eq!(name, "rename_keys"),
660 other => panic!("expected InvalidTransform, got {other:?}"),
661 }
662 }
663
664 #[test]
665 fn unknown_transform_errors() {
666 let specs = vec![TransformSpec {
667 kind: "make_uppercase".into(),
668 config: json!({}),
669 }];
670 let err = compile_transforms(&specs).unwrap_err();
671 match err {
672 CliError::UnknownTransform { name, .. } => assert_eq!(name, "make_uppercase"),
673 other => panic!("expected UnknownTransform, got {other:?}"),
674 }
675 }
676
677 #[cfg(feature = "transforms")]
678 #[test]
679 fn compiles_select_and_drop() {
680 let specs = vec![
681 TransformSpec {
682 kind: "select".into(),
683 config: json!({"fields": ["id", "name"]}),
684 },
685 TransformSpec {
686 kind: "drop".into(),
687 config: json!({"fields": ["secret"]}),
688 },
689 ];
690 let out = compile_transforms(&specs).unwrap();
691 assert_eq!(out.len(), 2);
692 }
693
694 #[cfg(feature = "transforms")]
695 #[test]
696 fn compiles_set_with_object_values() {
697 let specs = vec![TransformSpec {
698 kind: "set".into(),
699 config: json!({"values": {"_source": "api", "version": 1}}),
700 }];
701 let out = compile_transforms(&specs).unwrap();
702 assert_eq!(out.len(), 1);
703 }
704
705 #[cfg(feature = "transforms")]
706 #[test]
707 fn compiles_rename_field() {
708 let specs = vec![TransformSpec {
709 kind: "rename_field".into(),
710 config: json!({"fields": {"old": "new"}}),
711 }];
712 let out = compile_transforms(&specs).unwrap();
713 assert_eq!(out.len(), 1);
714 }
715
716 #[cfg(feature = "transforms")]
717 #[test]
718 fn compiles_cast_with_default_on_error() {
719 let specs = vec![TransformSpec {
720 kind: "cast".into(),
721 config: json!({"fields": {"age": "int", "price": "float"}}),
722 }];
723 let out = compile_transforms(&specs).unwrap();
724 assert_eq!(out.len(), 1);
725 }
726
727 #[cfg(feature = "transforms")]
728 #[test]
729 fn cast_rejects_unknown_target_type() {
730 let specs = vec![TransformSpec {
731 kind: "cast".into(),
732 config: json!({"fields": {"x": "uuid"}}),
733 }];
734 let err = compile_transforms(&specs).unwrap_err();
735 match err {
736 CliError::InvalidTransform { name, .. } => assert_eq!(name, "cast"),
737 other => panic!("expected InvalidTransform, got {other:?}"),
738 }
739 }
740
741 #[cfg(feature = "transforms")]
742 #[test]
743 fn cast_rejects_unknown_on_error_mode() {
744 let specs = vec![TransformSpec {
745 kind: "cast".into(),
746 config: json!({"fields": {"x": "int"}, "on_error": "explode"}),
747 }];
748 let err = compile_transforms(&specs).unwrap_err();
749 match err {
750 CliError::InvalidTransform { name, .. } => assert_eq!(name, "cast"),
751 other => panic!("expected InvalidTransform, got {other:?}"),
752 }
753 }
754
755 #[cfg(feature = "transforms")]
756 #[test]
757 fn redact_uses_default_mask_when_omitted() {
758 let specs = vec![TransformSpec {
759 kind: "redact".into(),
760 config: json!({"fields": ["ssn"]}),
761 }];
762 let out = compile_transforms(&specs).unwrap();
763 assert_eq!(out.len(), 1);
764 }
765
766 #[cfg(feature = "transforms")]
767 #[test]
768 fn value_case_requires_mode() {
769 let specs = vec![TransformSpec {
770 kind: "value_case".into(),
771 config: json!({"fields": ["email"]}),
772 }];
773 let err = compile_transforms(&specs).unwrap_err();
774 match err {
775 CliError::InvalidTransform { name, .. } => assert_eq!(name, "value_case"),
776 other => panic!("expected InvalidTransform, got {other:?}"),
777 }
778 }
779
780 #[cfg(feature = "transforms")]
781 #[test]
782 fn available_transforms_lists_every_kind() {
783 let names = available_transforms();
784 for expected in [
785 "flatten",
786 "rename_keys",
787 "keys_case",
788 "select",
789 "drop",
790 "set",
791 "rename_field",
792 "cast",
793 "redact",
794 "value_case",
795 "spell_symbols",
796 "filter",
797 "explode",
798 ] {
799 assert!(names.contains(&expected), "missing {expected}");
800 }
801 assert!(
802 !names.contains(&"snake_case"),
803 "snake_case must be removed in favour of keys_case"
804 );
805 }
806
807 #[cfg(feature = "transforms")]
808 #[test]
809 fn transform_descriptions_covers_every_compiled_kind() {
810 let names = available_transforms();
813 let desc_names: Vec<&'static str> = transform_descriptions()
814 .into_iter()
815 .map(|(n, _)| n)
816 .collect();
817 assert_eq!(names, desc_names);
818 for (_, desc) in transform_descriptions() {
819 assert!(!desc.is_empty(), "every transform needs a description");
820 }
821 }
822
823 #[cfg(feature = "transforms")]
824 #[test]
825 fn transform_schema_returns_object_for_every_kind() {
826 for name in available_transforms() {
827 let schema = transform_schema(name).unwrap_or_else(|e| {
828 panic!("schema lookup failed for {name}: {e}");
829 });
830 assert!(schema.is_object(), "schema for {name} must be an object");
831 }
832 }
833
834 #[cfg(feature = "transforms")]
835 #[test]
836 fn transform_schema_select_and_drop_share_shape() {
837 let select = transform_schema("select").unwrap();
840 let drop = transform_schema("drop").unwrap();
841 assert_eq!(select, drop);
842 }
843
844 #[test]
845 fn transform_schema_unknown_errors_with_available_list() {
846 let err = transform_schema("make_uppercase").unwrap_err();
847 match err {
848 CliError::UnknownTransform { name, available } => {
849 assert_eq!(name, "make_uppercase");
850 #[cfg(feature = "transforms")]
851 assert!(available.contains("flatten"), "{available}");
852 #[cfg(not(feature = "transforms"))]
853 assert!(available.contains("rebuild"), "{available}");
854 }
855 other => panic!("expected UnknownTransform, got {other:?}"),
856 }
857 }
858
859 #[cfg(feature = "transform-filter")]
860 #[test]
861 fn compiles_filter_eq() {
862 let specs = vec![TransformSpec {
863 kind: "filter".into(),
864 config: json!({"path": "status", "op": "eq", "value": "active"}),
865 }];
866 let out = compile_transforms(&specs).unwrap();
867 assert_eq!(out.len(), 1);
868 assert!(matches!(out[0], TransformStage::Filter(_)));
869 }
870
871 #[cfg(feature = "transform-filter")]
872 #[test]
873 fn filter_rejects_in_with_non_array_value() {
874 let specs = vec![TransformSpec {
875 kind: "filter".into(),
876 config: json!({"path": "v", "op": "in", "value": "scalar"}),
877 }];
878 let err = compile_transforms(&specs).unwrap_err();
879 match err {
880 CliError::InvalidTransform { name, message } => {
881 assert_eq!(name, "filter");
882 assert!(message.contains("requires an array"), "{message}");
883 }
884 other => panic!("expected InvalidTransform, got {other:?}"),
885 }
886 }
887
888 #[cfg(feature = "transform-filter")]
889 #[test]
890 fn filter_rejects_exists_with_value() {
891 let specs = vec![TransformSpec {
892 kind: "filter".into(),
893 config: json!({"path": "v", "op": "exists", "value": "x"}),
894 }];
895 let err = compile_transforms(&specs).unwrap_err();
896 match err {
897 CliError::InvalidTransform { name, .. } => assert_eq!(name, "filter"),
898 other => panic!("expected InvalidTransform, got {other:?}"),
899 }
900 }
901
902 #[cfg(feature = "transform-filter")]
903 #[test]
904 fn filter_rejects_bad_path() {
905 let specs = vec![TransformSpec {
906 kind: "filter".into(),
907 config: json!({"path": "$..items", "op": "exists"}),
908 }];
909 let err = compile_transforms(&specs).unwrap_err();
910 match err {
911 CliError::InvalidTransform { name, .. } => assert_eq!(name, "filter"),
912 other => panic!("expected InvalidTransform, got {other:?}"),
913 }
914 }
915
916 #[cfg(feature = "transform-explode")]
917 #[test]
918 fn compiles_explode_with_defaults() {
919 let specs = vec![TransformSpec {
920 kind: "explode".into(),
921 config: json!({"path": "items"}),
922 }];
923 let out = compile_transforms(&specs).unwrap();
924 assert_eq!(out.len(), 1);
925 assert!(matches!(out[0], TransformStage::Explode(_)));
926 }
927
928 #[cfg(feature = "transform-explode")]
929 #[test]
930 fn compiles_explode_with_custom_prefix_and_on_missing() {
931 let specs = vec![TransformSpec {
932 kind: "explode".into(),
933 config: json!({
934 "path": "items",
935 "prefix": "item",
936 "separator": "_",
937 "on_missing": "drop"
938 }),
939 }];
940 let out = compile_transforms(&specs).unwrap();
941 assert_eq!(out.len(), 1);
942 }
943
944 #[cfg(feature = "transform-explode")]
945 #[test]
946 fn explode_rejects_bad_path() {
947 let specs = vec![TransformSpec {
948 kind: "explode".into(),
949 config: json!({"path": "$..items"}),
950 }];
951 let err = compile_transforms(&specs).unwrap_err();
952 match err {
953 CliError::InvalidTransform { name, .. } => assert_eq!(name, "explode"),
954 other => panic!("expected InvalidTransform, got {other:?}"),
955 }
956 }
957
958 #[cfg(feature = "transform-explode")]
959 #[test]
960 fn explode_rejects_invalid_on_missing() {
961 let specs = vec![TransformSpec {
962 kind: "explode".into(),
963 config: json!({"path": "items", "on_missing": "explode_harder"}),
964 }];
965 let err = compile_transforms(&specs).unwrap_err();
966 match err {
967 CliError::InvalidTransform { name, .. } => assert_eq!(name, "explode"),
968 other => panic!("expected InvalidTransform, got {other:?}"),
969 }
970 }
971
972 #[cfg(feature = "transform-sql")]
973 #[test]
974 fn compiles_sql_to_page_fn() {
975 let specs = vec![TransformSpec {
976 kind: "sql".into(),
977 config: json!({"query": "SELECT * FROM batch"}),
978 }];
979 let out = compile_transforms(&specs).unwrap();
980 assert_eq!(out.len(), 1);
981 assert!(matches!(out[0], faucet_core::TransformStage::PageFn(_)));
982 }
983
984 #[cfg(feature = "transform-sql")]
985 #[test]
986 fn sql_bad_query_is_invalid_transform() {
987 let specs = vec![TransformSpec {
988 kind: "sql".into(),
989 config: json!({"query": "SELEKT bad"}),
990 }];
991 let err = compile_transforms(&specs).unwrap_err();
992 match err {
993 CliError::InvalidTransform { name, .. } => assert_eq!(name, "sql"),
994 other => panic!("expected InvalidTransform, got {other:?}"),
995 }
996 }
997
998 #[cfg(feature = "transform-sql")]
999 #[test]
1000 fn sql_schema_and_listing_present() {
1001 assert!(transform_schema("sql").is_ok());
1002 assert!(available_transforms().contains(&"sql"));
1003 }
1004
1005 #[cfg(feature = "transform-cdc-unwrap")]
1006 #[test]
1007 fn compiles_cdc_unwrap_with_defaults() {
1008 let specs = vec![TransformSpec {
1009 kind: "cdc_unwrap".into(),
1010 config: json!({}),
1011 }];
1012 let out = compile_transforms(&specs).unwrap();
1013 assert_eq!(out.len(), 1);
1014 assert!(matches!(out[0], faucet_core::TransformStage::CdcUnwrap(_)));
1015 }
1016
1017 #[cfg(feature = "transform-cdc-unwrap")]
1018 #[test]
1019 fn cdc_unwrap_schema_and_listing_present() {
1020 assert!(transform_schema("cdc_unwrap").is_ok());
1021 assert!(available_transforms().contains(&"cdc_unwrap"));
1022 }
1023
1024 #[cfg(feature = "quality")]
1025 #[test]
1026 fn quality_descriptions_has_one_entry_per_check() {
1027 #[cfg(feature = "quality-jsonschema")]
1033 assert_eq!(quality_descriptions().len(), 13);
1034 #[cfg(not(feature = "quality-jsonschema"))]
1035 assert_eq!(quality_descriptions().len(), 12);
1036 }
1037}