use crate::config::TransformSpec;
use crate::error::{CliError, CliResult};
#[cfg(feature = "transforms")]
use faucet_core::{CastOnError, CastType, KeyCaseMode, ValueCaseMode};
#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
use faucet_core::{JsonSchema, schema_for};
use faucet_core::{RecordTransform, TransformStage};
#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
use serde::Deserialize;
use serde_json::Value;
#[cfg(feature = "transforms")]
use std::collections::HashMap;
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct FlattenConfig {
#[serde(default = "default_separator")]
separator: String,
}
#[cfg(feature = "transforms")]
fn default_separator() -> String {
"__".to_owned()
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct RenameKeysConfig {
pattern: String,
replacement: String,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct FieldsConfig {
fields: Vec<String>,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct SetConfig {
values: serde_json::Map<String, Value>,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct RenameFieldConfig {
fields: HashMap<String, String>,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct CastConfig {
fields: HashMap<String, CastType>,
#[serde(default)]
on_error: CastOnError,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct RedactConfig {
fields: Vec<String>,
#[serde(default = "default_mask")]
mask: Value,
}
#[cfg(feature = "transforms")]
fn default_mask() -> Value {
Value::String("***".to_owned())
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct ValueCaseConfig {
fields: Vec<String>,
mode: ValueCaseMode,
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct SpellSymbolsConfig {
#[serde(default)]
extra: HashMap<String, String>,
#[serde(default = "default_spell_separator")]
separator: String,
}
#[cfg(feature = "transforms")]
fn default_spell_separator() -> String {
" ".to_owned()
}
#[cfg(feature = "transforms")]
#[derive(Debug, Deserialize, JsonSchema)]
struct KeysCaseConfig {
mode: KeyCaseMode,
}
#[cfg(feature = "transform-filter")]
#[derive(Debug, Deserialize, JsonSchema)]
struct FilterConfig {
path: String,
op: faucet_core::FilterOp,
#[serde(default, skip_serializing_if = "Option::is_none")]
value: Option<Value>,
}
#[cfg(feature = "transform-explode")]
#[derive(Debug, Deserialize, JsonSchema)]
struct ExplodeConfig {
path: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
prefix: Option<String>,
#[serde(default = "default_explode_separator_cli")]
separator: String,
#[serde(default)]
on_missing: faucet_core::OnMissing,
}
#[cfg(feature = "transform-explode")]
fn default_explode_separator_cli() -> String {
"_".to_owned()
}
#[cfg(feature = "transform-cdc-unwrap")]
#[derive(Debug, Deserialize, JsonSchema)]
struct CdcUnwrapConfig {
#[serde(default = "cdc_op_field")]
op_field: String,
#[serde(default = "cdc_after_field")]
after_field: String,
#[serde(default = "cdc_before_field")]
before_field: String,
#[serde(default = "cdc_key_field")]
key_field: String,
#[serde(default = "cdc_marker_field")]
marker_field: String,
#[serde(default = "cdc_delete_ops")]
delete_ops: Vec<String>,
#[serde(default = "cdc_drop_ops")]
drop_ops: Vec<String>,
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_op_field() -> String {
"op".into()
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_after_field() -> String {
"after".into()
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_before_field() -> String {
"before".into()
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_key_field() -> String {
"document_key".into()
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_marker_field() -> String {
"__op".into()
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_delete_ops() -> Vec<String> {
vec!["d".into(), "delete".into()]
}
#[cfg(feature = "transform-cdc-unwrap")]
fn cdc_drop_ops() -> Vec<String> {
vec!["ddl".into(), "truncate".into()]
}
struct TransformDef {
kind: &'static str,
description: &'static str,
schema_fn: fn() -> Value,
compile_fn: fn(&str, Value) -> CliResult<TransformStage>,
}
fn registry() -> Vec<TransformDef> {
#[allow(unused_mut)]
let mut defs: Vec<TransformDef> = Vec::new();
#[cfg(feature = "transforms")]
{
defs.extend(vec![
TransformDef {
kind: "flatten",
description: "Flatten nested objects into a single level (configurable separator).",
schema_fn: || schema::<FlattenConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<FlattenConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Flatten {
separator: cfg.separator,
}))
},
},
TransformDef {
kind: "rename_keys",
description: "Rewrite every key via a regex pattern + replacement.",
schema_fn: || schema::<RenameKeysConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<RenameKeysConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::RenameKeys {
pattern: cfg.pattern,
replacement: cfg.replacement,
}))
},
},
TransformDef {
kind: "keys_case",
description: "Re-case every key (snake / camel / pascal / kebab / screaming_snake).",
schema_fn: || schema::<KeysCaseConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<KeysCaseConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::KeysCase {
mode: cfg.mode,
}))
},
},
TransformDef {
kind: "select",
description: "Keep only the listed top-level fields; drop the rest.",
schema_fn: || schema::<FieldsConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<FieldsConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Select {
fields: cfg.fields,
}))
},
},
TransformDef {
kind: "drop",
description: "Remove the listed top-level fields.",
schema_fn: || schema::<FieldsConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<FieldsConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Drop {
fields: cfg.fields,
}))
},
},
TransformDef {
kind: "set",
description: "Set named fields to constant values on every record.",
schema_fn: || schema::<SetConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<SetConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Set {
values: cfg.values,
}))
},
},
TransformDef {
kind: "rename_field",
description: "Rename specific top-level fields by name.",
schema_fn: || schema::<RenameFieldConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<RenameFieldConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::RenameField {
fields: cfg.fields,
}))
},
},
TransformDef {
kind: "cast",
description: "Coerce named fields to int / float / bool / string / timestamp.",
schema_fn: || schema::<CastConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<CastConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Cast {
fields: cfg.fields,
on_error: cfg.on_error,
}))
},
},
TransformDef {
kind: "redact",
description: "Overwrite the listed fields with a mask value (default `***`).",
schema_fn: || schema::<RedactConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<RedactConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::Redact {
fields: cfg.fields,
mask: cfg.mask,
}))
},
},
TransformDef {
kind: "value_case",
description: "Lowercase, uppercase, or trim the value of named string fields.",
schema_fn: || schema::<ValueCaseConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<ValueCaseConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::ValueCase {
fields: cfg.fields,
mode: cfg.mode,
}))
},
},
TransformDef {
kind: "spell_symbols",
description: "Replace punctuation/symbols in string values with their spelled-out words.",
schema_fn: || schema::<SpellSymbolsConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<SpellSymbolsConfig>(kind, config)?;
Ok(TransformStage::Map(RecordTransform::SpellSymbols {
extra: cfg.extra,
separator: cfg.separator,
}))
},
},
#[cfg(feature = "transform-filter")]
TransformDef {
kind: "filter",
description: "Keep records where a JSONPath predicate is true.",
schema_fn: || schema::<FilterConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<FilterConfig>(kind, config)?;
let stage = TransformStage::Filter(faucet_core::FilterSpec {
path: cfg.path,
op: cfg.op,
value: cfg.value,
});
faucet_core::compile_stage(&stage).map_err(|e| match e {
faucet_core::FaucetError::Transform(msg) => CliError::InvalidTransform {
name: kind.to_owned(),
message: msg,
},
other => CliError::InvalidTransform {
name: kind.to_owned(),
message: format!("{other}"),
},
})?;
Ok(stage)
},
},
#[cfg(feature = "transform-explode")]
TransformDef {
kind: "explode",
description: "Expand an array field into one record per element.",
schema_fn: || schema::<ExplodeConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<ExplodeConfig>(kind, config)?;
let stage = TransformStage::Explode(faucet_core::ExplodeSpec {
path: cfg.path,
prefix: cfg.prefix,
separator: cfg.separator,
on_missing: cfg.on_missing,
});
faucet_core::compile_stage(&stage).map_err(|e| match e {
faucet_core::FaucetError::Transform(msg) => CliError::InvalidTransform {
name: kind.to_owned(),
message: msg,
},
other => CliError::InvalidTransform {
name: kind.to_owned(),
message: format!("{other}"),
},
})?;
Ok(stage)
},
},
]);
}
#[cfg(feature = "transform-sql")]
{
defs.push(TransformDef {
kind: "sql",
description: "Run DuckDB SQL over the whole page; records are the `batch` relation.",
schema_fn: || schema_sql(),
compile_fn: |kind, config| {
let cfg: faucet_transform_sql::SqlTransformConfig = decode_sql(kind, config)?;
let transform = faucet_transform_sql::SqlTransform::compile(&cfg).map_err(|e| {
let message = match &e {
faucet_core::FaucetError::Transform(m)
| faucet_core::FaucetError::Config(m) => m.clone(),
other => format!("{other}"),
};
CliError::InvalidTransform {
name: kind.to_owned(),
message,
}
})?;
Ok(transform.into_page_stage())
},
});
}
#[cfg(feature = "transform-cdc-unwrap")]
{
defs.push(TransformDef {
kind: "cdc_unwrap",
description: "Normalize a CDC envelope into a flat row + delete marker (for upsert sinks).",
schema_fn: || schema::<CdcUnwrapConfig>(),
compile_fn: |kind, config| {
let cfg = decode::<CdcUnwrapConfig>(kind, config)?;
Ok(faucet_core::TransformStage::CdcUnwrap(faucet_core::CdcUnwrapSpec {
op_field: cfg.op_field,
after_field: cfg.after_field,
before_field: cfg.before_field,
key_field: cfg.key_field,
marker_field: cfg.marker_field,
delete_ops: cfg.delete_ops,
drop_ops: cfg.drop_ops,
}))
},
});
}
defs
}
pub fn compile_transforms(specs: &[TransformSpec]) -> CliResult<Vec<TransformStage>> {
let mut out = Vec::with_capacity(specs.len());
for s in specs {
out.push(compile_one(s)?);
}
Ok(out)
}
fn compile_one(spec: &TransformSpec) -> CliResult<TransformStage> {
match registry().into_iter().find(|t| t.kind == spec.kind) {
Some(def) => (def.compile_fn)(&spec.kind, spec.config.clone()),
None => Err(unknown_transform(&spec.kind)),
}
}
pub fn transform_descriptions() -> Vec<(&'static str, &'static str)> {
registry()
.into_iter()
.map(|t| (t.kind, t.description))
.collect()
}
pub fn available_transforms() -> Vec<&'static str> {
registry().into_iter().map(|t| t.kind).collect()
}
#[cfg(feature = "quality")]
pub fn quality_descriptions() -> Vec<(&'static str, &'static str)> {
let mut checks = vec![
("not_null", "field present and non-null"),
("not_empty", "string non-empty after trim"),
("regex_match", "string matches a regex"),
("value_in_set", "value is in an allowed set"),
("not_in_set", "value is not in a forbidden set"),
("compare", "numeric/scalar comparison (gt/gte/lt/lte/eq/ne)"),
("type_is", "value is of an expected JSON type"),
("string_length", "string length within [min,max]"),
];
#[cfg(feature = "quality-jsonschema")]
checks.push((
"json_schema",
"record validates against a JSON Schema (feature-gated)",
));
checks.extend([
("row_count", "batch row count within [min,max]"),
("null_rate", "batch null rate of a field <= max"),
("unique", "composite key unique within the batch"),
(
"distinct_count",
"distinct values of a field within [min,max]",
),
]);
checks
}
pub fn transform_schema(name: &str) -> CliResult<Value> {
registry()
.into_iter()
.find(|t| t.kind == name)
.map(|t| (t.schema_fn)())
.ok_or_else(|| unknown_transform(name))
}
fn unknown_transform(name: &str) -> CliError {
let available = available_transforms();
CliError::UnknownTransform {
name: name.to_owned(),
available: if available.is_empty() {
"(none — rebuild faucet-cli with the `transforms` feature enabled)".to_owned()
} else {
available.join(", ")
},
}
}
#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
fn schema<T: JsonSchema>() -> Value {
serde_json::to_value(schema_for!(T)).unwrap_or_else(|_| serde_json::json!({"type": "object"}))
}
#[cfg(any(feature = "transforms", feature = "transform-cdc-unwrap"))]
fn decode<T: serde::de::DeserializeOwned>(name: &str, config: Value) -> CliResult<T> {
serde_json::from_value(config).map_err(|e| CliError::InvalidTransform {
name: name.to_owned(),
message: e.to_string(),
})
}
#[cfg(feature = "transform-sql")]
fn schema_sql() -> Value {
serde_json::to_value(faucet_core::schema_for!(
faucet_transform_sql::SqlTransformConfig
))
.unwrap_or(Value::Null)
}
#[cfg(feature = "transform-sql")]
fn decode_sql(kind: &str, config: Value) -> CliResult<faucet_transform_sql::SqlTransformConfig> {
serde_json::from_value(config).map_err(|e| CliError::InvalidTransform {
name: kind.to_owned(),
message: e.to_string(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn empty_list_compiles_to_empty() {
let out = compile_transforms(&[]).unwrap();
assert!(out.is_empty());
}
#[cfg(feature = "transforms")]
#[test]
fn compiles_keys_case_and_flatten() {
let specs = vec![
TransformSpec {
kind: "keys_case".into(),
config: json!({"mode": "snake"}),
},
TransformSpec {
kind: "flatten".into(),
config: json!({"separator": "."}),
},
];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 2);
}
#[cfg(feature = "transforms")]
#[test]
fn keys_case_rejects_unknown_mode() {
let specs = vec![TransformSpec {
kind: "keys_case".into(),
config: json!({"mode": "spongebob"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "keys_case"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn keys_case_requires_mode() {
let specs = vec![TransformSpec {
kind: "keys_case".into(),
config: json!({}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "keys_case"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn snake_case_kind_is_no_longer_recognized() {
let specs = vec![TransformSpec {
kind: "snake_case".into(),
config: json!({}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::UnknownTransform { name, .. } => assert_eq!(name, "snake_case"),
other => panic!("expected UnknownTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn rename_keys_requires_pattern_and_replacement() {
let specs = vec![TransformSpec {
kind: "rename_keys".into(),
config: json!({"pattern": "^_"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "rename_keys"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[test]
fn unknown_transform_errors() {
let specs = vec![TransformSpec {
kind: "make_uppercase".into(),
config: json!({}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::UnknownTransform { name, .. } => assert_eq!(name, "make_uppercase"),
other => panic!("expected UnknownTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn compiles_select_and_drop() {
let specs = vec![
TransformSpec {
kind: "select".into(),
config: json!({"fields": ["id", "name"]}),
},
TransformSpec {
kind: "drop".into(),
config: json!({"fields": ["secret"]}),
},
];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 2);
}
#[cfg(feature = "transforms")]
#[test]
fn compiles_set_with_object_values() {
let specs = vec![TransformSpec {
kind: "set".into(),
config: json!({"values": {"_source": "api", "version": 1}}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
}
#[cfg(feature = "transforms")]
#[test]
fn compiles_rename_field() {
let specs = vec![TransformSpec {
kind: "rename_field".into(),
config: json!({"fields": {"old": "new"}}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
}
#[cfg(feature = "transforms")]
#[test]
fn compiles_cast_with_default_on_error() {
let specs = vec![TransformSpec {
kind: "cast".into(),
config: json!({"fields": {"age": "int", "price": "float"}}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
}
#[cfg(feature = "transforms")]
#[test]
fn cast_rejects_unknown_target_type() {
let specs = vec![TransformSpec {
kind: "cast".into(),
config: json!({"fields": {"x": "uuid"}}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "cast"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn cast_rejects_unknown_on_error_mode() {
let specs = vec![TransformSpec {
kind: "cast".into(),
config: json!({"fields": {"x": "int"}, "on_error": "explode"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "cast"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn redact_uses_default_mask_when_omitted() {
let specs = vec![TransformSpec {
kind: "redact".into(),
config: json!({"fields": ["ssn"]}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
}
#[cfg(feature = "transforms")]
#[test]
fn value_case_requires_mode() {
let specs = vec![TransformSpec {
kind: "value_case".into(),
config: json!({"fields": ["email"]}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "value_case"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transforms")]
#[test]
fn available_transforms_lists_every_kind() {
let names = available_transforms();
for expected in [
"flatten",
"rename_keys",
"keys_case",
"select",
"drop",
"set",
"rename_field",
"cast",
"redact",
"value_case",
"spell_symbols",
"filter",
"explode",
] {
assert!(names.contains(&expected), "missing {expected}");
}
assert!(
!names.contains(&"snake_case"),
"snake_case must be removed in favour of keys_case"
);
}
#[cfg(feature = "transforms")]
#[test]
fn transform_descriptions_covers_every_compiled_kind() {
let names = available_transforms();
let desc_names: Vec<&'static str> = transform_descriptions()
.into_iter()
.map(|(n, _)| n)
.collect();
assert_eq!(names, desc_names);
for (_, desc) in transform_descriptions() {
assert!(!desc.is_empty(), "every transform needs a description");
}
}
#[cfg(feature = "transforms")]
#[test]
fn transform_schema_returns_object_for_every_kind() {
for name in available_transforms() {
let schema = transform_schema(name).unwrap_or_else(|e| {
panic!("schema lookup failed for {name}: {e}");
});
assert!(schema.is_object(), "schema for {name} must be an object");
}
}
#[cfg(feature = "transforms")]
#[test]
fn transform_schema_select_and_drop_share_shape() {
let select = transform_schema("select").unwrap();
let drop = transform_schema("drop").unwrap();
assert_eq!(select, drop);
}
#[test]
fn transform_schema_unknown_errors_with_available_list() {
let err = transform_schema("make_uppercase").unwrap_err();
match err {
CliError::UnknownTransform { name, available } => {
assert_eq!(name, "make_uppercase");
#[cfg(feature = "transforms")]
assert!(available.contains("flatten"), "{available}");
#[cfg(not(feature = "transforms"))]
assert!(available.contains("rebuild"), "{available}");
}
other => panic!("expected UnknownTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-filter")]
#[test]
fn compiles_filter_eq() {
let specs = vec![TransformSpec {
kind: "filter".into(),
config: json!({"path": "status", "op": "eq", "value": "active"}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
assert!(matches!(out[0], TransformStage::Filter(_)));
}
#[cfg(feature = "transform-filter")]
#[test]
fn filter_rejects_in_with_non_array_value() {
let specs = vec![TransformSpec {
kind: "filter".into(),
config: json!({"path": "v", "op": "in", "value": "scalar"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, message } => {
assert_eq!(name, "filter");
assert!(message.contains("requires an array"), "{message}");
}
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-filter")]
#[test]
fn filter_rejects_exists_with_value() {
let specs = vec![TransformSpec {
kind: "filter".into(),
config: json!({"path": "v", "op": "exists", "value": "x"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "filter"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-filter")]
#[test]
fn filter_rejects_bad_path() {
let specs = vec![TransformSpec {
kind: "filter".into(),
config: json!({"path": "$..items", "op": "exists"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "filter"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-explode")]
#[test]
fn compiles_explode_with_defaults() {
let specs = vec![TransformSpec {
kind: "explode".into(),
config: json!({"path": "items"}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
assert!(matches!(out[0], TransformStage::Explode(_)));
}
#[cfg(feature = "transform-explode")]
#[test]
fn compiles_explode_with_custom_prefix_and_on_missing() {
let specs = vec![TransformSpec {
kind: "explode".into(),
config: json!({
"path": "items",
"prefix": "item",
"separator": "_",
"on_missing": "drop"
}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
}
#[cfg(feature = "transform-explode")]
#[test]
fn explode_rejects_bad_path() {
let specs = vec![TransformSpec {
kind: "explode".into(),
config: json!({"path": "$..items"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "explode"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-explode")]
#[test]
fn explode_rejects_invalid_on_missing() {
let specs = vec![TransformSpec {
kind: "explode".into(),
config: json!({"path": "items", "on_missing": "explode_harder"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "explode"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-sql")]
#[test]
fn compiles_sql_to_page_fn() {
let specs = vec![TransformSpec {
kind: "sql".into(),
config: json!({"query": "SELECT * FROM batch"}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
assert!(matches!(out[0], faucet_core::TransformStage::PageFn(_)));
}
#[cfg(feature = "transform-sql")]
#[test]
fn sql_bad_query_is_invalid_transform() {
let specs = vec![TransformSpec {
kind: "sql".into(),
config: json!({"query": "SELEKT bad"}),
}];
let err = compile_transforms(&specs).unwrap_err();
match err {
CliError::InvalidTransform { name, .. } => assert_eq!(name, "sql"),
other => panic!("expected InvalidTransform, got {other:?}"),
}
}
#[cfg(feature = "transform-sql")]
#[test]
fn sql_schema_and_listing_present() {
assert!(transform_schema("sql").is_ok());
assert!(available_transforms().contains(&"sql"));
}
#[cfg(feature = "transform-cdc-unwrap")]
#[test]
fn compiles_cdc_unwrap_with_defaults() {
let specs = vec![TransformSpec {
kind: "cdc_unwrap".into(),
config: json!({}),
}];
let out = compile_transforms(&specs).unwrap();
assert_eq!(out.len(), 1);
assert!(matches!(out[0], faucet_core::TransformStage::CdcUnwrap(_)));
}
#[cfg(feature = "transform-cdc-unwrap")]
#[test]
fn cdc_unwrap_schema_and_listing_present() {
assert!(transform_schema("cdc_unwrap").is_ok());
assert!(available_transforms().contains(&"cdc_unwrap"));
}
#[cfg(feature = "quality")]
#[test]
fn quality_descriptions_has_one_entry_per_check() {
#[cfg(feature = "quality-jsonschema")]
assert_eq!(quality_descriptions().len(), 13);
#[cfg(not(feature = "quality-jsonschema"))]
assert_eq!(quality_descriptions().len(), 12);
}
}