Skip to main content

faucet_core/
create_table.rs

1//! Shared auto-create-table planning for table-based sinks (#580).
2//!
3//! A first-ever sync cannot assume the destination table exists. The BigQuery
4//! sink solved this in #578 by inferring a schema from the first written page
5//! and creating the table; every other table destination has the same problem,
6//! and before this module each one answered it differently (or not at all) —
7//! different field names, different defaults, some with no create path at all.
8//!
9//! This module holds the **dialect-neutral half**: turn a page of records into
10//! an ordered, deduplicated column plan. Each sink supplies only the mapping
11//! from [`SqlBaseType`] to its own type keyword — the same mapping its
12//! `evolve_schema` already needs, so the two agree by construction rather than
13//! by inspection.
14//!
15//! ## Why column order is sorted, not source order
16//!
17//! Source order would read better in a `SELECT *`, but it is not recoverable:
18//! `serde_json::Map` is a `BTreeMap` or an `IndexMap` depending on whether the
19//! `preserve_order` feature is unified into the build, so "the order the
20//! record listed its keys" is a *build-time* property. A created table's
21//! column order would then differ between a `-p crate` build and an
22//! `--all-features` one — the kind of difference that only shows up as a
23//! diffed schema in production. Sorted is worse to read and always the same.
24
25use crate::drift::{SqlBaseType, json_schema_base_type};
26use serde_json::Value;
27
28/// One column of a planned `CREATE TABLE`.
29#[derive(Debug, Clone, PartialEq, Eq)]
30pub struct PlannedColumn {
31    /// Column name, exactly as it appeared in the record. Quoting is the
32    /// sink's job (each dialect quotes differently).
33    pub name: String,
34    /// Dialect-neutral type, to be mapped by the sink.
35    pub base_type: SqlBaseType,
36    /// Whether any observed value for this column was null, or the column was
37    /// absent from at least one record. Planned columns are always created
38    /// nullable in practice — see [`plan_columns`] — but the flag is kept so a
39    /// sink that wants to emit `NOT NULL` for a provably-present column can.
40    pub nullable: bool,
41}
42
43/// Plan the columns for a table created from `page`.
44///
45/// Returns `None` when the page carries nothing to infer from (empty, or no
46/// record is a JSON object) — the caller must then leave the table
47/// uncreated rather than emit a zero-column `CREATE TABLE`, and try again on
48/// the next page.
49///
50/// **Every column is planned nullable.** A column that happened to be present
51/// and non-null in the first page is not thereby required forever, and a
52/// `NOT NULL` inferred from one page turns the *second* page into a hard write
53/// failure the moment a real-world record omits the field. Narrowing later is
54/// the schema-drift policy's job (#194), which can see more than one page.
55pub fn plan_columns(page: &[Value]) -> Option<Vec<PlannedColumn>> {
56    let schema = crate::schema::infer_schema(page);
57    let props = schema.get("properties")?.as_object()?;
58    if props.is_empty() {
59        return None;
60    }
61
62    // Sorted, for the determinism reason in the module docs. `infer_schema`
63    // already unions every record's keys, so a column that appears only in a
64    // later record is still planned.
65    let mut ordered: Vec<String> = props.keys().cloned().collect();
66    ordered.sort_unstable();
67
68    let columns: Vec<PlannedColumn> = ordered
69        .into_iter()
70        .map(|name| {
71            let fragment = &props[&name];
72            PlannedColumn {
73                // A column whose every observed value was null carries no type
74                // information; TEXT is the one choice that can hold whatever
75                // shows up next without a lossy cast.
76                base_type: json_schema_base_type(fragment).unwrap_or(SqlBaseType::Text),
77                nullable: true,
78                name,
79            }
80        })
81        .collect();
82
83    (!columns.is_empty()).then_some(columns)
84}
85
86/// Render a `CREATE TABLE` column list, given a dialect's type mapper and
87/// identifier quoter.
88///
89/// Kept separate from [`plan_columns`] so a sink whose `CREATE TABLE` needs
90/// extra clauses (a primary key, a partition spec, an engine) composes this
91/// into its own statement rather than fighting a one-size template.
92pub fn render_columns<Q, T>(columns: &[PlannedColumn], quote: Q, ty: T) -> String
93where
94    Q: Fn(&str) -> String,
95    T: Fn(SqlBaseType) -> &'static str,
96{
97    columns
98        .iter()
99        .map(|c| format!("{} {}", quote(&c.name), ty(c.base_type)))
100        .collect::<Vec<_>>()
101        .join(", ")
102}
103
104/// Plan the columns for a table whose writes dedup on `key` (#676).
105///
106/// Like [`plan_columns`], plus: every key column is planned even when the
107/// first page never carries it (typed TEXT, the type that holds anything), and
108/// key columns are marked non-nullable. A created table can then carry a
109/// primary key on `key` — which is what gives an upsert's `ON CONFLICT` /
110/// `MERGE` a target on the very first run.
111pub fn plan_keyed_columns(page: &[Value], key: &[String]) -> Option<Vec<PlannedColumn>> {
112    let mut columns = plan_columns(page)?;
113    for k in key {
114        match columns.iter_mut().find(|c| &c.name == k) {
115            Some(c) => c.nullable = false,
116            None => columns.push(PlannedColumn {
117                name: k.clone(),
118                base_type: SqlBaseType::Text,
119                nullable: false,
120            }),
121        }
122    }
123    Some(columns)
124}
125
126/// Render a column list where the type keyword may depend on the column, not
127/// just its base type — for dialects that cannot index an unbounded type and so
128/// need a bounded one for key columns (MySQL `TEXT`, SQL Server `NVARCHAR(MAX)`).
129pub fn render_column_defs<Q, T>(columns: &[PlannedColumn], quote: Q, ty: T) -> String
130where
131    Q: Fn(&str) -> String,
132    T: Fn(&PlannedColumn) -> &'static str,
133{
134    columns
135        .iter()
136        .map(|c| format!("{} {}", quote(&c.name), ty(c)))
137        .collect::<Vec<_>>()
138        .join(", ")
139}
140
141/// `PRIMARY KEY (<key…>)` for a created table, or `None` for an empty key.
142pub fn render_primary_key<Q>(key: &[String], quote: Q) -> Option<String>
143where
144    Q: Fn(&str) -> String,
145{
146    (!key.is_empty()).then(|| {
147        format!(
148            "PRIMARY KEY ({})",
149            key.iter().map(|k| quote(k)).collect::<Vec<_>>().join(", ")
150        )
151    })
152}
153
154/// The uniform error a sink raises when its target is missing and
155/// `create_table: false` (#580).
156///
157/// One wording across every sink, because the fix is always the same two
158/// choices and an operator hitting it on connector number three should not
159/// have to re-read a different sentence.
160pub fn missing_target_error(connector: &str, target: &str) -> crate::error::FaucetError {
161    crate::error::FaucetError::Sink(format!(
162        "{connector}: target `{target}` does not exist and `create_table: false`. \
163         Create it first, or set `create_table: true` to have faucet create it from \
164         the first page's inferred schema."
165    ))
166}
167
168#[cfg(test)]
169mod tests {
170    use super::*;
171    use serde_json::json;
172
173    #[test]
174    fn plans_every_column_with_its_inferred_type() {
175        let page = vec![
176            json!({ "id": 1, "name": "a", "amount": 1.5, "ok": true }),
177            json!({ "id": 2, "name": "b", "amount": 2.5, "ok": false }),
178        ];
179        let cols = plan_columns(&page).expect("a plan");
180        let by_name: std::collections::HashMap<&str, SqlBaseType> = cols
181            .iter()
182            .map(|c| (c.name.as_str(), c.base_type))
183            .collect();
184        assert_eq!(by_name["id"], SqlBaseType::Integer);
185        assert_eq!(by_name["name"], SqlBaseType::Text);
186        assert_eq!(by_name["amount"], SqlBaseType::Double);
187        assert_eq!(by_name["ok"], SqlBaseType::Boolean);
188    }
189
190    #[test]
191    fn column_order_is_deterministic_regardless_of_record_key_order() {
192        // `serde_json::Map` is a BTreeMap or an IndexMap depending on feature
193        // unification, so a plan that read order off the record would produce
194        // a different table under `-p crate` than under `--all-features`.
195        let a = plan_columns(&[json!({ "z": 1, "a": 2, "m": 3 })]).expect("a");
196        let b = plan_columns(&[json!({ "a": 2, "m": 3, "z": 1 })]).expect("b");
197        assert_eq!(a, b, "the same record set must plan the same table");
198        assert_eq!(
199            a.iter().map(|c| c.name.as_str()).collect::<Vec<_>>(),
200            vec!["a", "m", "z"]
201        );
202    }
203
204    #[test]
205    fn a_column_appearing_only_in_a_later_record_is_still_planned() {
206        // A sparse first record must not silently drop a column — the next
207        // page would then fail to write it, on a table faucet itself created.
208        let page = vec![json!({ "id": 1 }), json!({ "id": 2, "note": "hi" })];
209        let cols = plan_columns(&page).expect("a plan");
210        assert_eq!(
211            cols.iter().map(|c| c.name.as_str()).collect::<Vec<_>>(),
212            vec!["id", "note"]
213        );
214    }
215
216    #[test]
217    fn every_planned_column_is_nullable() {
218        // Inferring NOT NULL from one page turns the second page into a hard
219        // write failure the first time a record omits the field.
220        let page = vec![json!({ "id": 1, "name": "a" })];
221        let cols = plan_columns(&page).expect("a plan");
222        assert!(cols.iter().all(|c| c.nullable), "{cols:?}");
223    }
224
225    #[test]
226    fn nested_values_plan_as_json() {
227        let page = vec![json!({ "obj": {"a": 1}, "arr": [1, 2] })];
228        let cols = plan_columns(&page).expect("a plan");
229        assert!(
230            cols.iter().all(|c| c.base_type == SqlBaseType::Json),
231            "{cols:?}"
232        );
233    }
234
235    #[test]
236    fn an_all_null_column_falls_back_to_text() {
237        // No type information at all; TEXT is the only choice that can hold
238        // whatever turns up next without a lossy cast.
239        let page = vec![json!({ "id": 1, "maybe": Value::Null })];
240        let cols = plan_columns(&page).expect("a plan");
241        let maybe = cols.iter().find(|c| c.name == "maybe").expect("column");
242        assert_eq!(maybe.base_type, SqlBaseType::Text);
243    }
244
245    #[test]
246    fn an_empty_or_non_object_page_plans_nothing() {
247        // A zero-column CREATE TABLE is invalid in every dialect, so the sink
248        // must wait for a page it can actually learn from.
249        assert!(plan_columns(&[]).is_none());
250        assert!(plan_columns(&[json!(1), json!("x")]).is_none());
251        assert!(plan_columns(&[json!({})]).is_none());
252    }
253
254    #[test]
255    fn mixed_int_and_float_widens_to_double() {
256        // Creating an INTEGER column and then writing 1.5 into it is exactly
257        // the silent-truncation class this inference exists to avoid.
258        let page = vec![json!({ "n": 1 }), json!({ "n": 1.5 })];
259        let cols = plan_columns(&page).expect("a plan");
260        assert_eq!(cols[0].base_type, SqlBaseType::Double);
261    }
262
263    #[test]
264    fn render_columns_uses_the_dialects_quoting_and_types() {
265        let cols = plan_columns(&[json!({ "id": 1, "name": "a" })]).expect("a plan");
266        let sql = render_columns(
267            &cols,
268            |n| format!("\"{}\"", n.replace('"', "\"\"")),
269            |t| match t {
270                SqlBaseType::Integer => "BIGINT",
271                SqlBaseType::Double => "DOUBLE PRECISION",
272                SqlBaseType::Boolean => "BOOLEAN",
273                SqlBaseType::Text => "TEXT",
274                SqlBaseType::Json => "JSONB",
275            },
276        );
277        assert_eq!(sql, "\"id\" BIGINT, \"name\" TEXT");
278    }
279
280    #[test]
281    fn the_missing_target_error_names_both_ways_out() {
282        let e = missing_target_error("postgres sink", "public.orders");
283        let msg = e.to_string();
284        assert!(msg.contains("public.orders"), "{msg}");
285        assert!(
286            msg.contains("create_table: true"),
287            "must name the fix: {msg}"
288        );
289        assert!(
290            matches!(e, crate::error::FaucetError::Sink(_)),
291            "a missing destination is a sink failure, not a config one — the config \
292             was legal, the destination was not there"
293        );
294    }
295
296    #[test]
297    fn render_columns_quotes_a_hostile_identifier() {
298        // The quoter is the sink's, but the renderer must actually route every
299        // name through it — a name reaching the SQL unquoted is an injection.
300        let cols = plan_columns(&[json!({ "we\"ird": 1 })]).expect("a plan");
301        let sql = render_columns(
302            &cols,
303            |n| format!("\"{}\"", n.replace('"', "\"\"")),
304            |_| "TEXT",
305        );
306        assert_eq!(sql, "\"we\"\"ird\" TEXT");
307    }
308
309    #[test]
310    fn keyed_plan_marks_key_columns_required_and_adds_missing_ones() {
311        let page = vec![json!({ "id": 1, "name": "a" })];
312        let cols = plan_keyed_columns(&page, &["id".into(), "tenant".into()]).expect("a plan");
313        let id = cols.iter().find(|c| c.name == "id").unwrap();
314        assert!(!id.nullable);
315        assert_eq!(id.base_type, SqlBaseType::Integer);
316        let tenant = cols.iter().find(|c| c.name == "tenant").unwrap();
317        assert!(!tenant.nullable);
318        assert_eq!(tenant.base_type, SqlBaseType::Text);
319        assert!(cols.iter().find(|c| c.name == "name").unwrap().nullable);
320    }
321
322    #[test]
323    fn keyed_plan_is_none_for_an_empty_page() {
324        assert!(plan_keyed_columns(&[], &["id".into()]).is_none());
325    }
326
327    #[test]
328    fn column_defs_can_type_by_column() {
329        let cols = vec![
330            PlannedColumn {
331                name: "id".into(),
332                base_type: SqlBaseType::Text,
333                nullable: false,
334            },
335            PlannedColumn {
336                name: "note".into(),
337                base_type: SqlBaseType::Text,
338                nullable: true,
339            },
340        ];
341        let sql = render_column_defs(
342            &cols,
343            |n| format!("`{n}`"),
344            |c| {
345                if c.nullable { "TEXT" } else { "VARCHAR(255)" }
346            },
347        );
348        assert_eq!(sql, "`id` VARCHAR(255), `note` TEXT");
349    }
350
351    #[test]
352    fn primary_key_renders_every_key_column_in_order() {
353        let q = |n: &str| format!("\"{n}\"");
354        assert_eq!(
355            render_primary_key(&["a".into(), "b".into()], q).as_deref(),
356            Some("PRIMARY KEY (\"a\", \"b\")")
357        );
358        assert!(render_primary_key(&[], q).is_none());
359    }
360}