use std::collections::BTreeMap;
use crate::error::{Error, Result};
pub const QUERY_SERVER_VERSION: &str = "SELECT version() AS version FORMAT JSONEachRow";
pub const QUERY_CHANGED_SETTINGS: &str =
"SELECT name, value FROM system.settings WHERE changed FORMAT JSONEachRow";
pub const QUERY_TABLE_COLUMNS: &str = "SELECT name, type, default_kind, default_expression, position \
FROM system.columns WHERE database = {db:String} AND table = {table:String} \
ORDER BY position FORMAT JSONEachRow";
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ServerProfile {
pub version: String,
pub settings: Vec<(String, String)>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct DiscoveredColumn {
pub name: String,
pub r#type: String,
pub default_kind: String,
pub default_expression: String,
pub position: u64,
}
fn discovery_err<T>(message: impl Into<String>) -> Result<T> {
Err(Error::Discovery {
message: message.into(),
})
}
fn excerpt(line: &[u8]) -> String {
let head = &line[..line.len().min(60)];
String::from_utf8_lossy(head).into_owned()
}
fn json_each_row_docs(body: &[u8]) -> Result<Vec<serde_json::Map<String, serde_json::Value>>> {
let mut out = Vec::new();
for line in body.split(|&b| b == b'\n') {
let line = trim_ascii(line);
if line.is_empty() {
continue;
}
if line[0] != b'{' {
return discovery_err(format!("not a JSONEachRow line: {}", excerpt(line)));
}
match serde_json::from_slice::<serde_json::Value>(line) {
Ok(serde_json::Value::Object(map)) => out.push(map),
Ok(_) => return discovery_err(format!("not a JSON object: {}", excerpt(line))),
Err(e) => {
return discovery_err(format!("bad JSONEachRow line: {e}: {}", excerpt(line)));
}
}
}
Ok(out)
}
fn trim_ascii(mut b: &[u8]) -> &[u8] {
while let Some((f, rest)) = b.split_first() {
if f.is_ascii_whitespace() {
b = rest;
} else {
break;
}
}
while let Some((l, rest)) = b.split_last() {
if l.is_ascii_whitespace() {
b = rest;
} else {
break;
}
}
b
}
fn json_text(v: &serde_json::Value) -> String {
match v {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
}
}
fn text_field(row: &serde_json::Map<String, serde_json::Value>, key: &str) -> String {
row.get(key).map(json_text).unwrap_or_default()
}
pub fn parse_version_result(body: &[u8]) -> Result<String> {
let docs = json_each_row_docs(body)?;
if docs.len() != 1 {
return discovery_err(format!(
"version query returned {} rows, want 1",
docs.len()
));
}
let Some(v) = docs[0].get("version") else {
return discovery_err("version query row has no `version` field");
};
let out = json_text(v);
if out.is_empty() {
return discovery_err("version query returned an empty version");
}
Ok(out)
}
pub fn parse_changed_settings_result(body: &[u8]) -> Result<Vec<(String, String)>> {
let docs = json_each_row_docs(body)?;
let mut map = BTreeMap::new();
for row in &docs {
let Some(name) = row.get("name") else {
return discovery_err("settings row has no `name` field");
};
let Some(value) = row.get("value") else {
return discovery_err("settings row has no `value` field");
};
map.insert(json_text(name), json_text(value));
}
Ok(map.into_iter().collect())
}
pub fn parse_columns_result(body: &[u8]) -> Result<Vec<DiscoveredColumn>> {
let docs = json_each_row_docs(body)?;
let mut out = Vec::with_capacity(docs.len());
for row in &docs {
let col = DiscoveredColumn {
name: text_field(row, "name"),
r#type: text_field(row, "type"),
default_kind: text_field(row, "default_kind"),
default_expression: text_field(row, "default_expression"),
position: row
.get("position")
.and_then(|v| json_text(v).parse::<u64>().ok())
.unwrap_or(0),
};
if col.name.is_empty() || col.r#type.is_empty() {
return discovery_err(format!(
"columns row missing name/type: {}",
excerpt(
&serde_json::Value::Object(row.clone())
.to_string()
.into_bytes()
)
));
}
out.push(col);
}
if out.is_empty() {
return discovery_err(
"columns query returned no rows — wrong database/table, or no access",
);
}
Ok(out)
}
fn backquote_if_needed(name: &str) -> String {
let bytes = name.as_bytes();
let plain = !bytes.is_empty()
&& !bytes[0].is_ascii_digit()
&& bytes
.iter()
.all(|&c| c == b'_' || c.is_ascii_alphanumeric());
if plain {
return name.to_string();
}
format!("`{}`", name.replace('`', "``"))
}
pub fn reconstruct_ddl(cols: &[DiscoveredColumn]) -> Result<String> {
if cols.is_empty() {
return discovery_err("no columns to reconstruct");
}
let mut out = String::new();
for (i, c) in cols.iter().enumerate() {
if c.name.is_empty() || c.r#type.is_empty() {
return discovery_err(format!("column {i} has no name/type"));
}
if i > 0 {
out.push_str(", ");
}
out.push_str(&backquote_if_needed(&c.name));
out.push(' ');
out.push_str(&c.r#type);
match c.default_kind.as_str() {
"" => {
if !c.default_expression.is_empty() {
return discovery_err(format!(
"column {} has a default_expression but no default_kind",
c.name
));
}
}
kind @ ("DEFAULT" | "MATERIALIZED" | "ALIAS") => {
if c.default_expression.is_empty() {
return discovery_err(format!(
"column {} is {kind} but has no default_expression",
c.name
));
}
out.push(' ');
out.push_str(kind);
out.push(' ');
out.push_str(&c.default_expression);
}
"EPHEMERAL" => {
out.push_str(" EPHEMERAL");
if !c.default_expression.is_empty() {
out.push(' ');
out.push_str(&c.default_expression);
}
}
other => {
return discovery_err(format!(
"column {} has unknown default_kind {other:?}",
c.name
));
}
}
}
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_query_texts_are_the_specs_verbatim() {
assert_eq!(
QUERY_SERVER_VERSION,
"SELECT version() AS version FORMAT JSONEachRow"
);
assert_eq!(
QUERY_CHANGED_SETTINGS,
"SELECT name, value FROM system.settings WHERE changed FORMAT JSONEachRow"
);
assert_eq!(
QUERY_TABLE_COLUMNS,
"SELECT name, type, default_kind, default_expression, position \
FROM system.columns WHERE database = {db:String} AND table = {table:String} \
ORDER BY position FORMAT JSONEachRow"
);
}
#[test]
fn parse_version_result_reads_one_row_and_rejects_the_rest() {
let v = parse_version_result(b"{\"version\":\"25.8.28.1\"}\n").unwrap();
assert_eq!(v, "25.8.28.1");
assert!(parse_version_result(b"").is_err(), "empty body must error");
assert!(
parse_version_result(b"{\"version\":\"a\"}\n{\"version\":\"b\"}").is_err(),
"two rows must error"
);
assert!(
parse_version_result(b"{\"nope\":\"x\"}").is_err(),
"missing field must error"
);
}
#[test]
fn parse_changed_settings_result_reads_the_profile() {
let body = b"{\"name\":\"flatten_nested\",\"value\":\"0\"}\n\
{\"name\":\"date_time_input_format\",\"value\":\"best_effort\"}\n\
{\"name\":\"max_block_size\",\"value\":\"65409\"}\n";
let m = parse_changed_settings_result(body).unwrap();
assert_eq!(m.len(), 3);
let get = |k: &str| m.iter().find(|(n, _)| n == k).map(|(_, v)| v.as_str());
assert_eq!(get("flatten_nested"), Some("0"));
assert_eq!(get("date_time_input_format"), Some("best_effort"));
assert_eq!(get("max_block_size"), Some("65409"));
assert_eq!(parse_changed_settings_result(b"\n").unwrap(), vec![]);
}
#[test]
fn parse_changed_settings_duplicates_are_last_write_wins() {
let body = b"{\"name\":\"flatten_nested\",\"value\":\"1\"}\n\
{\"name\":\"max_block_size\",\"value\":\"65409\"}\n\
{\"name\":\"flatten_nested\",\"value\":\"0\"}\n";
let m = parse_changed_settings_result(body).unwrap();
assert_eq!(
m,
vec![
("flatten_nested".to_string(), "0".to_string()),
("max_block_size".to_string(), "65409".to_string()),
]
);
}
#[test]
fn parse_columns_result_survives_quoted_and_bare_positions() {
let body = br#"{"name":"id","type":"UInt64","default_kind":"","default_expression":"","position":"1"}
{"name":"ts","type":"DateTime","default_kind":"DEFAULT","default_expression":"now()","position":2}
{"name":"n.a","type":"Array(Int64)","default_kind":"","default_expression":"","position":"3"}
"#;
let cols = parse_columns_result(body).unwrap();
assert_eq!(cols.len(), 3);
assert_eq!(cols[0].position, 1);
assert_eq!(cols[1].position, 2);
assert_eq!(cols[2].position, 3);
assert_eq!(cols[1].default_kind, "DEFAULT");
assert_eq!(cols[1].default_expression, "now()");
assert!(
parse_columns_result(b"").is_err(),
"no rows must error — an empty table description is a wrong database/table"
);
}
#[test]
fn a_19_digit_position_never_routes_through_a_float() {
for body in [
&br#"{"name":"x","type":"UInt8","position":"18446744073709551615"}"#[..],
&br#"{"name":"x","type":"UInt8","position":18446744073709551615}"#[..],
] {
let cols = parse_columns_result(body).unwrap();
assert_eq!(cols[0].position, 18446744073709551615, "body {body:?}");
}
}
#[test]
fn reconstruct_ddl_spells_kinds_and_backquotes_only_where_needed() {
let cols = [
DiscoveredColumn {
name: "id".into(),
r#type: "UInt64".into(),
..Default::default()
},
DiscoveredColumn {
name: "ts".into(),
r#type: "DateTime".into(),
default_kind: "DEFAULT".into(),
default_expression: "now()".into(),
..Default::default()
},
DiscoveredColumn {
name: "n.a".into(),
r#type: "Array(Int64)".into(),
..Default::default()
},
DiscoveredColumn {
name: "e".into(),
r#type: "UInt8".into(),
default_kind: "EPHEMERAL".into(),
..Default::default()
},
DiscoveredColumn {
name: "m".into(),
r#type: "UInt64".into(),
default_kind: "MATERIALIZED".into(),
default_expression: "id + 1".into(),
..Default::default()
},
];
let ddl = reconstruct_ddl(&cols).unwrap();
assert!(
ddl.contains("`n.a` Array(Int64)"),
"flattened-Nested name not backquoted: {ddl}"
);
assert!(ddl.contains("ts DateTime DEFAULT now()"), "{ddl}");
assert!(ddl.contains("e UInt8 EPHEMERAL"), "{ddl}");
assert!(ddl.contains("m UInt64 MATERIALIZED id + 1"), "{ddl}");
assert_eq!(backquote_if_needed("a`b"), "`a``b`");
assert_eq!(backquote_if_needed("1x"), "`1x`");
assert_eq!(backquote_if_needed("with space"), "`with space`");
assert_eq!(backquote_if_needed("plain_Name9"), "plain_Name9");
}
#[test]
fn reconstruct_ddl_error_surfaces() {
for kind in ["DEFAULT", "MATERIALIZED", "ALIAS"] {
assert!(
reconstruct_ddl(&[DiscoveredColumn {
name: "x".into(),
r#type: "UInt8".into(),
default_kind: kind.into(),
..Default::default()
}])
.is_err(),
"{kind} without expression must error"
);
}
assert!(
reconstruct_ddl(&[DiscoveredColumn {
name: "x".into(),
r#type: "UInt8".into(),
default_kind: "WEIRD".into(),
..Default::default()
}])
.is_err()
);
assert!(
reconstruct_ddl(&[DiscoveredColumn {
name: "x".into(),
r#type: "UInt8".into(),
default_expression: "1".into(),
..Default::default()
}])
.is_err()
);
assert!(reconstruct_ddl(&[]).is_err());
assert_eq!(
reconstruct_ddl(&[DiscoveredColumn {
name: "e".into(),
r#type: "UInt8".into(),
default_kind: "EPHEMERAL".into(),
default_expression: "7".into(),
..Default::default()
}])
.unwrap(),
"e UInt8 EPHEMERAL 7"
);
}
#[test]
fn discovery_errors_carry_no_clickhouse_code() {
let err = parse_version_result(b"").unwrap_err();
assert!(matches!(err, Error::Discovery { .. }), "{err:?}");
assert_eq!(err.code(), None);
assert!(!err.is_unsupported());
}
}