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