Skip to main content

faucet_lineage/
column.rs

1//! Deterministic column-level lineage from a declarative transform chain.
2//!
3//! v1 supports the field-preserving / explicit-mapping transforms exactly.
4//! Any structure-changing or key-rewriting transform (`flatten`, `explode`,
5//! `keys_case`, `rename_keys`) and any `Custom`/unknown transform is
6//! [`ColumnOp::Opaque`]; if the chain contains one, [`derive()`] returns `None`
7//! and **no** column-lineage facet is emitted (never fabricated).
8
9use indexmap::IndexMap;
10use std::collections::BTreeSet;
11
12/// A lineage-relevant view of one transform stage. The CLI maps its resolved
13/// transform specs onto these.
14#[derive(Debug, Clone)]
15pub enum ColumnOp {
16    /// Keys unchanged (cast / redact / value_case / spell_symbols / filter).
17    Identity,
18    /// Explicit 1:1 key renames (old → new).
19    Rename(Vec<(String, String)>),
20    /// Keep only these top-level keys.
21    Select(Vec<String>),
22    /// Remove these top-level keys.
23    Drop(Vec<String>),
24    /// Add these new keys as literals (no upstream edge).
25    Set(Vec<String>),
26    /// Structure-changing / key-rewriting / unknown — lineage not derivable.
27    Opaque,
28}
29
30/// Output field → ordered list of originating input field names. A field with
31/// an empty source list is a literal (`set`).
32#[derive(Debug, Clone, PartialEq)]
33pub struct ColumnLineage {
34    pub edges: IndexMap<String, Vec<String>>,
35}
36
37/// Fold the transform chain over the input field set. Returns `None` if any op
38/// is [`ColumnOp::Opaque`].
39pub fn derive(input_fields: &[String], ops: &[ColumnOp]) -> Option<ColumnLineage> {
40    if ops.iter().any(|o| matches!(o, ColumnOp::Opaque)) {
41        return None;
42    }
43    // Working map: current field name → set of source fields.
44    let mut cur: IndexMap<String, BTreeSet<String>> = IndexMap::new();
45    for f in input_fields {
46        cur.insert(f.clone(), BTreeSet::from([f.clone()]));
47    }
48    for op in ops {
49        match op {
50            ColumnOp::Identity => {}
51            ColumnOp::Rename(pairs) => {
52                for (from, to) in pairs {
53                    if let Some(sources) = cur.shift_remove(from) {
54                        cur.insert(to.clone(), sources);
55                    }
56                }
57            }
58            ColumnOp::Select(keep) => {
59                let keep: BTreeSet<&String> = keep.iter().collect();
60                cur.retain(|k, _| keep.contains(k));
61            }
62            ColumnOp::Drop(remove) => {
63                let remove: BTreeSet<&String> = remove.iter().collect();
64                cur.retain(|k, _| !remove.contains(k));
65            }
66            ColumnOp::Set(added) => {
67                for f in added {
68                    cur.insert(f.clone(), BTreeSet::new());
69                }
70            }
71            ColumnOp::Opaque => unreachable!("guarded above"),
72        }
73    }
74    let edges = cur
75        .into_iter()
76        .map(|(k, v)| (k, v.into_iter().collect()))
77        .collect();
78    Some(ColumnLineage { edges })
79}
80
81#[cfg(test)]
82mod tests {
83    use super::*;
84
85    fn inputs() -> Vec<String> {
86        vec!["id".into(), "name".into(), "email".into()]
87    }
88
89    #[test]
90    fn identity_chain_maps_each_field_to_itself() {
91        let cl = derive(&inputs(), &[ColumnOp::Identity]).unwrap();
92        assert_eq!(cl.edges.get("id").unwrap(), &vec!["id".to_string()]);
93        assert_eq!(cl.edges.len(), 3);
94    }
95
96    #[test]
97    fn rename_field_rekeys_and_preserves_source() {
98        let ops = [ColumnOp::Rename(vec![("email".into(), "contact".into())])];
99        let cl = derive(&inputs(), &ops).unwrap();
100        assert_eq!(cl.edges.get("contact").unwrap(), &vec!["email".to_string()]);
101        assert!(!cl.edges.contains_key("email"));
102    }
103
104    #[test]
105    fn select_retains_only_listed() {
106        let cl = derive(&inputs(), &[ColumnOp::Select(vec!["id".into()])]).unwrap();
107        assert_eq!(
108            cl.edges.keys().cloned().collect::<Vec<_>>(),
109            vec!["id".to_string()]
110        );
111    }
112
113    #[test]
114    fn drop_removes_listed() {
115        let cl = derive(&inputs(), &[ColumnOp::Drop(vec!["email".into()])]).unwrap();
116        assert!(!cl.edges.contains_key("email"));
117        assert_eq!(cl.edges.len(), 2);
118    }
119
120    #[test]
121    fn set_adds_field_with_no_upstream() {
122        let cl = derive(&inputs(), &[ColumnOp::Set(vec!["created".into()])]).unwrap();
123        assert!(cl.edges.get("created").unwrap().is_empty());
124    }
125
126    #[test]
127    fn opaque_op_yields_no_lineage() {
128        assert!(derive(&inputs(), &[ColumnOp::Opaque]).is_none());
129        assert!(derive(&inputs(), &[ColumnOp::Identity, ColumnOp::Opaque]).is_none());
130    }
131}