Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,9 +60,12 @@ a `Concrete` origin naming an invented table. `Concrete` is a claim that the
column really comes from that table, so the resolver only emits it once it has
one. `Ambiguous` is reserved for a genuine choice between two or more known
relations, which is the only case a `CatalogProvider` is asked to settle.
`ColumnOrigin` is deliberately not `#[non_exhaustive]`: a consumer that starts
silently ignoring a new resolution state is the failure the enum exists to
prevent, so a new variant should break their build.
No public enum is `#[non_exhaustive]`. A consumer that starts silently
ignoring a new `ColumnOrigin` is the failure that enum exists to prevent, and
the same reasoning holds for the rest: while the crate is pre-1.0, a variant
that changes what a consumer should do is worth a compile error rather than a
`_` arm that swallows it. This matches the rule for sqlparser AST variants
above.

`resolve/topo.rs` validates that the graph is a DAG after removing recursive
CTE back-edges. `resolve/catalog.rs` applies the optional `CatalogProvider`
Expand Down
73 changes: 71 additions & 2 deletions sqllineage/src/build/expr.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use sqlparser::ast::{self, Expr, FunctionArguments, WindowType};
use sqlparser::ast::{self, AccessExpr, Expr, FunctionArguments, Subscript, WindowType};

use crate::build::LineageBuilder;
use crate::build::select::split_compound;
Expand Down Expand Up @@ -199,7 +199,9 @@ impl LineageBuilder {
v
}

Expr::CompoundFieldAccess { root, .. } => self.collect_ancestors(root),
Expr::CompoundFieldAccess { root, access_chain } => {
self.collect_compound_field_ancestors(root, access_chain)
}
Expr::JsonAccess { value, .. } => self.collect_ancestors(value),

Expr::Function(func) => {
Expand Down Expand Up @@ -301,6 +303,73 @@ impl LineageBuilder {
| Expr::MatchAgainst { .. } => vec![],
}
}

/// Collect the physical column at the root of a structured access chain.
///
/// `base.items[0]` is ambiguous at the syntax level: `base` can be a
/// visible relation binding, in which case `items` is its physical column,
/// or it can be an unqualified top-level column (`payload.items[0]`). The
/// scope binding, rather than rendered SQL text or dialect-specific names,
/// is the structural distinction between those cases.
fn collect_compound_field_ancestors(
&mut self,
root: &Expr,
access_chain: &[AccessExpr],
) -> Vec<NodeId> {
let mut ancestors = match (root, access_chain.first()) {
(Expr::Identifier(binding), Some(AccessExpr::Dot(Expr::Identifier(field))))
if self
.graph
.scopes
.lookup(self.current_scope, &binding.value)
.is_some() =>
{
let node = self.graph.add_ref(
field.value.clone(),
Some(binding.value.clone()),
self.current_scope,
);
vec![node]
}
(Expr::Identifier(column), Some(AccessExpr::Dot(Expr::Identifier(_)))) => {
vec![
self.graph
.add_unqualified(column.value.clone(), self.current_scope),
]
}
_ => self.collect_ancestors(root),
};

for access in access_chain {
if let AccessExpr::Subscript(subscript) = access {
ancestors.extend(self.collect_subscript_ancestors(subscript));
}
}
ancestors
}

fn collect_subscript_ancestors(&mut self, subscript: &Subscript) -> Vec<NodeId> {
match subscript {
Subscript::Index { index } => self.collect_ancestors(index),
Subscript::Slice {
lower_bound,
upper_bound,
stride,
} => {
let mut ancestors = Vec::new();
if let Some(lower) = lower_bound {
ancestors.extend(self.collect_ancestors(lower));
}
if let Some(upper) = upper_bound {
ancestors.extend(self.collect_ancestors(upper));
}
if let Some(step) = stride {
ancestors.extend(self.collect_ancestors(step));
}
ancestors
}
}
}
}

/// Classify what an expression does to the values it reads.
Expand Down
7 changes: 7 additions & 0 deletions sqllineage/src/resolve/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,13 @@ fn expand_star(
mappings,
visited_scopes,
);
} else if let Some(Binding::Table(actual_table)) = binding {
// `t` is how the star was written, which for `SELECT a.* FROM
// real AS a` is the alias. The wildcard has to name the relation
// the alias stands for: the alias appears in no catalog and in no
// `tables.inputs`, so naming it here would both block expansion
// and claim a relation the table graph says does not exist.
mappings.push(wildcard_mapping(output_table, actual_table));
} else {
mappings.push(wildcard_mapping(output_table, t.clone()));
}
Expand Down
7 changes: 4 additions & 3 deletions sqllineage/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -222,10 +222,11 @@ impl Default for AnalyzeOptions {
/// dialect-specific rules such as `BigQuery`'s backtick quoting or T-SQL's
/// bracket quoting.
///
/// Marked `#[non_exhaustive]`: `sqlparser` gains dialects over time, and
/// adding one here should not be a breaking change for downstream matches.
/// Exhaustive on purpose. `sqlparser` gains dialects over time and each one
/// added here breaks a downstream match, which is the point: a consumer that
/// maps dialects to its own vocabulary should be told about a new one rather
/// than route it silently through a `_` arm.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum Dialect {
#[default]
Generic,
Expand Down
37 changes: 37 additions & 0 deletions sqllineage/tests/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -80,6 +80,43 @@ fn select_star_without_catalog_preserved() {
}
}

/// A star written on an alias names the relation the alias stands for. The
/// alias itself appears in no catalog and in no `tables.inputs`, so naming it
/// would both block expansion and claim a relation that does not exist.
#[test]
fn qualified_alias_star_names_the_aliased_relation() {
let result = analyze("SELECT u.* FROM users AS u", AnalyzeOptions::default())
.expect("parse")
.into_iter()
.next()
.unwrap();
assert_eq!(result.columns.mappings.len(), 1);
match &result.columns.mappings[0].sources[0] {
ColumnOrigin::Wildcard { table } => assert_eq!(table.table, "users"),
other => panic!("expected Wildcard, got {other:?}"),
}
}

#[test]
fn qualified_alias_star_expands_from_catalog() {
for sql in [
"SELECT u.* FROM users AS u",
"WITH x AS (SELECT u.* FROM users AS u) SELECT * FROM x",
] {
let result = analyze(sql, opts_with_catalog())
.expect("parse")
.into_iter()
.next()
.unwrap();
assert_eq!(result.columns.mappings.len(), 3, "{sql}");
assert_eq!(
concrete_sources(find_mapping(&result.columns.mappings, "email")),
vec![("users".into(), "email".into())],
"{sql}"
);
}
}

#[test]
fn ambiguous_column_resolved_by_catalog() {
let sql = "SELECT name FROM users JOIN orders ON users.id = orders.user_id";
Expand Down
81 changes: 80 additions & 1 deletion sqllineage/tests/expr_coverage.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,21 @@
mod common;

use common::{analyze_one, concrete_sources, find_mapping};
use sqllineage::TransformKind;
use sqllineage::{AnalyzeOptions, Dialect, TransformKind, analyze};

fn analyze_with_dialect(sql: &str, dialect: Dialect) -> sqllineage::AnalyzeResult {
analyze(
sql,
AnalyzeOptions {
dialect,
..AnalyzeOptions::default()
},
)
.expect("SQL should parse")
.into_iter()
.next()
.unwrap_or_default()
}

#[test]
fn extract_year() {
Expand Down Expand Up @@ -107,3 +121,68 @@ fn json_access() {
let m = find_mapping(&result.columns.mappings, "val");
assert_eq!(concrete_sources(m), vec![("t".into(), "data".into())]);
}

#[test]
fn qualified_compound_field_access_uses_binding_column() {
let result = analyze_one("SELECT base.items_array[1] AS item FROM actual_table AS base");
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(
concrete_sources(m),
vec![("actual_table".into(), "items_array".into())]
);
}

#[test]
fn compound_field_access_retains_column_dependent_index() {
let result = analyze_one("SELECT base.items_array[idx] AS item FROM actual_table AS base");
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(
concrete_sources(m),
vec![
("actual_table".into(), "idx".into()),
("actual_table".into(), "items_array".into()),
]
);
}

#[test]
fn nested_qualified_compound_field_access_keeps_top_level_column() {
let result = analyze_one("SELECT base.payload.items[1] AS item FROM actual_table AS base");
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(
concrete_sources(m),
vec![("actual_table".into(), "payload".into())]
);
}

#[test]
fn cte_compound_field_access_uses_cte_binding_column() {
let result = analyze_one(
"WITH base AS (SELECT items_array FROM actual_table) SELECT base.items_array[1] AS item FROM base",
);
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(
concrete_sources(m),
vec![("actual_table".into(), "items_array".into())]
);
}

#[test]
fn unqualified_compound_field_access_uses_top_level_column() {
let result = analyze_one("SELECT payload.items[1] AS item FROM t");
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(concrete_sources(m), vec![("t".into(), "payload".into())]);
}

#[test]
fn bigquery_offset_compound_field_access_uses_binding_column() {
let result = analyze_with_dialect(
"SELECT base.items_array[OFFSET(0)] AS item FROM actual_table AS base",
Dialect::BigQuery,
);
let m = find_mapping(&result.columns.mappings, "item");
assert_eq!(
concrete_sources(m),
vec![("actual_table".into(), "items_array".into())]
);
}
Loading