Skip to main content

faucet_core/
discover.rs

1//! Live source introspection (`faucet discover`, #211).
2//!
3//! A [`Source`](crate::Source) that can enumerate the datasets living behind
4//! its connection (tables in a database schema, collections in MongoDB,
5//! indices in Elasticsearch, key prefixes in an object store) overrides
6//! [`Source::discover`](crate::Source::discover) to return one
7//! [`DatasetDescriptor`] per dataset. The CLI turns that list into a
8//! ready-to-run config: one matrix row per dataset, each row deep-merging the
9//! descriptor's [`config_patch`](DatasetDescriptor::config_patch) over the
10//! connection config.
11
12use serde::{Deserialize, Serialize};
13use serde_json::{Value, json};
14
15/// One dataset discovered behind a source's connection.
16#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
17pub struct DatasetDescriptor {
18    /// Human-readable dataset identity — a qualified table name
19    /// (`public.orders`), a collection, an index, or an object-key prefix.
20    pub name: String,
21    /// What kind of dataset this is: `table`, `collection`, `index`,
22    /// `prefix`, or `object`. Informational (rendered in output); the
23    /// selection mechanics live entirely in [`config_patch`](Self::config_patch).
24    pub kind: String,
25    /// The dataset's column shape as an
26    /// [`infer_schema`](crate::schema::infer_schema)-shaped object
27    /// (`{"type":"object","properties":{…}}`), or `None` when the source
28    /// cannot cheaply introspect it (object stores, schemaless collections
29    /// with no sample).
30    #[serde(default, skip_serializing_if = "Option::is_none")]
31    pub schema: Option<Value>,
32    /// Approximate row/object count when the catalog exposes one cheaply
33    /// (e.g. `pg_class.reltuples`). Never computed via a full scan.
34    #[serde(default, skip_serializing_if = "Option::is_none")]
35    pub estimated_rows: Option<u64>,
36    /// Partial source-config **override** selecting exactly this dataset —
37    /// deep-merged over the connection config by a matrix row (e.g.
38    /// `{"query": "SELECT * FROM \"public\".\"orders\""}` for a SQL source,
39    /// `{"collection": "orders"}` for MongoDB, `{"prefix": "raw/orders/"}`
40    /// for an object store). Must never contain credentials.
41    pub config_patch: Value,
42    /// Optional partial **sink**-config override routing this dataset to its own
43    /// destination — deep-merged over the sink template by a matrix row (e.g.
44    /// `{"table_id": "account"}` so a fan-out of discovered objects lands one
45    /// table per object). `None` (the common case) leaves the sink untouched.
46    /// Must never contain credentials.
47    #[serde(default, skip_serializing_if = "Option::is_none")]
48    pub sink_patch: Option<Value>,
49    /// The dataset's primary-key columns in key order, when the catalog
50    /// declares one (`None` when it has none or the source cannot tell). A
51    /// multi-table `faucet mirror` keys each table's upsert on it (#731).
52    #[serde(default, skip_serializing_if = "Option::is_none")]
53    pub primary_key: Option<Vec<String>>,
54}
55
56impl DatasetDescriptor {
57    /// Construct a descriptor with no schema / row estimate.
58    pub fn new(name: impl Into<String>, kind: impl Into<String>, config_patch: Value) -> Self {
59        Self {
60            name: name.into(),
61            kind: kind.into(),
62            schema: None,
63            estimated_rows: None,
64            config_patch,
65            sink_patch: None,
66            primary_key: None,
67        }
68    }
69
70    /// Attach the catalog's primary-key columns (in key order). An empty list
71    /// means "no primary key" and is stored as `None`.
72    pub fn with_primary_key(mut self, columns: Vec<String>) -> Self {
73        self.primary_key = (!columns.is_empty()).then_some(columns);
74        self
75    }
76
77    /// Attach an inferred/introspected schema.
78    pub fn with_schema(mut self, schema: Value) -> Self {
79        self.schema = Some(schema);
80        self
81    }
82
83    /// Attach a sink-config override routing this dataset to its own destination.
84    pub fn with_sink_patch(mut self, sink_patch: Value) -> Self {
85        self.sink_patch = Some(sink_patch);
86        self
87    }
88
89    /// Attach a catalog-provided row estimate.
90    pub fn with_estimated_rows(mut self, rows: u64) -> Self {
91        self.estimated_rows = Some(rows);
92        self
93    }
94}
95
96/// Attach primary keys to catalog descriptors (#731). `key_columns` yields
97/// `(dataset name, column)` pairs **in key order** (as a catalog query ordered
98/// by the constraint's column position returns them); a dataset with no pair
99/// keeps `primary_key: None`.
100pub fn attach_primary_keys(
101    descriptors: Vec<DatasetDescriptor>,
102    key_columns: impl IntoIterator<Item = (String, String)>,
103) -> Vec<DatasetDescriptor> {
104    let mut keys: std::collections::HashMap<String, Vec<String>> = Default::default();
105    for (name, column) in key_columns {
106        keys.entry(name).or_default().push(column);
107    }
108    descriptors
109        .into_iter()
110        .map(|d| match keys.remove(&d.name) {
111            Some(cols) => d.with_primary_key(cols),
112            None => d,
113        })
114        .collect()
115}
116
117/// Map a SQL catalog type name (as reported by `information_schema.columns`
118/// or an equivalent) to a JSON-Schema type fragment matching the shape
119/// [`infer_schema`](crate::schema::infer_schema) produces.
120///
121/// The mapping is deliberately conservative and dialect-tolerant — it matches
122/// on lowercase substrings so `TIMESTAMP WITH TIME ZONE`, `timestamptz`,
123/// `Nullable(Int64)` etc. all land on a sensible JSON type. Unknown types map
124/// to `string` (every SQL value has a textual form, so `string` is the safe
125/// over-approximation — never a lossy one like `number`).
126pub fn sql_type_to_json_schema(data_type: &str) -> Value {
127    let t = data_type.to_ascii_lowercase();
128    let ty = if t.contains("bool") || t == "bit" {
129        "boolean"
130    } else if t.contains("json") || t.contains("variant") || t == "object" || t.contains("struct") {
131        // checked before the integer branch so `STRUCT<a INT>` maps to object
132        "object"
133    } else if t.contains("array") || t.starts_with('_') {
134        // Postgres reports array types as `ARRAY` (information_schema) or
135        // `_int4`-style internal names.
136        "array"
137    } else if (t.contains("int") && !t.contains("interval") && !t.contains("point"))
138        || t == "serial"
139        || t == "bigserial"
140        || t == "smallserial"
141    {
142        // covers int2/4/8, integer, bigint, smallint, tinyint, mediumint —
143        // but not `interval` / `point`, whose substring would false-match
144        "integer"
145    } else if t.contains("double")
146        || t.contains("float")
147        || t.contains("real")
148        || t.contains("decimal")
149        || t.contains("numeric")
150        || t.contains("money")
151        || t == "number"
152    {
153        "number"
154    } else {
155        // text, varchar, char, uuid, date, time, timestamp, bytea, blob,
156        // enum, inet, … — all serialized as JSON strings by the sources.
157        "string"
158    };
159    json!({ "type": ty })
160}
161
162/// Wrap a type fragment as nullable (`{"type": ["T", "null"]}`), matching the
163/// nullable shape [`infer_schema`](crate::schema::infer_schema) emits.
164pub fn nullable_type(mut fragment: Value) -> Value {
165    let single = match fragment.get("type") {
166        Some(Value::String(t)) => t.clone(),
167        _ => return fragment,
168    };
169    // Wrap only the `type` in place so sibling keys (e.g. a `format: date-time`
170    // hint a typed sink maps to TIMESTAMP/DATE) survive the nullability wrap.
171    if let Some(obj) = fragment.as_object_mut() {
172        obj.insert("type".into(), json!([single, "null"]));
173    }
174    fragment
175}
176
177/// Assemble an `infer_schema`-shaped object schema from
178/// `(column_name, type_fragment)` pairs, preserving input order semantics
179/// (`serde_json` map ordering applies).
180pub fn columns_to_schema(columns: impl IntoIterator<Item = (String, Value)>) -> Value {
181    let mut properties = serde_json::Map::new();
182    for (name, fragment) in columns {
183        properties.insert(name, fragment);
184    }
185    json!({ "type": "object", "properties": Value::Object(properties) })
186}
187
188#[cfg(test)]
189mod tests {
190    use super::*;
191
192    #[test]
193    fn descriptor_builder_round_trips() {
194        let d = DatasetDescriptor::new("public.orders", "table", json!({"query": "SELECT 1"}))
195            .with_schema(json!({"type": "object", "properties": {}}))
196            .with_estimated_rows(42);
197        assert_eq!(d.name, "public.orders");
198        assert_eq!(d.kind, "table");
199        assert_eq!(d.estimated_rows, Some(42));
200        let v = serde_json::to_value(&d).unwrap();
201        let back: DatasetDescriptor = serde_json::from_value(v).unwrap();
202        assert_eq!(back, d);
203    }
204
205    #[test]
206    fn primary_key_builder_keeps_order_and_treats_empty_as_none() {
207        let d = DatasetDescriptor::new("public.orders", "table", json!({}))
208            .with_primary_key(vec!["tenant".into(), "id".into()]);
209        assert_eq!(d.primary_key, Some(vec!["tenant".into(), "id".into()]));
210        let v = serde_json::to_value(&d).unwrap();
211        assert_eq!(v["primary_key"], json!(["tenant", "id"]));
212        let none = DatasetDescriptor::new("t", "table", json!({})).with_primary_key(vec![]);
213        assert_eq!(none.primary_key, None);
214        assert!(
215            serde_json::to_value(&none)
216                .unwrap()
217                .get("primary_key")
218                .is_none()
219        );
220    }
221
222    #[test]
223    fn attach_primary_keys_matches_by_name_in_key_order() {
224        let ds = vec![
225            DatasetDescriptor::new("public.a", "table", json!({})),
226            DatasetDescriptor::new("public.b", "table", json!({})),
227        ];
228        let out = attach_primary_keys(
229            ds,
230            [
231                ("public.a".to_string(), "tenant".to_string()),
232                ("public.a".to_string(), "id".to_string()),
233                ("public.zzz".to_string(), "id".to_string()),
234            ],
235        );
236        assert_eq!(out[0].primary_key, Some(vec!["tenant".into(), "id".into()]));
237        assert_eq!(out[1].primary_key, None);
238    }
239
240    #[test]
241    fn descriptor_omits_absent_optionals_in_json() {
242        let d = DatasetDescriptor::new("t", "table", json!({}));
243        let v = serde_json::to_value(&d).unwrap();
244        assert!(v.get("schema").is_none());
245        assert!(v.get("estimated_rows").is_none());
246    }
247
248    #[test]
249    fn sql_types_map_to_json_types() {
250        for (sql, want) in [
251            ("integer", "integer"),
252            ("BIGINT", "integer"),
253            ("smallint", "integer"),
254            ("tinyint(1)", "integer"),
255            ("serial", "integer"),
256            ("double precision", "number"),
257            ("NUMERIC(10,2)", "number"),
258            ("decimal", "number"),
259            ("float8", "number"),
260            ("money", "number"),
261            ("boolean", "boolean"),
262            ("BOOL", "boolean"),
263            ("bit", "boolean"),
264            ("json", "object"),
265            ("JSONB", "object"),
266            ("VARIANT", "object"),
267            ("STRUCT<a INT>", "object"),
268            ("ARRAY", "array"),
269            ("_int4", "array"),
270            ("text", "string"),
271            ("character varying", "string"),
272            ("uuid", "string"),
273            ("timestamp with time zone", "string"),
274            ("date", "string"),
275            ("bytea", "string"),
276            ("some_exotic_type", "string"),
277        ] {
278            assert_eq!(
279                sql_type_to_json_schema(sql),
280                json!({ "type": want }),
281                "for SQL type {sql:?}"
282            );
283        }
284    }
285
286    #[test]
287    fn nullable_wraps_scalar_type() {
288        assert_eq!(
289            nullable_type(json!({"type": "integer"})),
290            json!({"type": ["integer", "null"]})
291        );
292        // Already-complex fragments pass through untouched.
293        let complex = json!({"type": ["string", "null"]});
294        assert_eq!(nullable_type(complex.clone()), complex);
295        // Sibling keys survive the wrap — a temporal `format` hint must not be
296        // lost when the column is nullable (a typed sink maps it to TIMESTAMP).
297        assert_eq!(
298            nullable_type(json!({"type": "string", "format": "date-time"})),
299            json!({"type": ["string", "null"], "format": "date-time"})
300        );
301    }
302
303    #[test]
304    fn columns_to_schema_shapes_like_infer_schema() {
305        let schema = columns_to_schema(vec![
306            ("id".to_string(), json!({"type": "integer"})),
307            ("name".to_string(), json!({"type": ["string", "null"]})),
308        ]);
309        assert_eq!(schema["type"], "object");
310        assert_eq!(schema["properties"]["id"]["type"], "integer");
311        assert_eq!(schema["properties"]["name"]["type"][1], "null");
312    }
313}