1use serde::{Deserialize, Serialize};
13use serde_json::{Value, json};
14
15#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
17pub struct DatasetDescriptor {
18 pub name: String,
21 pub kind: String,
25 #[serde(default, skip_serializing_if = "Option::is_none")]
31 pub schema: Option<Value>,
32 #[serde(default, skip_serializing_if = "Option::is_none")]
35 pub estimated_rows: Option<u64>,
36 pub config_patch: Value,
42 #[serde(default, skip_serializing_if = "Option::is_none")]
48 pub sink_patch: Option<Value>,
49 #[serde(default, skip_serializing_if = "Option::is_none")]
53 pub primary_key: Option<Vec<String>>,
54}
55
56impl DatasetDescriptor {
57 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 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 pub fn with_schema(mut self, schema: Value) -> Self {
79 self.schema = Some(schema);
80 self
81 }
82
83 pub fn with_sink_patch(mut self, sink_patch: Value) -> Self {
85 self.sink_patch = Some(sink_patch);
86 self
87 }
88
89 pub fn with_estimated_rows(mut self, rows: u64) -> Self {
91 self.estimated_rows = Some(rows);
92 self
93 }
94}
95
96pub 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
117pub 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 "object"
133 } else if t.contains("array") || t.starts_with('_') {
134 "array"
137 } else if (t.contains("int") && !t.contains("interval") && !t.contains("point"))
138 || t == "serial"
139 || t == "bigserial"
140 || t == "smallserial"
141 {
142 "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 "string"
158 };
159 json!({ "type": ty })
160}
161
162pub 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 if let Some(obj) = fragment.as_object_mut() {
172 obj.insert("type".into(), json!([single, "null"]));
173 }
174 fragment
175}
176
177pub 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 let complex = json!({"type": ["string", "null"]});
294 assert_eq!(nullable_type(complex.clone()), complex);
295 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}