Skip to main content

faucet_cli/
transforms.rs

1//! Compile YAML/JSON transform declarations into `TransformStage` values.
2//!
3//! The built-in transforms are exposed via config — custom closure
4//! transforms require Rust code and are reserved for the library API.
5
6use 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/// Inline-config schema for the `flatten` transform.
20#[cfg(feature = "transforms")]
21#[derive(Debug, Deserialize, JsonSchema)]
22struct FlattenConfig {
23    /// Separator joining nested keys (default: `"__"`).
24    #[serde(default = "default_separator")]
25    separator: String,
26}
27
28#[cfg(feature = "transforms")]
29fn default_separator() -> String {
30    "__".to_owned()
31}
32
33/// Inline-config schema for the `rename_keys` transform.
34#[cfg(feature = "transforms")]
35#[derive(Debug, Deserialize, JsonSchema)]
36struct RenameKeysConfig {
37    /// Rust regex matched against every key.
38    pattern: String,
39    /// Replacement string. May reference capture groups (`$1`, `${name}`).
40    replacement: String,
41}
42
43#[cfg(feature = "transforms")]
44#[derive(Debug, Deserialize, JsonSchema)]
45struct FieldsConfig {
46    /// Top-level field names to act on.
47    fields: Vec<String>,
48}
49
50#[cfg(feature = "transforms")]
51#[derive(Debug, Deserialize, JsonSchema)]
52struct SetConfig {
53    /// Map of field name → constant value to set on every record.
54    values: serde_json::Map<String, Value>,
55}
56
57#[cfg(feature = "transforms")]
58#[derive(Debug, Deserialize, JsonSchema)]
59struct RenameFieldConfig {
60    /// Map of old field name → new field name.
61    fields: HashMap<String, String>,
62}
63
64#[cfg(feature = "transforms")]
65#[derive(Debug, Deserialize, JsonSchema)]
66struct CastConfig {
67    /// Map of field name → target type.
68    fields: HashMap<String, CastType>,
69    /// What to do when a value cannot be cast. Default: `error`.
70    #[serde(default)]
71    on_error: CastOnError,
72}
73
74#[cfg(feature = "transforms")]
75#[derive(Debug, Deserialize, JsonSchema)]
76struct RedactConfig {
77    /// Top-level field names to overwrite with `mask`.
78    fields: Vec<String>,
79    /// Replacement value. Default: the string `"***"`.
80    #[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    /// String-valued fields to re-case.
93    fields: Vec<String>,
94    /// Casing convention to apply to each listed field.
95    mode: ValueCaseMode,
96}
97
98#[cfg(feature = "transforms")]
99#[derive(Debug, Deserialize, JsonSchema)]
100struct SpellSymbolsConfig {
101    /// Extra symbol → word overrides layered on top of the built-in map.
102    #[serde(default)]
103    extra: HashMap<String, String>,
104    /// String inserted between expanded words. Default: a single space.
105    #[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    /// Output convention for every key in the record.
118    mode: KeyCaseMode,
119}
120
121#[cfg(feature = "transform-filter")]
122#[derive(Debug, Deserialize, JsonSchema)]
123struct FilterConfig {
124    /// JSONPath subset: bare key, dot path, or bracketed string key.
125    path: String,
126    /// One of `eq`, `ne`, `exists`, `in`, `not_in`.
127    op: faucet_core::FilterOp,
128    /// Required for `eq`/`ne`/`in`/`not_in`. For `in`/`not_in`, must be an array.
129    #[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    /// JSONPath subset: bare key, dot path, or bracketed string key.
137    path: String,
138    /// Prefix prepended to object-element fields. Defaults to the last
139    /// segment of `path`. Empty string = pure LATERAL FLATTEN (no prefix).
140    #[serde(default, skip_serializing_if = "Option::is_none")]
141    prefix: Option<String>,
142    /// Separator between prefix and element field key. Default `"_"`.
143    #[serde(default = "default_explode_separator_cli")]
144    separator: String,
145    /// `passthrough` (default), `drop`, or `error` when path doesn't yield a
146    /// non-empty array.
147    #[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    /// Envelope field holding the operation code. Default `"op"`.
160    #[serde(default = "cdc_op_field")]
161    op_field: String,
162    /// Envelope field holding the post-image row. Default `"after"`.
163    #[serde(default = "cdc_after_field")]
164    after_field: String,
165    /// Envelope field holding the pre-image row. Default `"before"`.
166    #[serde(default = "cdc_before_field")]
167    before_field: String,
168    /// Fallback key field for deletes when `before` is absent. Default `"document_key"`.
169    #[serde(default = "cdc_key_field")]
170    key_field: String,
171    /// Field stamped onto every emitted row with the op value. Default `"__op"`.
172    #[serde(default = "cdc_marker_field")]
173    marker_field: String,
174    /// Op values that mean delete. Default `["d", "delete"]`.
175    #[serde(default = "cdc_delete_ops")]
176    delete_ops: Vec<String>,
177    /// Op values that cause the record to be dropped (1→0). Default `["ddl", "truncate"]`.
178    #[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
211/// One row in the transform registry — the single source of truth for every
212/// built-in transform's kind, one-line description, JSON Schema, and
213/// `TransformSpec → TransformStage` decoder. `compile_one`,
214/// `transform_descriptions`, and `transform_schema` all read from this list
215/// so adding a new transform means appending one entry (no parallel match
216/// arms to keep in sync).
217struct TransformDef {
218    kind: &'static str,
219    description: &'static str,
220    schema_fn: fn() -> Value,
221    compile_fn: fn(&str, Value) -> CliResult<TransformStage>,
222}
223
224/// Every transform compiled into this build, in display order.
225///
226/// Non-capturing closures coerce to the `fn` pointers held by `TransformDef`,
227/// so each row stays a single self-contained record next to its sibling
228/// entries.
229fn 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                    // Re-use stage's compile-time validation so error messages match.
369                    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
461/// Compile a list of [`TransformSpec`]s into [`TransformStage`]s in the
462/// declared order. Most built-ins compile to a [`TransformStage::Map`];
463/// richer stages (e.g. `filter`, future fan-outs) compile to other
464/// variants. Unknown or malformed entries surface as a `CliError`.
465pub 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
480/// One-line summary of every transform compiled into this build. Used by
481/// `faucet list`.
482pub fn transform_descriptions() -> Vec<(&'static str, &'static str)> {
483    registry()
484        .into_iter()
485        .map(|t| (t.kind, t.description))
486        .collect()
487}
488
489/// Names of every transform compiled into this build.
490pub fn available_transforms() -> Vec<&'static str> {
491    registry().into_iter().map(|t| t.kind).collect()
492}
493
494// Keep in sync with faucet_core::{RecordCheck, BatchCheck} — one entry per check variant.
495/// One-line descriptions of the available quality checks, for `faucet list`.
496/// The `json_schema` entry only appears when the `quality-jsonschema` feature
497/// is enabled, mirroring `faucet schema quality` so `list` and `schema` agree.
498#[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
527/// Return the JSON Schema for the named transform's config. Mirrors
528/// `registry::source_schema` / `sink_schema` so `faucet schema transform <name>`
529/// reads symmetrically with the connector variants.
530pub 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        // Removed in favour of `keys_case { mode: snake }`.
639        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        // descriptions and available_transforms must never drift — `faucet list`
811        // and the `UnknownTransform` "Available:" line both read from this.
812        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        // Both accept `{ fields: Vec<String> }` — the schema is the same object,
838        // just titled `FieldsConfig`.
839        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        // 8 always-on per-record checks + 4 per-batch checks = 12; the
1028        // `json_schema` per-record check is only present (→ 13) when the
1029        // `quality-jsonschema` feature is enabled, matching `faucet schema
1030        // quality`. If you add a RecordCheck/BatchCheck variant in faucet-core,
1031        // add its description here too.
1032        #[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}