use crate::facts::FactsDb;
use crate::{Options, Result};
#[must_use]
pub fn source_table(opts: &Options) -> &'static str {
if opts.time_bucket.is_some() {
"changes_bucketed"
} else if opts.use_canonical_lineage {
"changes_lineage"
} else {
"changes"
}
}
pub fn materialize_if_needed(db: &FactsDb, opts: &Options) -> Result<()> {
if opts.use_canonical_lineage && opts.time_bucket.is_none() {
crate::facts::ingest::materialize_changes_lineage(db)?;
}
Ok(())
}
pub fn materialize_source(db: &FactsDb, opts: &Options) -> Result<()> {
if let Some(bucket) = opts.time_bucket {
crate::facts::ingest::materialize_changes_bucketed(db, bucket, opts.use_canonical_lineage)?;
} else if opts.use_canonical_lineage {
crate::facts::ingest::materialize_changes_lineage(db)?;
}
Ok(())
}
#[must_use]
pub fn rewrite(sql: &str, opts: &Options) -> String {
use std::sync::OnceLock;
static RE: OnceLock<regex::Regex> = OnceLock::new();
let src = source_table(opts);
if src == "changes" {
return sql.to_string();
}
let re =
RE.get_or_init(|| regex::Regex::new(r"(?i)\b(FROM|JOIN)\s+changes\b(\s*)(\w*)").unwrap());
re.replace_all(sql, |caps: ®ex::Captures<'_>| {
let kw = &caps[1];
let ws = &caps[2];
let next = &caps[3];
let needs_alias = next.is_empty() || is_sql_keyword(next);
if needs_alias {
format!("{kw} {src} AS changes{ws}{next}")
} else {
format!("{kw} {src}{ws}{next}")
}
})
.into_owned()
}
fn is_sql_keyword(token: &str) -> bool {
const KEYWORDS: &[&str] = &[
"WHERE",
"GROUP",
"HAVING",
"ORDER",
"LIMIT",
"OFFSET",
"JOIN",
"INNER",
"LEFT",
"RIGHT",
"FULL",
"OUTER",
"CROSS",
"NATURAL",
"ON",
"USING",
"UNION",
"INTERSECT",
"EXCEPT",
"WINDOW",
"QUALIFY",
"FETCH",
"SAMPLE",
"TABLESAMPLE",
"AS",
"WITH",
"ANTI",
"SEMI",
"ASOF",
];
let upper = token.to_ascii_uppercase();
KEYWORDS.contains(&upper.as_str())
}
#[cfg(test)]
mod tests {
use super::*;
fn opts_with(use_lineage: bool) -> Options {
Options {
use_canonical_lineage: use_lineage,
..Options::default()
}
}
#[test]
fn rewrite_adds_alias_when_no_existing_alias() {
let sql = "SELECT path FROM changes\nGROUP BY path";
let out = rewrite(sql, &opts_with(true));
assert!(out.contains("FROM changes_lineage AS changes"));
}
#[test]
fn rewrite_preserves_existing_alias() {
let sql = "SELECT c.path FROM changes c GROUP BY c.path";
let out = rewrite(sql, &opts_with(true));
assert!(
out.contains("FROM changes_lineage c"),
"existing alias `c` must survive: {out}"
);
assert!(!out.contains("AS changes c"));
}
#[test]
fn rewrite_join_with_qualified_refs() {
let sql = "SELECT a FROM commits INNER JOIN changes ON changes.rev = commits.rev";
let out = rewrite(sql, &opts_with(true));
assert!(out.contains("INNER JOIN changes_lineage AS changes ON"));
assert!(out.contains("changes.rev = commits.rev"));
}
#[test]
fn rewrite_handles_lowercase_sql_keywords() {
let sql = "select path from changes\ngroup by path";
let out = rewrite(sql, &opts_with(true));
let lower = out.to_lowercase();
assert!(
lower.contains("from changes_lineage as changes"),
"lowercase `from changes` must rewrite + add alias: {out}"
);
assert!(
lower.contains("group by path"),
"`group by` keyword sequence must survive verbatim: {out}"
);
assert!(
!lower.contains("from changes_lineage group"),
"lowercase `group` must NOT be treated as an alias: {out}"
);
}
#[test]
fn rewrite_lowercase_alias_preserved_in_lowercase_sql() {
let sql = "select c.path from changes c group by c.path";
let out = rewrite(sql, &opts_with(true));
assert!(
out.to_lowercase().contains("from changes_lineage c"),
"lowercase alias `c` must survive: {out}"
);
assert!(!out.contains("AS changes c"));
}
#[test]
fn rewrite_leaves_changes_bucketed_alone() {
let sql = "SELECT path FROM changes_bucketed GROUP BY path";
let out = rewrite(sql, &opts_with(true));
assert_eq!(out, sql, "must not touch identifiers like changes_bucketed");
}
#[test]
fn rewrite_noop_when_lineage_off() {
let sql = "SELECT path FROM changes\nGROUP BY path";
let out = rewrite(sql, &opts_with(false));
assert_eq!(out, sql);
}
}