From 1f6370af432ab8b2c25cf5095f97f7fcaf17d338 Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Thu, 24 Sep 2026 15:39:08 +0200 Subject: [PATCH 1/8] fix(core): stop long operator chains from tripping the depth guard The parser builds a chain of binary operators left-deep: `a || b || c` is `(a || b) || c`. The expression walkers recursed into both operands and counted every operator as one level of nesting, so a flat chain of more than about 100 operators passed MAX_RECURSION_DEPTH. The walkers then gave up at the guard, which sits at the leftmost end of the chain, and the statement was reported as APPROXIMATE_LINEAGE. Generated surrogate keys reach that length quickly, because every field is wrapped and joined with a separator: SELECT MD5( COALESCE(CAST(field_000 AS VARCHAR), '') || '-' || COALESCE(CAST(field_001 AS VARCHAR), '') || '-' || ... COALESCE(CAST(field_119 AS VARCHAR), '') ) AS row_key FROM events Over 120 fields, row_key was derived from the last 49 only: field_000 to field_070 were missing from its lineage. A key over 50 fields already loses its first operands. Walk the left spine of a BinaryOp chain in a loop and visit its operands left to right, each one level below the chain. The length of a chain no longer counts as depth, while nesting in right operands, function arguments and parentheses still does, so the guard keeps protecting the stack. The four walkers that share MAX_RECURSION_DEPTH use it: visit_expression_for_subqueries, collect_column_refs, find_aggregate_function and collect_simple_identifiers. Operands are still visited in source order. Tests: the 120-field key above keeps every field and raises no APPROXIMATE_LINEAGE; unit tests pin the operand order and check that nesting on the right still trips the guard. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../flowscope-core/src/analyzer/expression.rs | 111 ++++++++++++++++-- crates/flowscope-core/tests/lineage_engine.rs | 39 ++++++ 2 files changed, 138 insertions(+), 12 deletions(-) diff --git a/crates/flowscope-core/src/analyzer/expression.rs b/crates/flowscope-core/src/analyzer/expression.rs index 123d858c..77da1955 100644 --- a/crates/flowscope-core/src/analyzer/expression.rs +++ b/crates/flowscope-core/src/analyzer/expression.rs @@ -28,6 +28,26 @@ use tracing::debug; /// on maliciously crafted or deeply nested SQL expressions. pub(super) const MAX_RECURSION_DEPTH: usize = 100; +/// Returns the operands of a chain of binary operators, left to right. +/// +/// The parser builds `a || b || c` as `(a || b) || c`, so each operator of a +/// flat chain adds one level on the left. Walking that spine in a loop keeps +/// the length of a chain from counting against `MAX_RECURSION_DEPTH`, which +/// guards real nesting: otherwise a surrogate key concatenating fifty columns +/// and their separators loses its leftmost operands. Right operands are still +/// visited by recursion, so nesting on that side keeps counting. +fn binary_op_operands(expr: &Expr) -> Vec<&Expr> { + let mut operands = Vec::new(); + let mut current = expr; + while let Expr::BinaryOp { left, right, .. } = current { + operands.push(right.as_ref()); + current = left; + } + operands.push(current); + operands.reverse(); + operands +} + /// Analyzes SQL expressions to extract column references, detect aggregations, /// and capture filter predicates. /// @@ -114,9 +134,10 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { self.analyzer.analyze_query(self.ctx, subquery, None) } Expr::Exists { subquery, .. } => self.analyzer.analyze_query(self.ctx, subquery, None), - Expr::BinaryOp { left, right, .. } => { - self.visit_expression_for_subqueries(left, next_depth); - self.visit_expression_for_subqueries(right, next_depth); + Expr::BinaryOp { .. } => { + for operand in binary_op_operands(expr) { + self.visit_expression_for_subqueries(operand, next_depth); + } } Expr::UnaryOp { expr, .. } => self.visit_expression_for_subqueries(expr, next_depth), Expr::Nested(expr) => self.visit_expression_for_subqueries(expr, next_depth), @@ -230,9 +251,10 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { column, }); } - Expr::BinaryOp { left, right, .. } => { - depth_limited |= Self::collect_column_refs(left, refs, dialect, next_depth); - depth_limited |= Self::collect_column_refs(right, refs, dialect, next_depth); + Expr::BinaryOp { .. } => { + for operand in binary_op_operands(expr) { + depth_limited |= Self::collect_column_refs(operand, refs, dialect, next_depth); + } } Expr::UnaryOp { expr, .. } => { depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); @@ -424,9 +446,9 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { match expr { Expr::Function(func) => self.check_function_for_aggregate(func, next_depth), - Expr::BinaryOp { left, right, .. } => self - .find_aggregate_function(left, next_depth) - .or_else(|| self.find_aggregate_function(right, next_depth)), + Expr::BinaryOp { .. } => binary_op_operands(expr) + .into_iter() + .find_map(|operand| self.find_aggregate_function(operand, next_depth)), Expr::UnaryOp { expr, .. } | Expr::Nested(expr) | Expr::Cast { expr, .. } => { self.find_aggregate_function(expr, next_depth) } @@ -708,9 +730,12 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { } // Two expression patterns (left/right) - Expr::BinaryOp { left, right, .. } - | Expr::AnyOp { left, right, .. } - | Expr::AllOp { left, right, .. } => { + Expr::BinaryOp { .. } => { + for operand in binary_op_operands(expr) { + Self::collect_simple_identifiers(operand, identifiers, next_depth); + } + } + Expr::AnyOp { left, right, .. } | Expr::AllOp { left, right, .. } => { Self::collect_simple_identifiers(left, identifiers, next_depth); Self::collect_simple_identifiers(right, identifiers, next_depth); } @@ -852,4 +877,66 @@ mod tests { "no column refs should be recorded when guard triggers" ); } + + fn column(index: usize) -> Expr { + Expr::Identifier(ast::Ident::new(format!("c{index}"))) + } + + fn concat(left: Expr, right: Expr) -> Expr { + Expr::BinaryOp { + left: Box::new(left), + op: ast::BinaryOperator::StringConcat, + right: Box::new(right), + } + } + + /// `c0 || c1 || ... || c{len - 1}`, left-deep as the parser builds it. + fn flat_concat_chain(len: usize) -> Expr { + (1..len).fold(column(0), |chain, index| concat(chain, column(index))) + } + + #[test] + fn collect_column_refs_walks_long_operator_chains() { + let len = MAX_RECURSION_DEPTH * 3; + let mut refs = Vec::new(); + let hit = ExpressionAnalyzer::collect_column_refs( + &flat_concat_chain(len), + &mut refs, + Dialect::Generic, + 0, + ); + + assert!(!hit, "a flat operator chain should not trigger the guard"); + let columns: Vec = refs.into_iter().map(|col_ref| col_ref.column).collect(); + let expected: Vec = (0..len).map(|index| format!("c{index}")).collect(); + assert_eq!( + columns, expected, + "operands should be collected left to right" + ); + } + + #[test] + fn collect_column_refs_still_limits_nesting_on_the_right() { + let expr = (0..=MAX_RECURSION_DEPTH) + .rev() + .fold(column(MAX_RECURSION_DEPTH + 1), |nested, index| { + concat(column(index), nested) + }); + let mut refs = Vec::new(); + let hit = ExpressionAnalyzer::collect_column_refs(&expr, &mut refs, Dialect::Generic, 0); + + assert!(hit, "nesting past the limit should still trigger the guard"); + } + + #[test] + fn extract_simple_identifiers_walks_long_operator_chains() { + let len = MAX_RECURSION_DEPTH * 3; + let identifiers = ExpressionAnalyzer::extract_simple_identifiers(&flat_concat_chain(len)); + + assert_eq!( + identifiers.len(), + len, + "every operand of a flat chain should be collected" + ); + } } diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index 645d3ec9..ced3a2d0 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -7258,6 +7258,45 @@ fn deeply_nested_case_expressions() { assert!(col.is_some(), "nested CASE produces column"); } +#[test] +fn long_concatenation_keeps_every_operand() { + // A generated surrogate key joins every field with a separator, so 120 + // fields parse as a left-deep chain of 238 `||` operators. + let fields: Vec = (0..120).map(|i| format!("field_{i:03}")).collect(); + let key = fields + .iter() + .map(|field| format!("COALESCE(CAST({field} AS VARCHAR), '')")) + .collect::>() + .join(" || '-' || "); + let sql = format!("SELECT MD5({key}) AS row_key FROM events"); + + let result = run_analysis(&sql, Dialect::Generic, None); + assert!( + !issue_codes_list(&result).contains(&issue_codes::APPROXIMATE_LINEAGE.to_string()), + "a flat operator chain is not deep nesting, issues: {:?}", + result.issues + ); + + let stmt = first_statement(&result); + let row_key = find_column_node(&stmt, "row_key").expect("row_key column should exist"); + let sources: HashSet<&str> = stmt + .edges + .iter() + .filter(|edge| edge.edge_type == EdgeType::Derivation && edge.to == row_key.id) + .filter_map(|edge| stmt.nodes.iter().find(|node| node.id == edge.from)) + .map(|node| &*node.label) + .collect(); + let missing: Vec<&str> = fields + .iter() + .map(String::as_str) + .filter(|field| !sources.contains(field)) + .collect(); + assert!( + missing.is_empty(), + "every concatenated field should derive row_key, missing: {missing:?}" + ); +} + #[test] fn case_with_aggregate_in_condition() { // CASE with aggregate function in WHEN condition From 598e5629170efc63981433a9ca05dec23ced84c8 Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Thu, 24 Sep 2026 16:09:32 +0200 Subject: [PATCH 2/8] fix(core): read columns inside TRIM, SUBSTRING and other dedicated forms sqlparser gives TRIM, SUBSTRING, POSITION, CEIL and FLOOR, AT TIME ZONE, IS [NOT] DISTINCT FROM, the `col:path` and `col[i]` accessors, COLLATE, OVERLAY, SIMILAR TO, RLIKE, ANY and ALL, and a few more, their own `Expr` variants rather than function calls. collect_column_refs matched none of them and ended in a catch-all arm, so a column inside one of these forms was never a source, and no issue said so: SELECT o.id, TRIM(c.name) AS customer_name FROM orders o JOIN customers c ON o.customer_id = c.id customer_name came out with no source column. In `SUBSTR(code, 1, 2) || suffix` only suffix was read, and `WHERE TRIM(c.status) = 'active'` was not attached to customers. collect_simple_identifiers, which drives lateral column alias resolution, already walked SUBSTRING, CEIL and POSITION but not TRIM, so a key hashed from earlier aliases of the same SELECT list lost every source: SELECT MD5(UPPER(TRIM(CAST(order_id AS VARCHAR)))) AS order_key, MD5(UPPER(TRIM(CAST(customer_id AS VARCHAR)))) AS customer_key, MD5(CONCAT(UPPER(TRIM(CAST(order_key AS VARCHAR))), '||', UPPER(TRIM(CAST(customer_key AS VARCHAR))))) AS link_key FROM orders Make the match in collect_column_refs exhaustive. Every form descends into the operands it evaluates in the enclosing scope, including named arguments written `name => value`, the tested value of `x IN (SELECT ...)`, subscripts and bracket keys. Field names in a path step and lambda parameters are not read, and subqueries keep their own scope. The parser keeps `o.items[1]` as the root `o` followed by the steps `.items` and `[1]`, so the leading dot steps are folded back into the name: the column read is `o.items`, as without the subscript, and not a column `o` of orders. With no catch-all arm left, a variant added to sqlparser is a compile error rather than a column dropped from lineage. collect_simple_identifiers walks the same forms, so an alias hidden in one of them is replaced by its sources instead of being resolved as a column of a table in scope. The postgres_array_slicing snapshot changes accordingly: `a[:]`, `b[:1]`, `c[2:]` and `d[2:3]` now derive from the columns a, b, c and d rather than from the table node. Tests: the queries above, SUBSTR, CEIL, POSITION and `:` operands in Snowflake, the left operand of IN (subquery), a DuckDB lambda and a subscript of a qualified column in Postgres at the analysis level; a table of dedicated forms pins what collect_column_refs reads, and checks that extract_simple_identifiers sees the same bare names. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../flowscope-core/src/analyzer/expression.rs | 457 +++++++++++++++++- crates/flowscope-core/tests/lineage_engine.rs | 194 ++++++++ ..._snapshot_test@postgres_array_slicing.snap | 116 ++++- 3 files changed, 747 insertions(+), 20 deletions(-) diff --git a/crates/flowscope-core/src/analyzer/expression.rs b/crates/flowscope-core/src/analyzer/expression.rs index 77da1955..aff1f2e9 100644 --- a/crates/flowscope-core/src/analyzer/expression.rs +++ b/crates/flowscope-core/src/analyzer/expression.rs @@ -48,6 +48,30 @@ fn binary_op_operands(expr: &Expr) -> Vec<&Expr> { operands } +/// Splits a field access into the name it starts from and the steps after it. +/// +/// The parser keeps `o.items[1]` as the root `o` followed by the steps `.items` +/// and `[1]`. The name is `o.items`, read as it would be without the subscript, +/// and the steps after it are fields and subscripts of its value. The name is +/// empty when the root is not an identifier. +fn split_access_name<'e>( + root: &'e Expr, + access_chain: &'e [ast::AccessExpr], +) -> (Vec<&'e ast::Ident>, &'e [ast::AccessExpr]) { + let Expr::Identifier(first) = root else { + return (Vec::new(), access_chain); + }; + let mut name = vec![first]; + for step in access_chain { + match step { + ast::AccessExpr::Dot(Expr::Identifier(part)) => name.push(part), + _ => break, + } + } + let steps = &access_chain[name.len() - 1..]; + (name, steps) +} + /// Analyzes SQL expressions to extract column references, detect aggregations, /// and capture filter predicates. /// @@ -277,6 +301,10 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { FunctionArg::Named { arg: FunctionArgExpr::Expr(e), .. + } + | FunctionArg::ExprNamed { + arg: FunctionArgExpr::Expr(e), + .. } => { depth_limited |= Self::collect_column_refs(e, refs, dialect, next_depth); @@ -314,7 +342,7 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { Expr::Nested(inner) => { depth_limited |= Self::collect_column_refs(inner, refs, dialect, next_depth); } - Expr::Subquery(_) => { + Expr::Subquery(_) | Expr::Exists { .. } => { // Subquery columns are handled separately } Expr::InList { expr, list, .. } => { @@ -348,9 +376,207 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { Expr::Extract { expr, .. } => { depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); } - _ => { - // Other expressions don't contain column references or are handled elsewhere + // The parser gives TRIM, SUBSTRING, CEIL, `col:path` and the like + // variants of their own rather than function calls. There is no + // catch-all arm, so a variant added to sqlparser is a compile error + // here instead of a column that silently drops out of lineage. + Expr::IsUnknown(e) + | Expr::IsNotUnknown(e) + | Expr::IsNormalized { expr: e, .. } + | Expr::Ceil { expr: e, .. } + | Expr::Floor { expr: e, .. } + | Expr::Collate { expr: e, .. } + | Expr::Prefixed { value: e, .. } + | Expr::Named { expr: e, .. } + | Expr::OuterJoin(e) + | Expr::Prior(e) => { + depth_limited |= Self::collect_column_refs(e, refs, dialect, next_depth); + } + Expr::Interval(interval) => { + depth_limited |= + Self::collect_column_refs(&interval.value, refs, dialect, next_depth); + } + Expr::IsDistinctFrom(left, right) + | Expr::IsNotDistinctFrom(left, right) + | Expr::AnyOp { left, right, .. } + | Expr::AllOp { left, right, .. } + | Expr::SimilarTo { + expr: left, + pattern: right, + .. + } + | Expr::RLike { + expr: left, + pattern: right, + .. + } + | Expr::Position { + expr: left, + r#in: right, + } + | Expr::AtTimeZone { + timestamp: left, + time_zone: right, + } + | Expr::InUnnest { + expr: left, + array_expr: right, + .. + } => { + depth_limited |= Self::collect_column_refs(left, refs, dialect, next_depth); + depth_limited |= Self::collect_column_refs(right, refs, dialect, next_depth); + } + Expr::MemberOf(member_of) => { + depth_limited |= + Self::collect_column_refs(&member_of.value, refs, dialect, next_depth); + depth_limited |= + Self::collect_column_refs(&member_of.array, refs, dialect, next_depth); + } + Expr::Trim { + expr, + trim_what, + trim_characters, + .. + } => { + depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); + if let Some(what) = trim_what { + depth_limited |= Self::collect_column_refs(what, refs, dialect, next_depth); + } + for characters in trim_characters.iter().flatten() { + depth_limited |= + Self::collect_column_refs(characters, refs, dialect, next_depth); + } + } + Expr::Substring { + expr, + substring_from, + substring_for, + .. + } => { + depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); + for bound in [substring_from, substring_for].into_iter().flatten() { + depth_limited |= Self::collect_column_refs(bound, refs, dialect, next_depth); + } + } + Expr::Overlay { + expr, + overlay_what, + overlay_from, + overlay_for, + } => { + depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); + depth_limited |= Self::collect_column_refs(overlay_what, refs, dialect, next_depth); + depth_limited |= Self::collect_column_refs(overlay_from, refs, dialect, next_depth); + if let Some(length) = overlay_for { + depth_limited |= Self::collect_column_refs(length, refs, dialect, next_depth); + } + } + Expr::Convert { expr, styles, .. } => { + depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); + for style in styles { + depth_limited |= Self::collect_column_refs(style, refs, dialect, next_depth); + } + } + // A dot step names a field of the value, not a column of the query; + // a subscript or bracket key is an expression and can read one. + Expr::JsonAccess { value, path } => { + depth_limited |= Self::collect_column_refs(value, refs, dialect, next_depth); + for elem in &path.path { + if let ast::JsonPathElem::Bracket { key } = elem { + depth_limited |= Self::collect_column_refs(key, refs, dialect, next_depth); + } + } + } + Expr::CompoundFieldAccess { root, access_chain } => { + let (name, steps) = split_access_name(root, access_chain); + match name.split_last() { + Some((column, table)) if !table.is_empty() => refs.push(ColumnRef { + table: Some( + table + .iter() + .map(|part| part.value.as_str()) + .collect::>() + .join("."), + ), + column: column.value.clone(), + }), + _ => { + depth_limited |= Self::collect_column_refs(root, refs, dialect, next_depth); + } + } + for access in steps { + match access { + ast::AccessExpr::Dot(_) => {} + ast::AccessExpr::Subscript(ast::Subscript::Index { index }) => { + depth_limited |= + Self::collect_column_refs(index, refs, dialect, next_depth); + } + ast::AccessExpr::Subscript(ast::Subscript::Slice { + lower_bound, + upper_bound, + stride, + }) => { + for bound in [lower_bound, upper_bound, stride].into_iter().flatten() { + depth_limited |= + Self::collect_column_refs(bound, refs, dialect, next_depth); + } + } + } + } + } + Expr::Struct { values: exprs, .. } | Expr::Array(ast::Array { elem: exprs, .. }) => { + for e in exprs { + depth_limited |= Self::collect_column_refs(e, refs, dialect, next_depth); + } + } + Expr::GroupingSets(sets) | Expr::Cube(sets) | Expr::Rollup(sets) => { + for e in sets.iter().flatten() { + depth_limited |= Self::collect_column_refs(e, refs, dialect, next_depth); + } + } + Expr::Dictionary(fields) => { + for field in fields { + depth_limited |= + Self::collect_column_refs(&field.value, refs, dialect, next_depth); + } + } + Expr::Map(map) => { + for entry in &map.entries { + depth_limited |= + Self::collect_column_refs(&entry.key, refs, dialect, next_depth); + depth_limited |= + Self::collect_column_refs(&entry.value, refs, dialect, next_depth); + } + } + Expr::MatchAgainst { columns, .. } => { + for name in columns { + let mut parts = name.0.iter().filter_map(|part| part.as_ident()); + let Some(column) = parts.next_back() else { + continue; + }; + let table = parts.map(|i| i.value.as_str()).collect::>(); + refs.push(ColumnRef { + table: (!table.is_empty()).then(|| table.join(".")), + column: column.value.clone(), + }); + } + } + Expr::InSubquery { expr, .. } => { + // The value tested is read here; the subquery has its own scope + // and is analysed separately. + depth_limited |= Self::collect_column_refs(expr, refs, dialect, next_depth); + } + Expr::Lambda(_) => { + // Reading the body would take the lambda's parameters for + // columns of the enclosing query. } + Expr::CompoundIdentifier(_) => { + // The parser never builds a compound name of fewer than two parts. + } + Expr::Value(_) + | Expr::TypedString(_) + | Expr::Wildcard(_) + | Expr::QualifiedWildcard(..) => {} } depth_limited @@ -693,6 +919,9 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { identifiers } + /// Walks the same forms as `collect_column_refs`. A lateral alias that + /// lineage reads but this walk misses is never substituted, so it stays a + /// column reference and is resolved against the tables in scope. fn collect_simple_identifiers(expr: &Expr, identifiers: &mut HashSet, depth: usize) { if depth > MAX_RECURSION_DEPTH { #[cfg(feature = "tracing")] @@ -725,9 +954,17 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { | Expr::IsNotTrue(e) | Expr::IsUnknown(e) | Expr::IsNotUnknown(e) - | Expr::JsonAccess { value: e, .. } => { + | Expr::IsNormalized { expr: e, .. } + | Expr::Collate { expr: e, .. } + | Expr::Prefixed { value: e, .. } + | Expr::Named { expr: e, .. } + | Expr::OuterJoin(e) + | Expr::Prior(e) => { Self::collect_simple_identifiers(e, identifiers, next_depth); } + Expr::Interval(interval) => { + Self::collect_simple_identifiers(&interval.value, identifiers, next_depth); + } // Two expression patterns (left/right) Expr::BinaryOp { .. } => { @@ -771,6 +1008,10 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { Self::collect_simple_identifiers(e1, identifiers, next_depth); Self::collect_simple_identifiers(e2, identifiers, next_depth); } + Expr::MemberOf(member_of) => { + Self::collect_simple_identifiers(&member_of.value, identifiers, next_depth); + Self::collect_simple_identifiers(&member_of.array, identifiers, next_depth); + } // Three expression patterns Expr::Between { @@ -782,7 +1023,9 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { } // List patterns - Expr::Tuple(exprs) => { + Expr::Tuple(exprs) + | Expr::Struct { values: exprs, .. } + | Expr::Array(ast::Array { elem: exprs, .. }) => { for e in exprs { Self::collect_simple_identifiers(e, identifiers, next_depth); } @@ -793,6 +1036,22 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { Self::collect_simple_identifiers(e, identifiers, next_depth); } } + Expr::GroupingSets(sets) | Expr::Cube(sets) | Expr::Rollup(sets) => { + for e in sets.iter().flatten() { + Self::collect_simple_identifiers(e, identifiers, next_depth); + } + } + Expr::Dictionary(fields) => { + for field in fields { + Self::collect_simple_identifiers(&field.value, identifiers, next_depth); + } + } + Expr::Map(map) => { + for entry in &map.entries { + Self::collect_simple_identifiers(&entry.key, identifiers, next_depth); + Self::collect_simple_identifiers(&entry.value, identifiers, next_depth); + } + } // Function arguments Expr::Function(func) => { @@ -803,6 +1062,10 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { | FunctionArg::Named { arg: FunctionArgExpr::Expr(e), .. + } + | FunctionArg::ExprNamed { + arg: FunctionArgExpr::Expr(e), + .. } => Self::collect_simple_identifiers(e, identifiers, next_depth), _ => {} } @@ -845,8 +1108,84 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { } } + Expr::Trim { + expr, + trim_what, + trim_characters, + .. + } => { + Self::collect_simple_identifiers(expr, identifiers, next_depth); + if let Some(what) = trim_what { + Self::collect_simple_identifiers(what, identifiers, next_depth); + } + for characters in trim_characters.iter().flatten() { + Self::collect_simple_identifiers(characters, identifiers, next_depth); + } + } + Expr::Overlay { + expr, + overlay_what, + overlay_from, + overlay_for, + } => { + Self::collect_simple_identifiers(expr, identifiers, next_depth); + Self::collect_simple_identifiers(overlay_what, identifiers, next_depth); + Self::collect_simple_identifiers(overlay_from, identifiers, next_depth); + if let Some(length) = overlay_for { + Self::collect_simple_identifiers(length, identifiers, next_depth); + } + } + Expr::Convert { expr, styles, .. } => { + Self::collect_simple_identifiers(expr, identifiers, next_depth); + for style in styles { + Self::collect_simple_identifiers(style, identifiers, next_depth); + } + } + Expr::JsonAccess { value, path } => { + Self::collect_simple_identifiers(value, identifiers, next_depth); + for elem in &path.path { + if let ast::JsonPathElem::Bracket { key } = elem { + Self::collect_simple_identifiers(key, identifiers, next_depth); + } + } + } + Expr::CompoundFieldAccess { root, access_chain } => { + let (name, steps) = split_access_name(root, access_chain); + // `o.items[1]` names a qualified column, which is no simple identifier. + if name.len() < 2 { + Self::collect_simple_identifiers(root, identifiers, next_depth); + } + for access in steps { + match access { + ast::AccessExpr::Dot(_) => {} + ast::AccessExpr::Subscript(ast::Subscript::Index { index }) => { + Self::collect_simple_identifiers(index, identifiers, next_depth); + } + ast::AccessExpr::Subscript(ast::Subscript::Slice { + lower_bound, + upper_bound, + stride, + }) => { + for bound in [lower_bound, upper_bound, stride].into_iter().flatten() { + Self::collect_simple_identifiers(bound, identifiers, next_depth); + } + } + } + } + } + Expr::MatchAgainst { columns, .. } => { + for name in columns { + if let [ast::ObjectNamePart::Identifier(ident)] = name.0.as_slice() { + identifiers.insert(ident.value.clone()); + } + } + } + Expr::InSubquery { expr, .. } => { + Self::collect_simple_identifiers(expr, identifiers, next_depth); + } + // Skip subqueries - they have their own scope - Expr::Subquery(_) | Expr::InSubquery { .. } | Expr::Exists { .. } => {} + Expr::Subquery(_) | Expr::Exists { .. } => {} // Skip qualified names (table.column) - not simple identifiers Expr::CompoundIdentifier(_) => {} @@ -928,6 +1267,112 @@ mod tests { assert!(hit, "nesting past the limit should still trigger the guard"); } + /// Expressions whose operands the parser keeps outside function arguments, + /// with the columns each one reads, qualified as written. + const DEDICATED_FORMS: &[(Dialect, &str, &[&str])] = &[ + ( + Dialect::Generic, + "TRIM(BOTH pad FROM c.name)", + &["c.name", "pad"], + ), + (Dialect::Snowflake, "TRIM(name, pad)", &["name", "pad"]), + ( + Dialect::Generic, + "SUBSTRING(name FROM start FOR len)", + &["name", "start", "len"], + ), + ( + Dialect::Generic, + "POSITION(needle IN hay)", + &["needle", "hay"], + ), + ( + Dialect::Generic, + "CEIL(amount) + FLOOR(fee)", + &["amount", "fee"], + ), + (Dialect::Generic, "ts AT TIME ZONE zone", &["ts", "zone"]), + (Dialect::Generic, "a IS DISTINCT FROM b", &["a", "b"]), + (Dialect::Generic, "flag IS UNKNOWN", &["flag"]), + ( + Dialect::Snowflake, + "payload:items[pos]", + &["payload", "pos"], + ), + (Dialect::Postgres, "items[pos].field", &["items", "pos"]), + (Dialect::Postgres, "t.items[pos]", &["t.items", "pos"]), + (Dialect::Duckdb, "s.t.items[1].field", &["s.t.items"]), + (Dialect::Postgres, "name COLLATE \"C\"", &["name"]), + ( + Dialect::Postgres, + "OVERLAY(name PLACING patch FROM 2 FOR len)", + &["name", "patch", "len"], + ), + ( + Dialect::Postgres, + "name SIMILAR TO pattern", + &["name", "pattern"], + ), + (Dialect::Mysql, "name RLIKE pattern", &["name", "pattern"]), + (Dialect::Postgres, "a = ANY(b)", &["a", "b"]), + (Dialect::Mysql, "CONVERT(amount, CHAR)", &["amount"]), + (Dialect::Postgres, "ARRAY[a, b]", &["a", "b"]), + (Dialect::Postgres, "f(arg => col)", &["col"]), + (Dialect::Mysql, "DATE_ADD(d, INTERVAL n DAY)", &["d", "n"]), + ( + Dialect::Mysql, + "MATCH (title, t.body) AGAINST ('x')", + &["title", "t.body"], + ), + (Dialect::Generic, "x IN (SELECT y FROM t)", &["x"]), + (Dialect::Generic, "EXISTS (SELECT y FROM t)", &[]), + (Dialect::Duckdb, "list_transform(xs, x -> x + 1)", &["xs"]), + ]; + + fn parse_expr(sql: &str, dialect: Dialect) -> Expr { + sqlparser::parser::Parser::new(dialect.to_sqlparser_dialect().as_ref()) + .try_with_sql(sql) + .and_then(|mut parser| parser.parse_expr()) + .unwrap_or_else(|err| panic!("{sql} should parse: {err}")) + } + + #[test] + fn collect_column_refs_reads_operands_of_dedicated_forms() { + for &(dialect, sql, expected) in DEDICATED_FORMS { + let (refs, hit) = ExpressionAnalyzer::extract_column_refs_with_dialect( + &parse_expr(sql, dialect), + dialect, + ); + + assert!(!hit, "{sql} should not trigger the guard"); + let columns: Vec = refs + .into_iter() + .map(|col_ref| match col_ref.table { + Some(table) => format!("{table}.{}", col_ref.column), + None => col_ref.column, + }) + .collect(); + assert_eq!(columns, expected, "columns read by {sql}"); + } + } + + #[test] + fn extract_simple_identifiers_sees_what_collect_column_refs_reads() { + for &(dialect, sql, expected) in DEDICATED_FORMS { + let unqualified: HashSet = expected + .iter() + .filter(|column| !column.contains('.')) + .map(|column| column.to_string()) + .collect(); + + assert_eq!( + ExpressionAnalyzer::extract_simple_identifiers(&parse_expr(sql, dialect)), + unqualified, + "bare identifiers of {sql}" + ); + } + } + #[test] fn extract_simple_identifiers_walks_long_operator_chains() { let len = MAX_RECURSION_DEPTH * 3; diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index ced3a2d0..ae5f31f8 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -7297,6 +7297,200 @@ fn long_concatenation_keeps_every_operand() { ); } +/// Sources of an output column as `relation.column`, lowercased, following +/// derivation and data-flow edges to the relation that owns each source. +fn source_columns_of(stmt: &StmtView<'_>, output: &str) -> HashSet { + let target = stmt + .nodes + .iter() + .find(|node| node.node_type == NodeType::Column && node.label.eq_ignore_ascii_case(output)) + .unwrap_or_else(|| panic!("output column {output} should exist")); + stmt.edges + .iter() + .filter(|edge| { + edge.to == target.id + && matches!(edge.edge_type, EdgeType::Derivation | EdgeType::DataFlow) + }) + .filter_map(|edge| { + let source = stmt.nodes.iter().find(|node| node.id == edge.from)?; + let owner = stmt + .edges + .iter() + .find(|owned| owned.edge_type == EdgeType::Ownership && owned.to == source.id) + .and_then(|owned| stmt.nodes.iter().find(|node| node.id == owned.from))?; + Some(format!("{}.{}", owner.label, source.label).to_lowercase()) + }) + .collect() +} + +fn expected_sources(sources: &[&str]) -> HashSet { + sources.iter().map(|source| source.to_string()).collect() +} + +#[test] +fn trim_reads_the_column_of_the_relation_it_names() { + // Both relations have `name`. The column inside TRIM belongs to the joined + // relation, and must not be credited to the one the query is driven from. + let sql = r#" + SELECT o.id, TRIM(c.name) AS customer_name + FROM orders o + JOIN customers c ON o.customer_id = c.id + "#; + let schema = SchemaMetadata { + allow_implied: false, + default_catalog: None, + default_schema: None, + search_path: None, + case_sensitivity: None, + tables: vec![ + schema_table(None, None, "orders", &["id", "customer_id", "name"]), + schema_table(None, None, "customers", &["id", "name"]), + ], + }; + + let result = run_analysis(sql, Dialect::Generic, Some(schema)); + let stmt = first_statement(&result); + assert_eq!( + source_columns_of(&stmt, "customer_name"), + expected_sources(&["customers.name"]) + ); +} + +#[test] +fn dedicated_expression_forms_read_every_operand() { + // SUBSTR, CEIL, POSITION and the `:` path operator parse to their own + // expression variants rather than to function calls. + let sql = r#" + SELECT + SUBSTR(code, 1, 2) || suffix AS short_code, + CEIL(amount) + fee AS total, + payload:kind::STRING || suffix AS tagged, + POSITION('-' IN code) + fee AS dash_at + FROM events + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + for (output, sources) in [ + ("short_code", ["events.code", "events.suffix"]), + ("total", ["events.amount", "events.fee"]), + ("tagged", ["events.payload", "events.suffix"]), + ("dash_at", ["events.code", "events.fee"]), + ] { + assert_eq!( + source_columns_of(&stmt, output), + expected_sources(&sources), + "sources of {output}" + ); + } +} + +#[test] +fn lateral_alias_inside_trim_resolves_to_its_sources() { + // Keys hashed from earlier keys of the same SELECT list: the lateral + // aliases are only visible inside TRIM. + let sql = r#" + SELECT + MD5(UPPER(TRIM(CAST(order_id AS VARCHAR)))) AS order_key, + MD5(UPPER(TRIM(CAST(customer_id AS VARCHAR)))) AS customer_key, + MD5(CONCAT( + UPPER(TRIM(CAST(order_key AS VARCHAR))), '||', + UPPER(TRIM(CAST(customer_key AS VARCHAR))) + )) AS link_key + FROM orders + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert_eq!( + source_columns_of(&stmt, "link_key"), + expected_sources(&["orders.order_id", "orders.customer_id"]) + ); +} + +#[test] +fn filter_inside_trim_is_attached_to_its_table() { + let sql = r#" + SELECT o.id + FROM orders o + JOIN customers c ON o.customer_id = c.id + WHERE TRIM(c.status) = 'active' + "#; + + let result = run_analysis(sql, Dialect::Generic, None); + let stmt = first_statement(&result); + let customers = find_table_node(&stmt, "customers").expect("customers table not found"); + let orders = find_table_node(&stmt, "orders").expect("orders table not found"); + assert_eq!( + customers + .filters + .iter() + .map(|filter| filter.expression.as_str()) + .collect::>(), + ["TRIM(c.status) = 'active'"] + ); + assert!(orders.filters.is_empty(), "orders: {:?}", orders.filters); +} + +#[test] +fn in_subquery_reads_its_left_operand_only() { + // The subquery is analysed in its own scope, so its columns are not + // sources of the flag; the value tested against it is. + let sql = r#" + SELECT + o.id, + CASE WHEN o.customer_id IN (SELECT c.id FROM customers c) THEN 1 ELSE 0 END AS known + FROM orders o + "#; + + let result = run_analysis(sql, Dialect::Generic, None); + let stmt = first_statement(&result); + assert_eq!( + source_columns_of(&stmt, "known"), + expected_sources(&["orders.customer_id"]) + ); +} + +#[test] +fn lambda_parameters_are_not_read_as_columns() { + let sql = r#" + SELECT list_transform(prices, p -> p * rate) AS scaled + FROM products + "#; + + let result = run_analysis(sql, Dialect::Duckdb, None); + let stmt = first_statement(&result); + let sources = source_columns_of(&stmt, "scaled"); + assert!( + !sources.iter().any(|source| source.ends_with(".p")), + "a lambda parameter is not a column: {sources:?}" + ); + assert!(sources.contains("products.prices"), "sources: {sources:?}"); +} + +#[test] +fn subscript_of_a_qualified_column_reads_that_column() { + // `o.items[1]` parses as the root `o` with the steps `.items` and `[1]`. + // The root is the table alias, not a column of orders. + let sql = "SELECT o.id, o.items[1] AS first_item FROM orders o"; + let schema = SchemaMetadata { + allow_implied: false, + default_catalog: None, + default_schema: None, + search_path: None, + case_sensitivity: None, + tables: vec![schema_table(None, None, "orders", &["id", "items"])], + }; + + let result = run_analysis(sql, Dialect::Postgres, Some(schema)); + let stmt = first_statement(&result); + assert_eq!( + source_columns_of(&stmt, "first_item"), + expected_sources(&["orders.items"]) + ); + assert!(result.issues.is_empty(), "issues: {:?}", result.issues); +} + #[test] fn case_with_aggregate_in_condition() { // CASE with aggregate function in WHEN condition diff --git a/crates/flowscope-core/tests/snapshots/snapshots__run_postgres_snapshot_test@postgres_array_slicing.snap b/crates/flowscope-core/tests/snapshots/snapshots__run_postgres_snapshot_test@postgres_array_slicing.snap index 244d1a7a..64dd0f59 100644 --- a/crates/flowscope-core/tests/snapshots/snapshots__run_postgres_snapshot_test@postgres_array_slicing.snap +++ b/crates/flowscope-core/tests/snapshots/snapshots__run_postgres_snapshot_test@postgres_array_slicing.snap @@ -48,6 +48,32 @@ expression: clean_result ], "expression": "c[2:]" }, + { + "id": "column_95b826196d892dac", + "type": "column", + "label": "c", + "qualifiedName": "array_data.c", + "canonicalName": { + "schema": "array_data", + "name": "c" + }, + "statementIds": [ + 0 + ] + }, + { + "id": "column_b683e029cadda1d4", + "type": "column", + "label": "a", + "qualifiedName": "array_data.a", + "canonicalName": { + "schema": "array_data", + "name": "a" + }, + "statementIds": [ + 0 + ] + }, { "id": "column_bceb3c4d983131c5", "type": "column", @@ -60,6 +86,32 @@ expression: clean_result ], "expression": "d[2:3]" }, + { + "id": "column_c4a845893b35a6fa", + "type": "column", + "label": "b", + "qualifiedName": "array_data.b", + "canonicalName": { + "schema": "array_data", + "name": "b" + }, + "statementIds": [ + 0 + ] + }, + { + "id": "column_ea9f00345477a676", + "type": "column", + "label": "d", + "qualifiedName": "array_data.d", + "canonicalName": { + "schema": "array_data", + "name": "d" + }, + "statementIds": [ + 0 + ] + }, { "id": "output_74d22f755ec62dbc", "type": "output", @@ -100,8 +152,27 @@ expression: clean_result ], "edges": [ { - "id": "edge_550b27dba85a2465", + "id": "edge_16ba963d993199aa", + "from": "column_95b826196d892dac", + "to": "column_84c9f50405fdf52c", + "type": "derivation", + "expression": "c[2:]", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_1d5f16b94e24d605", "from": "table_eb6f82b7323d024e", + "to": "column_c4a845893b35a6fa", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_29829752b6d6e0ac", + "from": "column_ea9f00345477a676", "to": "column_bceb3c4d983131c5", "type": "derivation", "expression": "d[2:3]", @@ -110,17 +181,26 @@ expression: clean_result ] }, { - "id": "edge_77f5912a72ad7682", - "from": "output_74d22f755ec62dbc", - "to": "column_7748741d00eb4aa0", + "id": "edge_30772b2aeea20852", + "from": "table_eb6f82b7323d024e", + "to": "column_ea9f00345477a676", "type": "ownership", "statementIds": [ 0 ] }, { - "id": "edge_8128a9fd6e58dbf4", + "id": "edge_675b380196aa1111", "from": "table_eb6f82b7323d024e", + "to": "column_95b826196d892dac", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_75998b7fa2f4cdb8", + "from": "column_b683e029cadda1d4", "to": "column_7748741d00eb4aa0", "type": "derivation", "expression": "a[:]", @@ -129,11 +209,10 @@ expression: clean_result ] }, { - "id": "edge_91d9a3175fef27e4", - "from": "table_eb6f82b7323d024e", - "to": "column_6d76c5c7b69b69cb", - "type": "derivation", - "expression": "b[:1]", + "id": "edge_77f5912a72ad7682", + "from": "output_74d22f755ec62dbc", + "to": "column_7748741d00eb4aa0", + "type": "ownership", "statementIds": [ 0 ] @@ -148,11 +227,20 @@ expression: clean_result ] }, { - "id": "edge_ad3bc1a5b25c59cc", + "id": "edge_a4180dafa0264175", "from": "table_eb6f82b7323d024e", - "to": "column_84c9f50405fdf52c", + "to": "column_b683e029cadda1d4", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_dd32db879aa8c27f", + "from": "column_c4a845893b35a6fa", + "to": "column_6d76c5c7b69b69cb", "type": "derivation", - "expression": "c[2:]", + "expression": "b[:1]", "statementIds": [ 0 ] @@ -180,7 +268,7 @@ expression: clean_result "summary": { "statementCount": 1, "tableCount": 1, - "columnCount": 4, + "columnCount": 8, "joinCount": 0, "complexityScore": 5, "issueCount": { From 7b57f0d44ed4ad4686e65b84c0b3f98fc442caad Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:02:38 +0200 Subject: [PATCH 3/8] docs(changelog): note the expression lineage fixes Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4fffbd80..fe5380d7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- Walked left-deep chains of binary operators as chains, so that a generated surrogate key over more than about 100 operands keeps every operand in its lineage instead of losing the leftmost ones and reporting `APPROXIMATE_LINEAGE`. +- Read columns inside TRIM, SUBSTRING, POSITION, CEIL and FLOOR, AT TIME ZONE, IS [NOT] DISTINCT FROM, path and subscript accessors, COLLATE, OVERLAY, SIMILAR TO, RLIKE, ANY and ALL as lineage sources. + ## [0.9.2] - 2026-09-24 ### Added From 9e47889bb9f0a9fdf75ade10acc094a7894aefaf Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Thu, 24 Sep 2026 17:00:04 +0200 Subject: [PATCH 4/8] fix(core): read LATERAL FLATTEN as a relation with its own columns sqlparser gives `LATERAL FLATTEN(...)` its own table factor, `TableFactor::Function`, which the lineage visitor did not match. The alias was never registered and no issue was raised, so a column read through it lost its source: SELECT o.order_id, f.value::STRING AS tag FROM orders o, LATERAL FLATTEN(input => o.tags) f `f` was taken for a table name: tag came out derived from a column `value` of a table `F` that does not exist, the implied schema listed `F` beside ORDERS, and orders.tags was read neither for lineage nor for the implied schema. An unqualified `value` was reported as ambiguous across the tables in scope, and a FLATTEN over the value of another had nothing to follow. Register FLATTEN under its alias, or under its name when it has none, as a relation with the columns Snowflake documents: SEQ, KEY, PATH, INDEX, VALUE and THIS. VALUE and THIS carry the element, so they derive from the columns of the `INPUT =>` argument, or of the first positional one, with the call as their expression. SEQ, KEY, PATH and INDEX say where an element sits and read nothing. They are added without the relation-level edge a source-less projection gets, which would have made `'n' || f.index` depend on every base relation of the FROM clause. Each FLATTEN gets its own node, keyed by its scope, since CTEs commonly reuse one alias such as `f`. A star over the scope now includes the six columns. The alias is also recorded as a subquery alias, as a derived table's is, so that what is read through it stays out of the implied schema, and the qualified columns of the input are recorded there as projected columns are: the query above now implies ORDERS(order_id, tags) and nothing else. With column lineage disabled only that alias is recorded: at table level the input is a column of a relation already in the FROM clause, so a FLATTEN node would have no edge, and no column is emitted. Those six columns are all a FLATTEN has, so the scope records it as a relation with fixed columns. A bare name it lacks, with no schema to place it, still resolves to the one other relation of the FROM clause, as it did while FLATTEN went unregistered: `id` in `SELECT id FROM orders o, LATERAL FLATTEN(input => o.tags) f` stays orders.id instead of turning ambiguous, and in `FROM customers, LATERAL FLATTEN(input => addresses) a, LATERAL FLATTEN(input => phones) p` the second input is read from customers. A warning for a name that is ambiguous between tables no longer lists the FLATTEN beside them. VALUE and THIS read the same input, so an input that cannot be resolved is reported once. Any other `LATERAL (...)` now raises UNSUPPORTED_SYNTAX, as `TABLE((...))` already does, instead of being dropped silently. The snowflake_lateral_flatten snapshot changes accordingly: `value AS p_id` now derives from b.cool_ids through the FLATTEN node, B's implied columns gain cool_ids, and the ambiguity warning for `value` is gone. Tests: the query above with a decoy `value` column on orders, its implied schema without a schema given, chained flattens back to the first input, the position columns reading nothing under a CTE joined to another table, a literal input, one alias reused by two CTEs with named and positional input, a star over FLATTEN, a bare column and a bare filter beside a FLATTEN with no schema, a bare input to a second FLATTEN, an ambiguous input reported once, SPLIT_TO_TABLE reported as unsupported, and no FLATTEN node or column with column lineage disabled. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/flowscope-core/src/analyzer/context.rs | 23 + crates/flowscope-core/src/analyzer/query.rs | 51 +- crates/flowscope-core/src/analyzer/visitor.rs | 171 +++++++ crates/flowscope-core/tests/lineage_engine.rs | 440 ++++++++++++++++++ ...apshot_test@snowflake_lateral_flatten.snap | 209 ++++++++- 5 files changed, 878 insertions(+), 16 deletions(-) diff --git a/crates/flowscope-core/src/analyzer/context.rs b/crates/flowscope-core/src/analyzer/context.rs index 68cd4a62..92ebf609 100644 --- a/crates/flowscope-core/src/analyzer/context.rs +++ b/crates/flowscope-core/src/analyzer/context.rs @@ -58,6 +58,9 @@ pub(crate) struct Scope { /// True when the scope contains a table function relation whose output /// columns may be dialect-provided rather than schema-backed. pub(crate) has_table_function_relation: bool, + /// Node IDs of relations whose columns are all known, such as a FLATTEN. + /// An unqualified column missing from them belongs to another relation. + pub(crate) fixed_column_relations: HashSet>, } impl Scope { @@ -676,6 +679,26 @@ impl StatementContext { .is_some_and(|scope| scope.has_table_function_relation) } + /// Mark a relation of the current scope as having no columns beyond those + /// registered for it. + pub(crate) fn mark_fixed_columns_in_scope(&mut self, node_id: Arc) { + if let Some(scope) = self.current_scope_mut() { + scope.fixed_column_relations.insert(node_id); + } + } + + /// Relation instances of the current scope that may have a column missing + /// from every known column list, that is all but those with fixed columns. + pub(crate) fn open_relation_instances_in_current_scope(&self) -> Vec { + let Some(scope) = self.current_scope() else { + return Vec::new(); + }; + self.relation_instances_in_current_scope() + .into_iter() + .filter(|instance| !scope.fixed_column_relations.contains(&instance.node_id)) + .collect() + } + /// Register statement-global output columns for a named CTE definition. pub(crate) fn register_cte_output_columns(&mut self, name: String, columns: Vec) { self.aliased_subquery_columns.insert(name, columns); diff --git a/crates/flowscope-core/src/analyzer/query.rs b/crates/flowscope-core/src/analyzer/query.rs index 4569bca6..7f2af145 100644 --- a/crates/flowscope-core/src/analyzer/query.rs +++ b/crates/flowscope-core/src/analyzer/query.rs @@ -865,13 +865,20 @@ impl<'a> Analyzer<'a> { } } 0 => { - // No candidates found - if there's only one relation instance in scope, use it - // (the column might exist but not be in our schema) - if relation_instances.len() == 1 { - return Some(relation_instances[0].canonical.clone()); + // No candidates found - if only one relation instance in scope may have + // columns we do not know of, use it (the column might exist but not be in + // our schema). A FLATTEN has none, so it cannot own the column. + let owners = ctx.open_relation_instances_in_current_scope(); + if owners.len() == 1 { + return Some(owners[0].canonical.clone()); } // Multiple tables but column not found in any - ambiguous - let mut sorted_tables: Vec<_> = relation_instances + let listed = if owners.is_empty() { + &relation_instances + } else { + &owners + }; + let mut sorted_tables: Vec<_> = listed .iter() .map(|instance| instance.canonical.clone()) .collect(); @@ -1212,6 +1219,40 @@ impl<'a> Analyzer<'a> { }); } + /// Adds a column that `target_node` produces without reading any column, + /// such as the position of an element in a flattened array. + /// + /// Unlike a source-less projection, the column is not attributed to the + /// relations in scope: it describes the relation that owns it. + pub(super) fn add_unsourced_column( + &mut self, + ctx: &mut StatementContext, + name: &str, + target_node: &Arc, + ) { + let normalized_name = self.normalize_identifier(name); + let node_id = generate_column_node_id(Some(target_node), &normalized_name); + ctx.add_node(Node { + id: node_id.clone(), + node_type: NodeType::Column, + label: normalized_name.clone().into(), + ..Default::default() + }); + let edge_id = generate_edge_id(target_node, &node_id); + if !ctx.edge_ids.contains(&edge_id) { + ctx.add_edge(Edge::ownership( + edge_id, + target_node.clone(), + node_id.clone(), + )); + } + ctx.output_columns.push(OutputColumn { + name: normalized_name, + data_type: None, + node_id, + }); + } + fn is_definitely_ambiguous_unqualified_column( &self, ctx: &StatementContext, diff --git a/crates/flowscope-core/src/analyzer/visitor.rs b/crates/flowscope-core/src/analyzer/visitor.rs index 13ccbb31..acae5ab0 100644 --- a/crates/flowscope-core/src/analyzer/visitor.rs +++ b/crates/flowscope-core/src/analyzer/visitor.rs @@ -10,6 +10,7 @@ use super::helpers::{ alias_visibility_warning, find_cte_body_span, find_cte_definition_span, find_derived_table_alias_span, generate_statement_scoped_node_id, }; +use super::query::OutputColumnParams; use super::select_analyzer::SelectAnalyzer; use super::Analyzer; use crate::generated::is_value_table_function; @@ -465,6 +466,106 @@ impl<'a, 'b> LineageVisitor<'a, 'b> { } } + /// Registers Snowflake's `LATERAL FLATTEN(INPUT => expr) AS f` as a + /// relation with the columns FLATTEN returns, so that `f.value` resolves. + /// + /// VALUE and THIS carry the input, so they derive from the columns it + /// reads, with the whole call as their expression. SEQ, KEY, PATH and + /// INDEX say where an element sits rather than carry it, and read nothing. + fn visit_flatten( + &mut self, + name: &ast::ObjectName, + args: &[ast::FunctionArg], + alias: Option<&TableAlias>, + ) { + // Unaliased, the columns are still read unqualified (`value`), so the + // relation needs a name all the same. + let alias_name = alias.map_or_else(|| name.to_string(), |a| a.name.to_string()); + // What is read through the alias belongs to no table, so it stays out + // of the implied schema, as a derived table's columns do. + self.ctx + .register_subquery_alias_in_scope(alias_name.clone()); + // At table level the input is a column of a relation already in the + // FROM clause: a FLATTEN node would have nothing to connect to. + if !self.analyzer.column_lineage_enabled { + return; + } + // An alias such as `f` is often reused by every CTE that flattens, + // and each use is a different relation. + let scope_id = self.ctx.current_scope_id().unwrap_or_default(); + let node_id = self.ctx.add_node(Node { + id: generate_statement_scoped_node_id( + "flatten", + self.ctx.statement_index, + &format!("scope_{scope_id}::{alias_name}"), + ), + node_type: NodeType::Cte, + label: alias_name.clone().into(), + qualified_name: Some(alias_name.clone().into()), + ..Default::default() + }); + + let sources = match flatten_input(args) { + Some(input) => ExpressionAnalyzer::new(self.analyzer, self.ctx) + .extract_column_refs_with_warning(input), + None => Vec::new(), + }; + // The input is read as a projected column is, so its table learns of it. + for source in &sources { + let table = source.table.as_deref(); + if let Some(canonical) = table.and_then(|t| self.resolve_table_alias(Some(t))) { + self.ctx + .record_source_column(&canonical, &source.column, None); + } + } + let call = format!( + "{name}({})", + args.iter() + .map(ToString::to_string) + .collect::>() + .join(", ") + ); + + let checkpoint = self.ctx.projection_checkpoint(); + for column in ["SEQ", "KEY", "PATH", "INDEX", "VALUE", "THIS"] { + let carries_input = matches!(column, "VALUE" | "THIS"); + if carries_input && !sources.is_empty() { + let reported = self.analyzer.issues.len(); + self.analyzer.add_output_column_with_aggregation( + self.ctx, + OutputColumnParams { + name: column.to_string(), + sources: sources.clone(), + expression: Some(call.clone()), + data_type: None, + target_node: Some(node_id.to_string()), + approximate: false, + aggregation: None, + }, + ); + // THIS reads the input VALUE has just read, so whatever resolving + // it reports has been reported once already. + if column == "THIS" { + self.analyzer.issues.truncate(reported); + } + } else { + self.analyzer + .add_unsourced_column(self.ctx, column, &node_id); + } + } + let columns = self.ctx.take_output_columns_since(checkpoint); + + // The six columns are all FLATTEN returns: a bare name it lacks + // belongs to another relation of the FROM clause. + self.ctx.mark_fixed_columns_in_scope(node_id.clone()); + self.ctx + .register_table_in_scope(alias_name.clone(), node_id); + self.ctx + .register_alias_in_scope(alias_name.clone(), alias_name.clone()); + self.ctx + .register_subquery_columns_in_scope(alias_name, columns); + } + /// Emits a warning for unsupported alias usage in a clause. fn emit_alias_warning(&mut self, clause_name: &str, alias_name: &str) { let dialect = self.analyzer.request.dialect; @@ -749,6 +850,34 @@ impl<'a, 'b> Visitor for LineageVisitor<'a, 'b> { .with_statement(self.ctx.statement_index), ); } + TableFactor::Function { + name, args, alias, .. + } if name.to_string().eq_ignore_ascii_case("flatten") => { + self.visit_flatten(name, args, alias.as_ref()); + } + TableFactor::Function { + name, args, alias, .. + } => { + for expr in args.iter().filter_map(function_arg_expr) { + self.extract_identifiers_from_expr(expr); + } + if is_value_table_function(self.analyzer.request.dialect, &name.to_string()) { + self.ctx.mark_table_function_in_scope(); + } + if let Some(a) = alias { + self.ctx + .register_subquery_alias_in_scope(a.name.to_string()); + } + self.analyzer.issues.push( + Issue::info( + issue_codes::UNSUPPORTED_SYNTAX, + format!( + "Table function '{name}' lineage extracted with best-effort identifier matching" + ), + ) + .with_statement(self.ctx.statement_index), + ); + } TableFactor::Pivot { table, aggregate_functions, @@ -846,3 +975,45 @@ impl<'a, 'b> Visitor for LineageVisitor<'a, 'b> { } } } + +/// The expression passed to a table function argument, named or not. +fn function_arg_expr(arg: &ast::FunctionArg) -> Option<&Expr> { + match arg { + ast::FunctionArg::Named { + arg: ast::FunctionArgExpr::Expr(expr), + .. + } + | ast::FunctionArg::ExprNamed { + arg: ast::FunctionArgExpr::Expr(expr), + .. + } + | ast::FunctionArg::Unnamed(ast::FunctionArgExpr::Expr(expr)) => Some(expr), + _ => None, + } +} + +/// The value FLATTEN expands: its `INPUT =>` argument, or the first positional +/// one. PATH, OUTER, RECURSIVE and MODE are constants. +fn flatten_input(args: &[ast::FunctionArg]) -> Option<&Expr> { + let named = args.iter().find_map(|arg| { + let name = match arg { + ast::FunctionArg::Named { name, .. } => name, + ast::FunctionArg::ExprNamed { + name: Expr::Identifier(name), + .. + } => name, + _ => return None, + }; + if name.value.eq_ignore_ascii_case("input") { + function_arg_expr(arg) + } else { + None + } + }); + named.or_else(|| { + args.iter().find_map(|arg| match arg { + ast::FunctionArg::Unnamed(ast::FunctionArgExpr::Expr(expr)) => Some(expr), + _ => None, + }) + }) +} diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index ae5f31f8..9c4fd315 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -9824,6 +9824,446 @@ fn snowflake_lateral_flatten_tracks_source_table() { ); } +/// Table columns an output column is computed from, as `relation.column`, +/// lowercased, following derivation and data-flow edges through CTEs, +/// derived tables and table functions. +fn base_sources_of(stmt: &StmtView<'_>, output: &str) -> HashSet { + let owner_of = |id: &str| { + stmt.edges + .iter() + .find(|edge| edge.edge_type == EdgeType::Ownership && &*edge.to == id) + .and_then(|edge| stmt.nodes.iter().copied().find(|node| node.id == edge.from)) + }; + let target = stmt + .nodes + .iter() + .find(|node| { + node.node_type == NodeType::Column + && node.label.eq_ignore_ascii_case(output) + && owner_of(&node.id).is_some_and(|owner| owner.node_type == NodeType::Output) + }) + .unwrap_or_else(|| panic!("output column {output} should exist")); + + let mut sources = HashSet::new(); + let mut seen = HashSet::new(); + let mut pending = vec![target.id.clone()]; + while let Some(id) = pending.pop() { + if !seen.insert(id.clone()) { + continue; + } + for edge in stmt.edges.iter().filter(|edge| { + edge.to == id && matches!(edge.edge_type, EdgeType::Derivation | EdgeType::DataFlow) + }) { + let Some(source) = stmt + .nodes + .iter() + .find(|node| node.id == edge.from && node.node_type == NodeType::Column) + else { + continue; + }; + match owner_of(&source.id) { + Some(owner) if owner.node_type == NodeType::Table => { + sources.insert(format!("{}.{}", owner.label, source.label).to_lowercase()); + } + _ => pending.push(source.id.clone()), + } + } + } + sources +} + +/// Edges that carry data into the columns of the relation labelled `alias`, +/// keyed by the lowercased column label. +fn flows_into_columns_of<'a>(stmt: &StmtView<'a>, alias: &str) -> Vec<(String, &'a Edge)> { + let relation = find_cte_node(stmt, alias).unwrap_or_else(|| panic!("{alias} should exist")); + stmt.edges + .iter() + .filter(|owned| owned.edge_type == EdgeType::Ownership && owned.from == relation.id) + .flat_map(|owned| { + let label = stmt + .nodes + .iter() + .find(|node| node.id == owned.to) + .map(|node| node.label.to_lowercase()) + .unwrap_or_default(); + stmt.edges + .iter() + .copied() + .filter(move |edge| { + edge.to == owned.to + && matches!(edge.edge_type, EdgeType::Derivation | EdgeType::DataFlow) + }) + .map(move |edge| (label.clone(), edge)) + }) + .collect() +} + +fn schema_of(tables: Vec) -> SchemaMetadata { + SchemaMetadata { + allow_implied: false, + default_catalog: None, + default_schema: None, + search_path: None, + case_sensitivity: None, + tables, + } +} + +#[test] +fn flatten_value_comes_from_its_input_not_from_a_namesake() { + // `value` is also a column of orders: `f.value` must not resolve to it. + let sql = r#" + SELECT o.order_id, f.value::STRING AS tag + FROM orders o, + LATERAL FLATTEN(input => o.tags) f + "#; + let schema = schema_of(vec![schema_table( + None, + None, + "orders", + &["order_id", "tags", "value"], + )]); + + let result = run_analysis(sql, Dialect::Snowflake, Some(schema)); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "tag"), + expected_sources(&["orders.tags"]) + ); + + let flows = flows_into_columns_of(&stmt, "f"); + let expressions: HashSet<(&str, Option<&str>)> = flows + .iter() + .map(|(column, edge)| (column.as_str(), edge.expression.as_deref())) + .collect(); + assert_eq!( + expressions, + HashSet::from([ + ("value", Some("FLATTEN(input => o.tags)")), + ("this", Some("FLATTEN(input => o.tags)")), + ]) + ); +} + +#[test] +fn flatten_alias_is_not_an_implied_table() { + // Without a schema the tables are implied from what the query reads: + // orders has the flattened tags, and `f` is no table at all. + let sql = r#" + SELECT o.order_id, f.value::STRING AS tag, f.key AS k + FROM orders o, + LATERAL FLATTEN(input => o.tags) f + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let resolved = result + .resolved_schema + .expect("resolved schema should be present"); + let tables: Vec<(String, Vec)> = resolved + .tables + .iter() + .map(|table| { + ( + table.name.to_lowercase(), + table.columns.iter().map(|col| col.name.clone()).collect(), + ) + }) + .collect(); + assert_eq!( + tables, + vec![( + "orders".to_string(), + vec!["order_id".to_string(), "tags".to_string()] + )] + ); +} + +#[test] +fn chained_flattens_reach_the_first_input() { + let sql = r#" + SELECT m.value::STRING AS month + FROM events e, + LATERAL FLATTEN(input => PARSE_JSON(e.payload)) a, + LATERAL FLATTEN(input => SPLIT(a.value:"1"::STRING, ',')) m + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "month"), + expected_sources(&["events.payload"]) + ); +} + +#[test] +fn flatten_position_columns_read_nothing() { + // `f.index` numbers the elements of the array: it is not a value of + // orders, and the join to audit feeds nothing. + let sql = r#" + WITH base AS ( + SELECT o.id, o.a, o.b + FROM orders o + JOIN audit k ON o.id = k.id + ), + opts AS ( + SELECT 'n' || (f.index + 1) AS n, b.* + FROM base b, + LATERAL FLATTEN(input => ARRAY_CONSTRUCT(b.a, b.b)) f + ) + SELECT CONCAT_WS('/', id, n) AS label + FROM opts + "#; + let schema = schema_of(vec![ + schema_table(None, None, "orders", &["id", "a", "b"]), + schema_table(None, None, "audit", &["id", "note"]), + ]); + + let result = run_analysis(sql, Dialect::Snowflake, Some(schema)); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "label"), + expected_sources(&["orders.id"]) + ); + + let fed: HashSet = flows_into_columns_of(&stmt, "f") + .into_iter() + .map(|(column, _)| column) + .collect(); + assert_eq!(fed, expected_sources(&["value", "this"])); +} + +#[test] +fn flatten_of_a_literal_reads_nothing() { + let sql = r#" + SELECT o.id, f.value AS step + FROM orders o, + LATERAL FLATTEN(input => ARRAY_CONSTRUCT(1, 2, 3)) f + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert!(flows_into_columns_of(&stmt, "f").is_empty()); + assert!(base_sources_of(&stmt, "step").is_empty()); +} + +#[test] +fn flatten_alias_reused_across_ctes_keeps_each_input() { + // Each CTE flattens its own column under the same alias. The second one + // passes the input positionally. + let sql = r#" + WITH order_tags AS ( + SELECT f.value::STRING AS tag + FROM orders o, + LATERAL FLATTEN(input => o.tags) f + ), + ticket_notes AS ( + SELECT f.value::STRING AS note + FROM tickets t, + LATERAL FLATTEN(t.notes) f + ) + SELECT tag, note + FROM order_tags, ticket_notes + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "tag"), + expected_sources(&["orders.tags"]) + ); + assert_eq!( + base_sources_of(&stmt, "note"), + expected_sources(&["tickets.notes"]) + ); +} + +#[test] +fn star_over_flatten_returns_its_columns() { + let sql = r#" + SELECT * + FROM items i, + LATERAL FLATTEN(input => i.arr) f + "#; + let schema = schema_of(vec![schema_table(None, None, "items", &["id", "arr"])]); + + let result = run_analysis(sql, Dialect::Snowflake, Some(schema)); + let stmt = first_statement(&result); + let output = stmt + .nodes + .iter() + .find(|node| node.node_type == NodeType::Output) + .expect("output node"); + let columns: HashSet = stmt + .edges + .iter() + .filter(|edge| edge.edge_type == EdgeType::Ownership && edge.from == output.id) + .filter_map(|edge| stmt.nodes.iter().find(|node| node.id == edge.to)) + .map(|node| node.label.to_lowercase()) + .collect(); + assert_eq!( + columns, + expected_sources(&["id", "arr", "seq", "key", "path", "index", "value", "this"]) + ); + for carried in ["value", "this"] { + assert_eq!( + base_sources_of(&stmt, carried), + expected_sources(&["items.arr"]), + "sources of {carried}" + ); + } + for position in ["seq", "key", "path", "index"] { + assert!( + base_sources_of(&stmt, position).is_empty(), + "{position} should read nothing" + ); + } +} + +#[test] +fn bare_column_beside_flatten_belongs_to_the_table() { + // No schema: `id` and `status` are not columns of FLATTEN, so they can + // only be columns of orders. + let sql = r#" + SELECT id, f.value AS v + FROM orders o, + LATERAL FLATTEN(input => o.tags) f + WHERE status = 'open' + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "id"), + expected_sources(&["orders.id"]) + ); + assert_eq!( + base_sources_of(&stmt, "v"), + expected_sources(&["orders.tags"]) + ); + let orders = stmt + .nodes + .iter() + .find(|node| node.node_type == NodeType::Table) + .expect("orders node"); + assert!( + orders + .filters + .iter() + .any(|filter| filter.expression.contains("status")), + "{:?}", + orders.filters + ); + assert!( + !issue_codes_list(&result).contains(&issue_codes::UNRESOLVED_REFERENCE.to_string()), + "{:?}", + result.issues + ); +} + +#[test] +fn bare_input_of_a_second_flatten_belongs_to_the_table() { + let sql = r#" + SELECT customer_id, phone.value::STRING AS phone_number + FROM customers, + LATERAL FLATTEN(input => addresses) addr, + LATERAL FLATTEN(input => phone_numbers) phone + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "customer_id"), + expected_sources(&["customers.customer_id"]) + ); + assert_eq!( + base_sources_of(&stmt, "phone_number"), + expected_sources(&["customers.phone_numbers"]) + ); + assert!( + !issue_codes_list(&result).contains(&issue_codes::UNRESOLVED_REFERENCE.to_string()), + "{:?}", + result.issues + ); +} + +#[test] +fn ambiguous_flatten_input_is_reported_once() { + // Without a schema `tags` may belong to either table. VALUE and THIS both + // carry it, yet it is one reference. FLATTEN has no `note`, so it is not + // named among the relations `note` may belong to. + let sql = r#" + SELECT note, f.value AS v + FROM orders o, customers c, + LATERAL FLATTEN(input => tags) f + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + let unresolved: Vec<&str> = result + .issues + .iter() + .filter(|issue| issue.code == issue_codes::UNRESOLVED_REFERENCE) + .map(|issue| issue.message.as_str()) + .collect(); + assert_eq!( + unresolved, + vec![ + "Column 'tags' is ambiguous across tables in scope: CUSTOMERS, ORDERS", + "Column 'note' is ambiguous across tables in scope: CUSTOMERS, ORDERS", + ] + ); +} + +#[test] +fn lateral_table_function_other_than_flatten_is_reported() { + let sql = r#" + SELECT s.value AS part + FROM files t, + LATERAL SPLIT_TO_TABLE(t.csv, ',') s + "#; + + let result = run_analysis(sql, Dialect::Snowflake, None); + assert!( + issue_codes_list(&result).contains(&issue_codes::UNSUPPORTED_SYNTAX.to_string()), + "{:?}", + result.issues + ); +} + +#[test] +fn flatten_adds_no_column_when_column_lineage_is_disabled() { + let sql = r#" + SELECT o.id, f.value AS v + FROM orders o, + LATERAL FLATTEN(input => o.tags) f + "#; + + let result = run_analysis_with_options( + sql, + Dialect::Snowflake, + None, + AnalysisOptions { + enable_column_lineage: Some(false), + ..Default::default() + }, + ); + let stmt = first_statement(&result); + + let tables: Vec = stmt + .nodes + .iter() + .filter(|node| node.node_type == NodeType::Table) + .map(|node| node.label.to_lowercase()) + .collect(); + assert_eq!(tables, vec!["orders".to_string()]); + assert!( + stmt.nodes + .iter() + .all(|node| matches!(node.node_type, NodeType::Table | NodeType::Output)), + "only relations should be emitted: {:?}", + stmt.nodes + ); +} + #[test] fn snowflake_higher_order_functions_track_source() { let sql = r#" diff --git a/crates/flowscope-core/tests/snapshots/snapshots__run_snowflake_snapshot_test@snowflake_lateral_flatten.snap b/crates/flowscope-core/tests/snapshots/snapshots__run_snowflake_snapshot_test@snowflake_lateral_flatten.snap index bd12afa9..46fd0e0b 100644 --- a/crates/flowscope-core/tests/snapshots/snapshots__run_snowflake_snapshot_test@snowflake_lateral_flatten.snap +++ b/crates/flowscope-core/tests/snapshots/snapshots__run_snowflake_snapshot_test@snowflake_lateral_flatten.snap @@ -8,10 +8,80 @@ expression: clean_result "statementIndex": 0, "statementType": "SELECT", "joinCount": 1, - "complexityScore": 20 + "complexityScore": 28 } ], "nodes": [ + { + "id": "column_1f95f87efeb20682", + "type": "column", + "label": "KEY", + "canonicalName": { + "name": "KEY" + }, + "statementIds": [ + 0 + ] + }, + { + "id": "column_268ea64f9e076ed3", + "type": "column", + "label": "cool_ids", + "qualifiedName": "B.cool_ids", + "canonicalName": { + "schema": "B", + "name": "cool_ids" + }, + "statementIds": [ + 0 + ] + }, + { + "id": "column_2f1758652518e9a9", + "type": "column", + "label": "THIS", + "canonicalName": { + "name": "THIS" + }, + "statementIds": [ + 0 + ], + "expression": "FLATTEN(input => b.cool_ids)" + }, + { + "id": "column_4e2ed24f2245005f", + "type": "column", + "label": "VALUE", + "canonicalName": { + "name": "VALUE" + }, + "statementIds": [ + 0 + ], + "expression": "FLATTEN(input => b.cool_ids)" + }, + { + "id": "column_7bfb651f3214d1d7", + "type": "column", + "label": "INDEX", + "canonicalName": { + "name": "INDEX" + }, + "statementIds": [ + 0 + ] + }, + { + "id": "column_8b8a696c1d07ba23", + "type": "column", + "label": "PATH", + "canonicalName": { + "name": "PATH" + }, + "statementIds": [ + 0 + ] + }, { "id": "column_9312b887682fcd76", "type": "column", @@ -23,6 +93,17 @@ expression: clean_result 0 ] }, + { + "id": "column_c2f34f4bce303b03", + "type": "column", + "label": "SEQ", + "canonicalName": { + "name": "SEQ" + }, + "statementIds": [ + 0 + ] + }, { "id": "column_ffa7c99c403f959a", "type": "column", @@ -34,6 +115,18 @@ expression: clean_result 0 ] }, + { + "id": "flatten_c1bebf2d8b83f913", + "type": "cte", + "label": "FLATTEN", + "qualifiedName": "FLATTEN", + "canonicalName": { + "name": "FLATTEN" + }, + "statementIds": [ + 0 + ] + }, { "id": "output_74d22f755ec62dbc", "type": "output", @@ -99,6 +192,42 @@ expression: clean_result } ], "edges": [ + { + "id": "edge_01f4f9bc314d367e", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_1f95f87efeb20682", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_09d6cf174d070d13", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_7bfb651f3214d1d7", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_10b868357df4cc6e", + "from": "column_4e2ed24f2245005f", + "to": "column_9312b887682fcd76", + "type": "data_flow", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_338f68a009f32040", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_2f1758652518e9a9", + "type": "ownership", + "statementIds": [ + 0 + ] + }, { "id": "edge_671cdc7a5a898ccd", "from": "output_74d22f755ec62dbc", @@ -108,6 +237,27 @@ expression: clean_result 0 ] }, + { + "id": "edge_6a0a97f0ea1ad2ab", + "from": "table_14fd3f715fc3ae74", + "to": "column_268ea64f9e076ed3", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_7823aeb8ad256f9d", + "from": "column_268ea64f9e076ed3", + "to": "column_4e2ed24f2245005f", + "type": "derivation", + "expression": "FLATTEN(input => b.cool_ids)", + "joinType": "INNER", + "joinCondition": "b.c_id = a.c_id", + "statementIds": [ + 0 + ] + }, { "id": "edge_7f6d7c7173448b60", "from": "table_14fd3f715fc3ae74", @@ -119,6 +269,15 @@ expression: clean_result 0 ] }, + { + "id": "edge_8ac209680be08157", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_8b8a696c1d07ba23", + "type": "ownership", + "statementIds": [ + 0 + ] + }, { "id": "edge_8d1b60030155f6f4", "from": "output_74d22f755ec62dbc", @@ -127,6 +286,36 @@ expression: clean_result "statementIds": [ 0 ] + }, + { + "id": "edge_9c212231bc72c292", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_c2f34f4bce303b03", + "type": "ownership", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_a4e289214a89eb54", + "from": "column_268ea64f9e076ed3", + "to": "column_2f1758652518e9a9", + "type": "derivation", + "expression": "FLATTEN(input => b.cool_ids)", + "joinType": "INNER", + "joinCondition": "b.c_id = a.c_id", + "statementIds": [ + 0 + ] + }, + { + "id": "edge_cb4d48c39a8c8c26", + "from": "flatten_c1bebf2d8b83f913", + "to": "column_4e2ed24f2245005f", + "type": "ownership", + "statementIds": [ + 0 + ] } ], "issues": [ @@ -135,23 +324,17 @@ expression: clean_result "code": "UNRESOLVED_REFERENCE", "message": "Column 'name' is ambiguous across tables in scope: A, B", "statementIndex": 0 - }, - { - "severity": "warning", - "code": "UNRESOLVED_REFERENCE", - "message": "Column 'value' is ambiguous across tables in scope: A, B", - "statementIndex": 0 } ], "summary": { "statementCount": 1, - "tableCount": 2, - "columnCount": 2, + "tableCount": 3, + "columnCount": 9, "joinCount": 1, - "complexityScore": 20, + "complexityScore": 28, "issueCount": { "errors": 0, - "warnings": 2, + "warnings": 1, "infos": 0 }, "hasErrors": false @@ -196,6 +379,10 @@ expression: clean_result "table": "A", "column": "c_id" } + }, + { + "name": "cool_ids", + "origin": "implied" } ], "origin": "implied", From 56fe139b71a75a4f0bb9f1932d89c8075c59c939 Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:50:59 +0200 Subject: [PATCH 5/8] fix(core): give each derived table its own node when an alias repeats A derived table's node was keyed by statement and alias. Two derived tables of one statement sharing an alias in different scopes, as a generated point in time query writes one `(...) AS v` per satellite, therefore shared one node and one set of column nodes, and each consumer received every producer's lineage: WITH a_intervals AS ( SELECT k, ts FROM (SELECT DISTINCT k, ts FROM sat_a) AS v ), b_intervals AS ( SELECT k, ts FROM (SELECT DISTINCT k, ts FROM sat_b) AS v ) SELECT a.ts AS a_ts, b.ts AS b_ts FROM a_intervals a JOIN b_intervals b ON a.k = b.k a_ts came out derived from sat_a.ts and sat_b.ts both, and so did b_ts, with no issue raised. Count the derived tables of each alias per statement. The first keeps the key it always had, so nothing changes for a statement that does not repeat an alias, and each later one is keyed `alias#n`, getting a node and column nodes of its own. The node keeps the alias as its label. Two derived tables of one alias in the same scope are invalid SQL, and lookups already go through the scope, so a per statement count is enough. A CTE name defined again in a nested WITH merges the same way, but CTEs are also looked up by name once the statement has been walked, so that takes more than a key and is left alone here. Tests: the query above, each output from its own satellite and two nodes labelled v; the same alias over twelve satellites, each output from its own. Both fail without the change. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/flowscope-core/src/analyzer/context.rs | 5 ++ crates/flowscope-core/src/analyzer/visitor.rs | 18 ++++- crates/flowscope-core/tests/lineage_engine.rs | 70 +++++++++++++++++++ 3 files changed, 92 insertions(+), 1 deletion(-) diff --git a/crates/flowscope-core/src/analyzer/context.rs b/crates/flowscope-core/src/analyzer/context.rs index 92ebf609..739ed2f5 100644 --- a/crates/flowscope-core/src/analyzer/context.rs +++ b/crates/flowscope-core/src/analyzer/context.rs @@ -140,6 +140,10 @@ pub(crate) struct StatementContext { pub(crate) table_aliases: HashMap, /// Subquery aliases (for reference tracking) pub(crate) subquery_aliases: HashSet, + /// How many derived tables of each alias the statement has opened, so that + /// a derived table reusing the alias of an earlier one in another scope + /// gets a node of its own. + pub(crate) derived_occurrences: HashMap, /// Last join/operation type for edge labeling pub(crate) last_operation: Option, /// Current join information (type + condition) for edge labeling @@ -235,6 +239,7 @@ impl StatementContext { relation_span_cursors: HashMap::new(), table_aliases: HashMap::new(), subquery_aliases: HashSet::new(), + derived_occurrences: HashMap::new(), last_operation: None, current_join_info: JoinInfo::default(), table_node_ids: HashMap::new(), diff --git a/crates/flowscope-core/src/analyzer/visitor.rs b/crates/flowscope-core/src/analyzer/visitor.rs index acae5ab0..64a2f4e6 100644 --- a/crates/flowscope-core/src/analyzer/visitor.rs +++ b/crates/flowscope-core/src/analyzer/visitor.rs @@ -784,12 +784,28 @@ impl<'a, 'b> Visitor for LineageVisitor<'a, 'b> { // We model derived tables as CTEs in the graph since they are conceptually // similar: both are ephemeral, named result sets scoped to a single query. // This avoids introducing a separate NodeType for a very similar concept. + // + // Two derived tables of one statement can share an alias in + // different scopes, as `(...) AS v` in two CTEs. Keyed by the alias + // alone, the second would find the first's node and its columns + // would feed both. The first keeps the plain key, each later one + // gets its own. let derived_node_id = alias_name.as_ref().map(|name| { + let occurrence = self + .ctx + .derived_occurrences + .entry(name.clone()) + .or_insert(0); + *occurrence += 1; + let key = match *occurrence { + 1 => name.clone(), + n => format!("{name}#{n}"), + }; self.ctx.add_node(Node { id: generate_statement_scoped_node_id( "derived", self.ctx.statement_index, - name, + &key, ), node_type: NodeType::Cte, label: name.clone().into(), diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index 9c4fd315..94f65dd1 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -12942,3 +12942,73 @@ fn oracle_insert_from_view_with_schema() { fn oracle_view_from_view_with_schema() { assert_oracle_no_unresolved_with_schema("view_from_view.sql"); } + +#[test] +fn derived_tables_sharing_an_alias_in_two_ctes_keep_their_own_columns() { + // A generated point in time query opens one `(...) AS v` per satellite. + let sql = r#" + WITH a_intervals AS ( + SELECT k, ts FROM (SELECT DISTINCT k, ts FROM sat_a) AS v + ), + b_intervals AS ( + SELECT k, ts FROM (SELECT DISTINCT k, ts FROM sat_b) AS v + ) + SELECT a.k, a.ts AS a_ts, b.ts AS b_ts + FROM a_intervals a + JOIN b_intervals b ON a.k = b.k + "#; + let schema = schema_of(vec![ + schema_table(None, None, "sat_a", &["k", "ts"]), + schema_table(None, None, "sat_b", &["k", "ts"]), + ]); + + let result = run_analysis(sql, Dialect::Snowflake, Some(schema)); + let stmt = first_statement(&result); + assert_eq!( + base_sources_of(&stmt, "a_ts"), + expected_sources(&["sat_a.ts"]) + ); + assert_eq!( + base_sources_of(&stmt, "b_ts"), + expected_sources(&["sat_b.ts"]) + ); + + let derived: Vec<&Node> = stmt + .nodes + .iter() + .copied() + .filter(|node| node.node_type == NodeType::Cte && &*node.label == "v") + .collect(); + assert_eq!(derived.len(), 2, "one node per derived table"); + assert_ne!(derived[0].id, derived[1].id); +} + +#[test] +fn a_derived_alias_reused_in_twelve_scopes_crosses_nothing() { + let ctes: Vec = (0..12) + .map(|i| format!("c{i} AS (SELECT ts FROM (SELECT DISTINCT ts FROM sat_{i}) AS v)")) + .collect(); + let outputs: Vec = (0..12).map(|i| format!("c{i}.ts AS ts_{i}")).collect(); + let from: Vec = (0..12).map(|i| format!("c{i}")).collect(); + let sql = format!( + "WITH {} SELECT {} FROM {}", + ctes.join(", "), + outputs.join(", "), + from.join(" CROSS JOIN ") + ); + let schema = schema_of( + (0..12) + .map(|i| schema_table(None, None, &format!("sat_{i}"), &["ts"])) + .collect(), + ); + + let result = run_analysis(&sql, Dialect::Snowflake, Some(schema)); + let stmt = first_statement(&result); + for i in 0..12 { + assert_eq!( + base_sources_of(&stmt, &format!("ts_{i}")), + expected_sources(&[&format!("sat_{i}.ts")]), + "ts_{i}" + ); + } +} From ce8eab66455f611e2d26fda3d6e28a1d02f2b933 Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Tue, 29 Sep 2026 00:58:24 +0200 Subject: [PATCH 6/8] fix(core): keep a predicate's subquery and an unaliased derived table out of the outputs A subquery a predicate reads was analysed with no target node, so its select list fell back to whatever the enclosing select targets: the statement's output at the top, and inside a CTE that CTE's column list, which a star over it then returned. SELECT k FROM t WHERE x >= (SELECT MAX(y) AS floor_y FROM u) returned k and floor_y. The same held for IN, for EXISTS, whose `SELECT 1` came back as col_0, and for a subquery in HAVING, a join condition or a grouping expression. No issue said so. Such a subquery is now analysed against a node of its own, labelled `(subquery)` and keyed per occurrence, which owns its columns and their lineage, and what it projects is taken back from the projection buffer once it has been read. A scalar subquery in the select list is not visited there and is returned as before. A derived table without an alias, which Snowflake and others accept, had no node at all: its projection hung on the enclosing target in the same way, and the outer select found no table in scope and reported so. It now gets a node like an aliased one, registered in its scope under a name no query can write, `(derived)`, so that a star over it expands and a bare column resolves. Tests: a scalar comparison, IN, a correlated NOT EXISTS, HAVING and a join condition each return only k; inside a CTE the subquery's column stays on its own node and out of the star; a select list scalar subquery is still returned; an unaliased derived table reads like an aliased one, with no issue. All but the select list case fail without the change. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/flowscope-core/src/analyzer/context.rs | 4 + .../flowscope-core/src/analyzer/expression.rs | 39 ++++++-- crates/flowscope-core/src/analyzer/visitor.rs | 18 +++- crates/flowscope-core/tests/lineage_engine.rs | 88 +++++++++++++++++++ 4 files changed, 139 insertions(+), 10 deletions(-) diff --git a/crates/flowscope-core/src/analyzer/context.rs b/crates/flowscope-core/src/analyzer/context.rs index 739ed2f5..c1b89b3a 100644 --- a/crates/flowscope-core/src/analyzer/context.rs +++ b/crates/flowscope-core/src/analyzer/context.rs @@ -144,6 +144,9 @@ pub(crate) struct StatementContext { /// a derived table reusing the alias of an earlier one in another scope /// gets a node of its own. pub(crate) derived_occurrences: HashMap, + /// How many subqueries a predicate reads the statement has opened, each + /// analysed against a node of its own. + pub(crate) predicate_subqueries: usize, /// Last join/operation type for edge labeling pub(crate) last_operation: Option, /// Current join information (type + condition) for edge labeling @@ -240,6 +243,7 @@ impl StatementContext { table_aliases: HashMap::new(), subquery_aliases: HashSet::new(), derived_occurrences: HashMap::new(), + predicate_subqueries: 0, last_operation: None, current_join_info: JoinInfo::default(), table_node_ids: HashMap::new(), diff --git a/crates/flowscope-core/src/analyzer/expression.rs b/crates/flowscope-core/src/analyzer/expression.rs index aff1f2e9..e912741b 100644 --- a/crates/flowscope-core/src/analyzer/expression.rs +++ b/crates/flowscope-core/src/analyzer/expression.rs @@ -13,10 +13,10 @@ use super::context::{ColumnRef, StatementContext}; use super::functions; -use super::helpers::check_expr_types; +use super::helpers::{check_expr_types, generate_statement_scoped_node_id}; use super::Analyzer; use crate::generated; -use crate::types::{AggregationInfo, FilterClauseType}; +use crate::types::{AggregationInfo, FilterClauseType, Node, NodeType}; use crate::Dialect; use sqlparser::ast::{self, Expr, FunctionArg, FunctionArgExpr}; use std::collections::HashSet; @@ -24,6 +24,10 @@ use std::sync::Arc; #[cfg(feature = "tracing")] use tracing::debug; +/// The label of the node a subquery read by a predicate is analyzed against. A +/// query cannot write it, so nothing can qualify a column with it. +const PREDICATE_SUBQUERY: &str = "(subquery)"; + /// Maximum recursion depth for expression traversal to prevent stack overflow /// on maliciously crafted or deeply nested SQL expressions. pub(super) const MAX_RECURSION_DEPTH: usize = 100; @@ -135,6 +139,29 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { } } + /// Analyzes a subquery an expression other than a projection reads: a + /// predicate compares with it or tests it for rows, and a join condition or + /// a grouping expression reads it the same way. What it projects is never + /// returned, so it is analyzed against a node of its own, labelled + /// `(subquery)`, and its columns are taken back once it has been read. + /// Analyzed with no target, they would hang on whatever the enclosing + /// select targets, the statement's output at the top, as if the query + /// returned them, and would join the column list of an enclosing CTE. + fn analyze_predicate_subquery(&mut self, query: &ast::Query) { + let checkpoint = self.ctx.projection_checkpoint(); + self.ctx.predicate_subqueries += 1; + let key = format!("{PREDICATE_SUBQUERY}#{}", self.ctx.predicate_subqueries); + let node_id = self.ctx.add_node(Node { + id: generate_statement_scoped_node_id("subquery", self.ctx.statement_index, &key), + node_type: NodeType::Cte, + label: PREDICATE_SUBQUERY.into(), + qualified_name: Some(PREDICATE_SUBQUERY.into()), + ..Default::default() + }); + self.analyzer.analyze_query(self.ctx, query, Some(&node_id)); + self.ctx.take_output_columns_since(checkpoint); + } + /// Recursively visits an expression to find and analyze subqueries. /// /// The `depth` parameter tracks recursion depth to prevent stack overflow @@ -153,11 +180,9 @@ impl<'a, 'b> ExpressionAnalyzer<'a, 'b> { let next_depth = depth + 1; match expr { - Expr::Subquery(query) => self.analyzer.analyze_query(self.ctx, query, None), - Expr::InSubquery { subquery, .. } => { - self.analyzer.analyze_query(self.ctx, subquery, None) - } - Expr::Exists { subquery, .. } => self.analyzer.analyze_query(self.ctx, subquery, None), + Expr::Subquery(query) => self.analyze_predicate_subquery(query), + Expr::InSubquery { subquery, .. } => self.analyze_predicate_subquery(subquery), + Expr::Exists { subquery, .. } => self.analyze_predicate_subquery(subquery), Expr::BinaryOp { .. } => { for operand in binary_op_operands(expr) { self.visit_expression_for_subqueries(operand, next_depth); diff --git a/crates/flowscope-core/src/analyzer/visitor.rs b/crates/flowscope-core/src/analyzer/visitor.rs index 64a2f4e6..89e1a2b9 100644 --- a/crates/flowscope-core/src/analyzer/visitor.rs +++ b/crates/flowscope-core/src/analyzer/visitor.rs @@ -21,6 +21,9 @@ use sqlparser::ast::{ }; use std::sync::Arc; +/// The name a derived table without an alias is registered under in its scope. +const UNALIASED_DERIVED_TABLE: &str = "(derived)"; + /// A visitor trait for traversing the SQL AST. /// /// This trait defines default behavior for visiting nodes (traversing children). @@ -775,11 +778,20 @@ impl<'a, 'b> Visitor for LineageVisitor<'a, 'b> { // We create a node for it in the graph, analyze its subquery to determine its // output columns, and then register its alias and columns in the current scope // so the outer query can reference it. - let alias_name = alias.as_ref().map(|a| a.name.to_string()); let projection_checkpoint = self.ctx.projection_checkpoint(); - let derived_span = alias_name + let derived_span = alias .as_ref() - .and_then(|name| self.locate_derived_alias_span(name)); + .and_then(|a| self.locate_derived_alias_span(&a.name.to_string())); + // A derived table without an alias, which Snowflake and others + // accept, is a relation the outer query reads all the same. + // Without a node its projection would hang on whatever the + // enclosing select targets, the statement's output at the top, and + // the outer select would find no table in scope. It is named so + // that no query can qualify a column with it. + let alias_name = Some(alias.as_ref().map_or_else( + || UNALIASED_DERIVED_TABLE.to_string(), + |a| a.name.to_string(), + )); // We model derived tables as CTEs in the graph since they are conceptually // similar: both are ephemeral, named result sets scoped to a single query. diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index 94f65dd1..b08d773e 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -13012,3 +13012,91 @@ fn a_derived_alias_reused_in_twelve_scopes_crosses_nothing() { ); } } + +/// Labels of the columns the statement returns, lowercased and sorted. +fn returned_columns(stmt: &StmtView<'_>) -> Vec { + let output_ids: HashSet<&str> = stmt + .nodes + .iter() + .filter(|node| node.node_type == NodeType::Output) + .map(|node| &*node.id) + .collect(); + let mut labels: Vec = stmt + .edges + .iter() + .filter(|edge| edge.edge_type == EdgeType::Ownership && output_ids.contains(&*edge.from)) + .filter_map(|edge| stmt.nodes.iter().find(|node| node.id == edge.to)) + .map(|node| node.label.to_lowercase()) + .collect(); + labels.sort(); + labels.dedup(); + labels +} + +fn predicate_schema() -> SchemaMetadata { + schema_of(vec![ + schema_table(None, None, "t", &["k", "x"]), + schema_table(None, None, "u", &["k", "y"]), + ]) +} + +#[test] +fn a_subquery_a_predicate_reads_returns_nothing() { + for sql in [ + "SELECT k FROM t WHERE x >= (SELECT MAX(y) AS floor_y FROM u)", + "SELECT k FROM t WHERE x IN (SELECT y FROM u)", + "SELECT k FROM t WHERE NOT EXISTS (SELECT 1 FROM u WHERE u.k = t.k)", + "SELECT k FROM t GROUP BY k HAVING MAX(x) > (SELECT MAX(y) FROM u)", + "SELECT t.k FROM t JOIN u ON u.k = t.k AND u.y IN (SELECT y FROM u)", + ] { + let result = run_analysis(sql, Dialect::Snowflake, Some(predicate_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["k"], "{sql}"); + } +} + +#[test] +fn a_subquery_a_predicate_reads_keeps_its_lineage_on_its_own_node() { + // Inside a CTE, the subquery's column must not join the CTE's columns + // either, or the star below would return it. + let sql = "WITH c AS (SELECT k FROM t WHERE x IN (SELECT y FROM u)) SELECT * FROM c"; + let result = run_analysis(sql, Dialect::Snowflake, Some(predicate_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["k"]); + let flows: Vec = flows_into_columns_of(&stmt, "(subquery)") + .into_iter() + .map(|(column, _)| column) + .collect(); + assert_eq!(flows, vec!["y"]); +} + +#[test] +fn a_scalar_subquery_in_the_select_list_is_still_returned() { + let sql = "SELECT k, (SELECT MAX(y) FROM u) AS m FROM t"; + let result = run_analysis(sql, Dialect::Snowflake, Some(predicate_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["k", "m"]); +} + +#[test] +fn a_derived_table_without_an_alias_is_read_like_one_with_it() { + let schema = schema_of(vec![schema_table(None, None, "orders", &["id", "amount"])]); + for sql in [ + "SELECT * FROM (SELECT id, amount FROM orders) WHERE amount > 0", + "SELECT * FROM (SELECT id, amount FROM orders) AS o WHERE amount > 0", + ] { + let result = run_analysis(sql, Dialect::Snowflake, Some(schema.clone())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["amount", "id"], "{sql}"); + assert_eq!( + base_sources_of(&stmt, "amount"), + expected_sources(&["orders.amount"]), + "{sql}" + ); + assert!( + issue_codes_list(&result).is_empty(), + "{sql}: {:?}", + issue_codes_list(&result) + ); + } +} From 124bb84f7412c2412e40a784871702112101a928 Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Tue, 29 Sep 2026 01:03:39 +0200 Subject: [PATCH 7/8] fix(core): honour EXCLUDE, EXCEPT and RENAME in a wildcard's expansion expand_wildcard emitted every column of the relation whatever options the wildcard carried: `SELECT * EXCLUDE (b) FROM p` returned b, and `SELECT * RENAME (b AS bee) FROM p` returned b under its old name. The common idiom of replacing a column fared worse: SELECT * EXCLUDE (b), UPPER(b) AS b FROM p The star copy of b and the computed b shared one output node, the copy came first, and the UPPER(b) derivation was merged away: b came out as a plain copy of p.b. The wildcard's options now reach the expansion. A column named by EXCLUDE, or by EXCEPT where the dialect writes it so, is not expanded; a column named by RENAME is expanded under its new name, still read from the old one. Names are compared as the expanded columns' own names are. REPLACE and ILIKE are still not read. A plain `SELECT *, f(x) AS x`, which does write two columns named x, is left as it was. Tests: EXCLUDE with and without parentheses, qualified and inside a CTE; EXCEPT; RENAME with and without parentheses, read from the old column; the idiom above, whose b carries UPPER(b). All fail without the change. Co-Authored-By: Claude Opus 5.5 (1M context) --- crates/flowscope-core/src/analyzer/query.rs | 50 +++++++++++- .../src/analyzer/select_analyzer.rs | 13 +++- crates/flowscope-core/tests/lineage_engine.rs | 76 +++++++++++++++++++ 3 files changed, 134 insertions(+), 5 deletions(-) diff --git a/crates/flowscope-core/src/analyzer/query.rs b/crates/flowscope-core/src/analyzer/query.rs index 7f2af145..9a233405 100644 --- a/crates/flowscope-core/src/analyzer/query.rs +++ b/crates/flowscope-core/src/analyzer/query.rs @@ -555,7 +555,47 @@ impl<'a> Analyzer<'a> { ctx: &mut StatementContext, table_qualifier: Option<&str>, target_node: Option<&str>, + options: &ast::WildcardAdditionalOptions, ) { + // What `* EXCLUDE (...)` or `* EXCEPT (...)` leaves out, and what + // `* RENAME (old AS new)` returns under another name, compared as the + // expanded columns' own names are. REPLACE and ILIKE are not read. + let mut excluded: HashSet = HashSet::new(); + match &options.opt_exclude { + Some(ast::ExcludeSelectItem::Single(ident)) => { + excluded.insert(self.normalize_identifier(&ident.to_string())); + } + Some(ast::ExcludeSelectItem::Multiple(idents)) => { + excluded.extend( + idents + .iter() + .map(|ident| self.normalize_identifier(&ident.to_string())), + ); + } + None => {} + } + if let Some(except) = &options.opt_except { + excluded.extend( + std::iter::once(&except.first_element) + .chain(&except.additional_elements) + .map(|ident| self.normalize_identifier(&ident.to_string())), + ); + } + let renames: Vec<&ast::IdentWithAlias> = match &options.opt_rename { + Some(ast::RenameSelectItem::Single(rename)) => vec![rename], + Some(ast::RenameSelectItem::Multiple(renames)) => renames.iter().collect(), + None => Vec::new(), + }; + let renamed: HashMap = renames + .into_iter() + .map(|rename| { + ( + self.normalize_identifier(&rename.ident.to_string()), + rename.alias.value.clone(), + ) + }) + .collect(); + // Resolve wildcard sources as (canonical, qualifier) pairs so repeated // relation instances in self-joins are expanded independently. let tables_to_expand: Vec<(String, String)> = if let Some(qualifier) = table_qualifier { @@ -615,13 +655,21 @@ impl<'a> Analyzer<'a> { if let Some(columns) = columns_to_add { // Expand from schema - NOT approximate. for col_info in columns { + let normalized = self.normalize_identifier(&col_info.name); + if excluded.contains(&normalized) { + continue; + } + let name = renamed + .get(&normalized) + .cloned() + .unwrap_or_else(|| col_info.name.clone()); let sources = vec![ColumnRef { table: Some(source_qualifier.clone()), column: col_info.name.clone(), }]; self.add_output_column( ctx, - &col_info.name, + &name, sources, None, col_info.data_type, diff --git a/crates/flowscope-core/src/analyzer/select_analyzer.rs b/crates/flowscope-core/src/analyzer/select_analyzer.rs index a1b55bfe..a14d5856 100644 --- a/crates/flowscope-core/src/analyzer/select_analyzer.rs +++ b/crates/flowscope-core/src/analyzer/select_analyzer.rs @@ -230,7 +230,7 @@ impl<'a, 'b> SelectAnalyzer<'a, 'b> { }, ); } - SelectItem::QualifiedWildcard(name, _) => { + SelectItem::QualifiedWildcard(name, options) => { let qualifier = name.to_string(); // SelectItemQualifiedWildcardKind::Display appends ".*" let qualifier = qualifier.strip_suffix(".*").unwrap_or(&qualifier); @@ -238,11 +238,16 @@ impl<'a, 'b> SelectAnalyzer<'a, 'b> { self.ctx, Some(qualifier), self.target_node.as_deref(), + options, ); } - SelectItem::Wildcard(_) => { - self.analyzer - .expand_wildcard(self.ctx, None, self.target_node.as_deref()); + SelectItem::Wildcard(options) => { + self.analyzer.expand_wildcard( + self.ctx, + None, + self.target_node.as_deref(), + options, + ); } } } diff --git a/crates/flowscope-core/tests/lineage_engine.rs b/crates/flowscope-core/tests/lineage_engine.rs index b08d773e..27cea89a 100644 --- a/crates/flowscope-core/tests/lineage_engine.rs +++ b/crates/flowscope-core/tests/lineage_engine.rs @@ -13100,3 +13100,79 @@ fn a_derived_table_without_an_alias_is_read_like_one_with_it() { ); } } + +fn star_options_schema() -> SchemaMetadata { + schema_of(vec![schema_table(None, None, "p", &["a", "b", "c"])]) +} + +#[test] +fn star_exclude_leaves_the_excluded_columns_out() { + for sql in [ + "SELECT * EXCLUDE (b) FROM p", + "SELECT * EXCLUDE b FROM p", + "SELECT p.* EXCLUDE (b) FROM p", + "WITH c AS (SELECT * EXCLUDE (b) FROM p) SELECT * FROM c", + ] { + let result = run_analysis(sql, Dialect::Snowflake, Some(star_options_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["a", "c"], "{sql}"); + } +} + +#[test] +fn a_column_excluded_and_computed_again_keeps_its_expression() { + // With the star copy of b still expanded, the two projections named b + // shared one node, the copy came first, and the expression was lost. + let sql = "SELECT * EXCLUDE (b), UPPER(b) AS b FROM p"; + let result = run_analysis(sql, Dialect::Snowflake, Some(star_options_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["a", "b", "c"]); + assert_eq!(base_sources_of(&stmt, "b"), expected_sources(&["p.b"])); + let output_b = stmt + .nodes + .iter() + .find(|node| { + node.node_type == NodeType::Column + && node.label.eq_ignore_ascii_case("b") + && stmt.edges.iter().any(|edge| { + edge.edge_type == EdgeType::Ownership + && edge.to == node.id + && stmt.nodes.iter().any(|owner| { + owner.id == edge.from && owner.node_type == NodeType::Output + }) + }) + }) + .expect("b is returned"); + let expressions: Vec> = stmt + .edges + .iter() + .filter(|edge| edge.to == output_b.id && edge.edge_type != EdgeType::Ownership) + .map(|edge| edge.expression.as_deref()) + .collect(); + assert_eq!(expressions, vec![Some("UPPER(b)")]); +} + +#[test] +fn star_except_leaves_the_excepted_columns_out() { + let sql = "SELECT * EXCEPT (b) FROM p"; + let result = run_analysis(sql, Dialect::Bigquery, Some(star_options_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["a", "c"]); +} + +#[test] +fn star_rename_returns_the_column_under_its_new_name() { + for sql in [ + "SELECT * RENAME (b AS bee) FROM p", + "SELECT * RENAME b AS bee FROM p", + ] { + let result = run_analysis(sql, Dialect::Snowflake, Some(star_options_schema())); + let stmt = first_statement(&result); + assert_eq!(returned_columns(&stmt), vec!["a", "bee", "c"], "{sql}"); + assert_eq!( + base_sources_of(&stmt, "bee"), + expected_sources(&["p.b"]), + "{sql}" + ); + } +} From b0514c714d8e2d53a8723c74b562f385362cc7fd Mon Sep 17 00:00:00 2001 From: William Horel <80535383+IL-William@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:02:45 +0200 Subject: [PATCH 8/8] docs(changelog): note the relation and scope lineage fixes Co-Authored-By: Claude Opus 5.5 (1M context) --- CHANGELOG.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index fe5380d7..ce95e924 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Walked left-deep chains of binary operators as chains, so that a generated surrogate key over more than about 100 operands keeps every operand in its lineage instead of losing the leftmost ones and reporting `APPROXIMATE_LINEAGE`. - Read columns inside TRIM, SUBSTRING, POSITION, CEIL and FLOOR, AT TIME ZONE, IS [NOT] DISTINCT FROM, path and subscript accessors, COLLATE, OVERLAY, SIMILAR TO, RLIKE, ANY and ALL as lineage sources. +- Read `LATERAL FLATTEN` as a relation with the SEQ, KEY, PATH, INDEX, VALUE and THIS columns, VALUE and THIS derived from its input, instead of as a table named after its alias. +- Kept two derived tables that share an alias in one statement apart, so that each consumer receives its own derived table's lineage. +- Kept what a subquery read by a predicate, a join condition or a grouping expression projects out of the statement's outputs, and gave a derived table without an alias a node of its own. +- Honoured `* EXCLUDE`, `* EXCEPT` and `* RENAME` in wildcard expansion. ## [0.9.2] - 2026-09-24