use crate::drift::{SqlBaseType, json_schema_base_type};
use serde_json::Value;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PlannedColumn {
pub name: String,
pub base_type: SqlBaseType,
pub nullable: bool,
}
pub fn plan_columns(page: &[Value]) -> Option<Vec<PlannedColumn>> {
let schema = crate::schema::infer_schema(page);
let props = schema.get("properties")?.as_object()?;
if props.is_empty() {
return None;
}
let mut ordered: Vec<String> = props.keys().cloned().collect();
ordered.sort_unstable();
let columns: Vec<PlannedColumn> = ordered
.into_iter()
.map(|name| {
let fragment = &props[&name];
PlannedColumn {
base_type: json_schema_base_type(fragment).unwrap_or(SqlBaseType::Text),
nullable: true,
name,
}
})
.collect();
(!columns.is_empty()).then_some(columns)
}
pub fn render_columns<Q, T>(columns: &[PlannedColumn], quote: Q, ty: T) -> String
where
Q: Fn(&str) -> String,
T: Fn(SqlBaseType) -> &'static str,
{
columns
.iter()
.map(|c| format!("{} {}", quote(&c.name), ty(c.base_type)))
.collect::<Vec<_>>()
.join(", ")
}
pub fn plan_keyed_columns(page: &[Value], key: &[String]) -> Option<Vec<PlannedColumn>> {
let mut columns = plan_columns(page)?;
for k in key {
match columns.iter_mut().find(|c| &c.name == k) {
Some(c) => c.nullable = false,
None => columns.push(PlannedColumn {
name: k.clone(),
base_type: SqlBaseType::Text,
nullable: false,
}),
}
}
Some(columns)
}
pub fn render_column_defs<Q, T>(columns: &[PlannedColumn], quote: Q, ty: T) -> String
where
Q: Fn(&str) -> String,
T: Fn(&PlannedColumn) -> &'static str,
{
columns
.iter()
.map(|c| format!("{} {}", quote(&c.name), ty(c)))
.collect::<Vec<_>>()
.join(", ")
}
pub fn render_primary_key<Q>(key: &[String], quote: Q) -> Option<String>
where
Q: Fn(&str) -> String,
{
(!key.is_empty()).then(|| {
format!(
"PRIMARY KEY ({})",
key.iter().map(|k| quote(k)).collect::<Vec<_>>().join(", ")
)
})
}
pub fn missing_target_error(connector: &str, target: &str) -> crate::error::FaucetError {
crate::error::FaucetError::Sink(format!(
"{connector}: target `{target}` does not exist and `create_table: false`. \
Create it first, or set `create_table: true` to have faucet create it from \
the first page's inferred schema."
))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn plans_every_column_with_its_inferred_type() {
let page = vec![
json!({ "id": 1, "name": "a", "amount": 1.5, "ok": true }),
json!({ "id": 2, "name": "b", "amount": 2.5, "ok": false }),
];
let cols = plan_columns(&page).expect("a plan");
let by_name: std::collections::HashMap<&str, SqlBaseType> = cols
.iter()
.map(|c| (c.name.as_str(), c.base_type))
.collect();
assert_eq!(by_name["id"], SqlBaseType::Integer);
assert_eq!(by_name["name"], SqlBaseType::Text);
assert_eq!(by_name["amount"], SqlBaseType::Double);
assert_eq!(by_name["ok"], SqlBaseType::Boolean);
}
#[test]
fn column_order_is_deterministic_regardless_of_record_key_order() {
let a = plan_columns(&[json!({ "z": 1, "a": 2, "m": 3 })]).expect("a");
let b = plan_columns(&[json!({ "a": 2, "m": 3, "z": 1 })]).expect("b");
assert_eq!(a, b, "the same record set must plan the same table");
assert_eq!(
a.iter().map(|c| c.name.as_str()).collect::<Vec<_>>(),
vec!["a", "m", "z"]
);
}
#[test]
fn a_column_appearing_only_in_a_later_record_is_still_planned() {
let page = vec![json!({ "id": 1 }), json!({ "id": 2, "note": "hi" })];
let cols = plan_columns(&page).expect("a plan");
assert_eq!(
cols.iter().map(|c| c.name.as_str()).collect::<Vec<_>>(),
vec!["id", "note"]
);
}
#[test]
fn every_planned_column_is_nullable() {
let page = vec![json!({ "id": 1, "name": "a" })];
let cols = plan_columns(&page).expect("a plan");
assert!(cols.iter().all(|c| c.nullable), "{cols:?}");
}
#[test]
fn nested_values_plan_as_json() {
let page = vec![json!({ "obj": {"a": 1}, "arr": [1, 2] })];
let cols = plan_columns(&page).expect("a plan");
assert!(
cols.iter().all(|c| c.base_type == SqlBaseType::Json),
"{cols:?}"
);
}
#[test]
fn an_all_null_column_falls_back_to_text() {
let page = vec![json!({ "id": 1, "maybe": Value::Null })];
let cols = plan_columns(&page).expect("a plan");
let maybe = cols.iter().find(|c| c.name == "maybe").expect("column");
assert_eq!(maybe.base_type, SqlBaseType::Text);
}
#[test]
fn an_empty_or_non_object_page_plans_nothing() {
assert!(plan_columns(&[]).is_none());
assert!(plan_columns(&[json!(1), json!("x")]).is_none());
assert!(plan_columns(&[json!({})]).is_none());
}
#[test]
fn mixed_int_and_float_widens_to_double() {
let page = vec![json!({ "n": 1 }), json!({ "n": 1.5 })];
let cols = plan_columns(&page).expect("a plan");
assert_eq!(cols[0].base_type, SqlBaseType::Double);
}
#[test]
fn render_columns_uses_the_dialects_quoting_and_types() {
let cols = plan_columns(&[json!({ "id": 1, "name": "a" })]).expect("a plan");
let sql = render_columns(
&cols,
|n| format!("\"{}\"", n.replace('"', "\"\"")),
|t| match t {
SqlBaseType::Integer => "BIGINT",
SqlBaseType::Double => "DOUBLE PRECISION",
SqlBaseType::Boolean => "BOOLEAN",
SqlBaseType::Text => "TEXT",
SqlBaseType::Json => "JSONB",
},
);
assert_eq!(sql, "\"id\" BIGINT, \"name\" TEXT");
}
#[test]
fn the_missing_target_error_names_both_ways_out() {
let e = missing_target_error("postgres sink", "public.orders");
let msg = e.to_string();
assert!(msg.contains("public.orders"), "{msg}");
assert!(
msg.contains("create_table: true"),
"must name the fix: {msg}"
);
assert!(
matches!(e, crate::error::FaucetError::Sink(_)),
"a missing destination is a sink failure, not a config one — the config \
was legal, the destination was not there"
);
}
#[test]
fn render_columns_quotes_a_hostile_identifier() {
let cols = plan_columns(&[json!({ "we\"ird": 1 })]).expect("a plan");
let sql = render_columns(
&cols,
|n| format!("\"{}\"", n.replace('"', "\"\"")),
|_| "TEXT",
);
assert_eq!(sql, "\"we\"\"ird\" TEXT");
}
#[test]
fn keyed_plan_marks_key_columns_required_and_adds_missing_ones() {
let page = vec![json!({ "id": 1, "name": "a" })];
let cols = plan_keyed_columns(&page, &["id".into(), "tenant".into()]).expect("a plan");
let id = cols.iter().find(|c| c.name == "id").unwrap();
assert!(!id.nullable);
assert_eq!(id.base_type, SqlBaseType::Integer);
let tenant = cols.iter().find(|c| c.name == "tenant").unwrap();
assert!(!tenant.nullable);
assert_eq!(tenant.base_type, SqlBaseType::Text);
assert!(cols.iter().find(|c| c.name == "name").unwrap().nullable);
}
#[test]
fn keyed_plan_is_none_for_an_empty_page() {
assert!(plan_keyed_columns(&[], &["id".into()]).is_none());
}
#[test]
fn column_defs_can_type_by_column() {
let cols = vec![
PlannedColumn {
name: "id".into(),
base_type: SqlBaseType::Text,
nullable: false,
},
PlannedColumn {
name: "note".into(),
base_type: SqlBaseType::Text,
nullable: true,
},
];
let sql = render_column_defs(
&cols,
|n| format!("`{n}`"),
|c| {
if c.nullable { "TEXT" } else { "VARCHAR(255)" }
},
);
assert_eq!(sql, "`id` VARCHAR(255), `note` TEXT");
}
#[test]
fn primary_key_renders_every_key_column_in_order() {
let q = |n: &str| format!("\"{n}\"");
assert_eq!(
render_primary_key(&["a".into(), "b".into()], q).as_deref(),
Some("PRIMARY KEY (\"a\", \"b\")")
);
assert!(render_primary_key(&[], q).is_none());
}
}