use std::path::PathBuf;
use std::sync::{Arc, OnceLock};
use chtypes::{
CODE_UNSUPPORTED, Format, NO_SETTINGS, Outcome, Registry, SETTING_CLOCK_OFFSET_NANOS,
SETTING_MAX_CLOCK_SKEW_NANOS, SETTING_NOW_EPOCH_NANOS, Schema, reason,
};
static REGISTRY: OnceLock<Option<Arc<Registry>>> = OnceLock::new();
fn registry_dir() -> PathBuf {
match std::env::var_os(chtypes::REGISTRY_ENV) {
Some(dir) => PathBuf::from(dir),
None => chtypes::default_registry_dir(),
}
}
fn announce(message: &str) {
use std::io::Write;
let _ = std::io::stderr().write_all(message.as_bytes());
let _ = std::io::stderr().flush();
}
fn registry() -> Option<&'static Arc<Registry>> {
REGISTRY
.get_or_init(|| {
let dir = registry_dir();
if !dir.is_dir() {
announce(&format!(
"\nSKIP: no chtypes artifact registry at {} — fetch one with scripts/fetch.sh \
25.8 (docs/guides/fetch.md), or point ${} at a registry. Every test in this file \
is skipped, each by name below.\n",
dir.display(),
chtypes::REGISTRY_ENV
));
return None;
}
match Registry::new(&dir) {
Ok(r) if r.libraries().is_empty() => {
announce(&format!(
"\nSKIP: registry {} holds no artifact — fetch one with scripts/fetch.sh \
25.8 (docs/guides/fetch.md). Every test in this file is skipped, each by name \
below.\n",
dir.display()
));
None
}
Ok(r) => Some(Arc::new(r)),
Err(e) => {
announce(&format!(
"\nSKIP: registry at {} did not load: {e}. Every test in this file is \
skipped.\n",
dir.display()
));
None
}
}
})
.as_ref()
}
macro_rules! test_name {
() => {{
fn here() {}
let full = std::any::type_name_of_val(&here);
full.strip_suffix("::here").unwrap_or(full)
}};
}
macro_rules! registry {
() => {
match registry() {
Some(r) => r,
None => {
announce(&format!(
"\nSKIP {}: needs an artifact registry\n",
test_name!()
));
return;
}
}
};
}
fn primary(reg: &Registry) -> Arc<chtypes::Library> {
reg.for_version("25.8")
.unwrap_or_else(|_| Arc::clone(®.libraries()[0]))
}
fn stored(schema: &Schema, format: Format, body: &[u8]) -> chtypes::BatchResult {
schema.rows(format, body, NO_SETTINGS).expect("rows")
}
#[test]
fn the_registry_loads_and_libraries_name_themselves() {
let reg = registry!();
assert!(!reg.versions().is_empty());
println!("registry {}: {:?}", reg.dir().display(), reg.versions());
let libraries = reg.libraries();
let minors: Vec<&str> = libraries.iter().map(|l| l.minor()).collect();
let numeric = |m: &str| -> (u64, u64) {
let mut it = m.split('.');
(
it.next().and_then(|s| s.parse().ok()).unwrap_or(0),
it.next().and_then(|s| s.parse().ok()).unwrap_or(0),
)
};
for w in minors.windows(2) {
assert!(
numeric(w[0]) < numeric(w[1]),
"libraries() out of release order: {minors:?}"
);
}
assert_eq!(
reg.versions(),
minors.iter().map(|m| m.to_string()).collect::<Vec<_>>(),
"versions() and libraries() disagree on order"
);
for lib in reg.libraries() {
assert!(!lib.version().is_empty(), "library reported no version");
let want_minor: String = lib
.version()
.splitn(3, '.')
.take(2)
.collect::<Vec<_>>()
.join(".");
assert_eq!(lib.minor(), want_minor, "minor line of {}", lib.version());
let dir = lib.path().parent().unwrap();
let manifest: serde_json::Value =
serde_json::from_str(&std::fs::read_to_string(dir.join("manifest.json")).unwrap())
.unwrap();
assert_eq!(
lib.path().file_name().unwrap().to_str().unwrap(),
manifest["library"].as_str().unwrap(),
"library file name must come from the manifest"
);
assert_eq!(
lib.version(),
manifest["clickhouse_version"].as_str().unwrap()
);
let families = lib.registered_families().expect("chs_registered_families");
assert!(
families.len() > 100,
"{} families on {}",
families.len(),
lib.version()
);
assert!(families.iter().any(|f| f == "UInt8"));
let schema = lib
.compile("a UInt8, b UInt8 DEFAULT a + 1")
.compile()
.unwrap();
let batch = stored(&schema, Format::JsonEachRow, br#"{"a":41}"#);
assert_eq!(batch.outcome, Outcome::Accepted, "{}", lib.version());
assert_eq!(batch.rows[0].values[1].text, "42", "on {}", lib.version());
}
let flags = primary(reg).function_flags().expect("chs_function_flags");
assert!(flags.lines().count() > 1_000, "function flags look empty");
assert!(
flags.lines().any(|l| l.starts_with("now\t")),
"now() must be audited"
);
}
#[test]
fn version_resolution_accepts_a_minor_line_a_patch_and_a_drifted_patch() {
let reg = registry!();
let lib = primary(reg);
assert_eq!(
reg.for_version(lib.version()).unwrap().version(),
lib.version()
);
assert_eq!(
reg.for_version(lib.minor()).unwrap().version(),
lib.version()
);
let drifted = format!("{}.99.7-lts", lib.minor());
assert_eq!(reg.for_version(&drifted).unwrap().version(), lib.version());
let err = reg.for_version("19.1").unwrap_err();
let msg = err.to_string();
for v in reg.versions() {
assert!(msg.contains(&v), "error must name loaded versions: {msg}");
}
assert!(err.code().is_none());
}
#[test]
fn several_versions_answer_in_one_process_with_their_own_semantics() {
let reg = registry!();
if reg.libraries().len() < 2 {
announce("\nSKIP: only one artifact loaded, nothing to isolate\n");
return;
}
let libraries = reg.libraries();
let mut versions: Vec<&str> = libraries.iter().map(|l| l.version()).collect();
versions.sort_unstable();
let count = versions.len();
versions.dedup();
assert_eq!(versions.len(), count, "two libraries reported one version");
let mut json_answers = Vec::new();
let mut default_answers = Vec::new();
for lib in reg.libraries() {
let json = match lib.compile("j JSON").compile() {
Ok(schema) => {
let batch = schema
.rows(Format::JsonEachRow, br#"{"j":{"k":1}}"#, NO_SETTINGS)
.expect("rows");
match batch.outcome {
Outcome::Accepted => "accepted".to_string(),
_ => format!("code {}", batch.err_code),
}
}
Err(e) => format!("code {:?}", e.code()),
};
let mixed = match lib
.compile("a UInt8, x Int64 DEFAULT if(1,2,'a')")
.compile()
{
Ok(schema) => {
let batch = stored(&schema, Format::JsonEachRow, br#"{"a":1}"#);
assert_eq!(batch.outcome, Outcome::Accepted, "{}", lib.version());
batch.rows[0]
.values
.iter()
.find(|v| v.column == "x")
.map(|v| v.text.to_string())
.unwrap()
}
Err(e) => format!("code {:?}", e.code()),
};
println!("{:>16} JSON={json:<12} if(1,2,'a')={mixed}", lib.version());
json_answers.push((lib.minor().to_string(), json));
default_answers.push((lib.minor().to_string(), mixed));
}
for (minor, answer) in &json_answers {
match minor.as_str() {
"24.8" => assert_eq!(answer, "code 44", "24.8 must gate JSON"),
_ => assert_eq!(answer, "accepted", "{minor} must accept JSON"),
}
}
for (minor, answer) in &default_answers {
match minor.as_str() {
"25.10" => assert_eq!(
answer, "code Some(386)",
"25.10 must reject the mixed DEFAULT"
),
_ => assert_eq!(answer, "2", "{minor} must accept the mixed DEFAULT"),
}
}
let minors: Vec<&str> = json_answers.iter().map(|(m, _)| m.as_str()).collect();
if !(minors.contains(&"24.8") || minors.contains(&"25.10")) {
announce(&format!(
"\nSKIP several_versions_answer_in_one_process_with_their_own_semantics (isolation \
half): the loaded lines {minors:?} are known to agree on both probes, so no \
divergence can be observed; the proof needs 24.8 or 25.10 beside another line \
— scripts/fetch.sh 24.8 (docs/guides/fetch.md). The per-version answers were asserted.\n"
));
return;
}
let json_distinct = distinct(&json_answers);
let default_distinct = distinct(&default_answers);
assert!(
json_distinct > 1 || default_distinct > 1,
"no divergence observed across {} versions — symbol isolation unproven",
reg.libraries().len()
);
}
fn distinct(answers: &[(String, String)]) -> usize {
let mut v: Vec<&str> = answers.iter().map(|(_, a)| a.as_str()).collect();
v.sort_unstable();
v.dedup();
v.len()
}
#[test]
fn canonical_types_round_trip_in_clickhouses_own_spelling() {
let reg = registry!();
let lib = primary(reg);
for (input, canonical) in [
("DECIMAL(18,4)", "Decimal(18, 4)"),
("Decimal64(4)", "Decimal(18, 4)"),
("Nullable(Decimal(18,4))", "Nullable(Decimal(18, 4))"),
("Enum8('a'=1,'b'=2)", "Enum8('a' = 1, 'b' = 2)"),
("Map(String,Array(UInt8))", "Map(String, Array(UInt8))"),
("LowCardinality( String )", "LowCardinality(String)"),
("BIGINT", "Int64"),
("INT", "Int32"),
("Variant(UInt8, String)", "Variant(String, UInt8)"),
("Int8(3)", "Int8"),
] {
assert_eq!(
lib.validate_type(input).unwrap(),
canonical,
"canonical spelling of {input} on {}",
lib.version()
);
}
let err = lib.validate_type("NotAType").unwrap_err();
assert_eq!(err.code(), Some(50), "{err}");
assert!(!err.is_unsupported());
assert!(err.to_string().contains("NotAType"), "{err}");
assert_eq!(
lib.reference_type("UInt8").unwrap().as_deref(),
Some("Int256")
);
assert_eq!(
lib.reference_type("DateTime").unwrap().as_deref(),
Some("DateTime64(0, 'UTC')")
);
assert_eq!(lib.reference_type("String").unwrap(), None);
}
#[test]
fn compiling_ddl_is_schema_aware_in_a_way_validate_type_cannot_be() {
let reg = registry!();
let lib = primary(reg);
let schema = lib
.compile("a UInt8, b Nullable(String) DEFAULT 'x', c DateTime MATERIALIZED now()")
.compile()
.unwrap();
let cols = schema.columns();
assert_eq!(cols.len(), 3, "column introspection is missing");
assert_eq!(cols[0].ty, "UInt8");
assert_eq!(cols[0].default_kind, chtypes::DefaultKind::None);
assert_eq!(cols[1].ty, "Nullable(String)");
assert_eq!(cols[1].default_kind, chtypes::DefaultKind::Default);
assert_eq!(cols[1].default_expr, "'x'");
assert!(cols[1].default_is_literal);
assert_eq!(cols[2].default_kind, chtypes::DefaultKind::Materialized);
assert_eq!(cols[2].default_expr, "now()");
assert!(!cols[2].default_is_literal);
let schema = lib.compile("x Int64 DEFAULT NULL").compile().unwrap();
assert_eq!(schema.columns()[0].ty, "Nullable(Int64)");
assert_eq!(schema.columns()[0].default_expr, "NULL");
assert!(schema.columns()[0].default_is_literal);
let schema = lib.compile("a UInt8, al ALIAS a + 1").compile().unwrap();
assert_eq!(schema.columns()[1].ty, "UInt16");
assert_eq!(
schema.columns()[1].default_kind,
chtypes::DefaultKind::Alias
);
}
#[test]
fn an_overflowing_uint8_is_stored_wrapped_and_reported() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("x UInt8").compile().unwrap();
let batch = stored(&schema, Format::JsonEachRow, br#"{"x":256}"#);
assert_eq!(batch.outcome, Outcome::Accepted);
assert_eq!(batch.rows_read, 1);
assert_eq!(batch.rows.len(), 1);
assert_eq!(batch.rows[0].values.len(), 1);
assert_eq!(batch.rows[0].values[0].column, "x");
assert_eq!(batch.rows[0].values[0].text, "0");
assert!(!batch.rows[0].values[0].null);
assert_eq!(batch.rows[0].values[0].source, "input");
assert_eq!(batch.transformed.len(), 1, "{:?}", batch.transformed);
let t = &batch.transformed[0];
assert_eq!(t.column, "x");
assert_eq!(t.input, "256");
assert_eq!(t.stored, "0");
assert_eq!(t.reason, reason::OVERFLOW_WRAP);
assert_eq!(t.row, 0);
assert!(t.lossy());
let batch = stored(
&schema,
Format::JsonEachRow,
b"{\"x\":1}\n{\"x\":300}\n{\"x\":2}\n",
);
assert_eq!(batch.rows_read, 3);
assert_eq!(batch.transformed.len(), 1);
assert_eq!(batch.transformed[0].row, 1);
assert_eq!(batch.transformed[0].stored, "44");
}
#[test]
fn a_rejected_row_surfaces_clickhouses_own_error_code() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("x UInt8").compile().unwrap();
let batch = stored(&schema, Format::JsonEachRow, br#"{"x":"abc"}"#);
assert_eq!(batch.outcome, Outcome::Rejected);
assert_eq!(batch.err_code, 27, "{}", batch.err_msg);
assert!(batch.err_msg.contains("Cannot parse"), "{}", batch.err_msg);
assert_ne!(batch.err_code, CODE_UNSUPPORTED);
let row = schema
.row(Format::JsonEachRow, br#"{"x":"abc"}"#)
.expect("row");
assert_eq!(row.outcome, Outcome::Rejected);
assert_eq!(row.err_code, 27);
assert!(row.values.is_empty());
}
#[test]
fn a_pinned_volatile_default_stores_an_exact_timestamp() {
let reg = registry!();
let lib = primary(reg);
let schema = lib
.compile("a UInt8, ts DateTime DEFAULT now()")
.compile()
.unwrap();
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"a":1}"#,
&[(SETTING_NOW_EPOCH_NANOS, "1700000000000000000")],
)
.expect("rows");
assert_eq!(batch.outcome, Outcome::Accepted);
let ts = batch.rows[0]
.values
.iter()
.find(|v| v.column == "ts")
.expect("ts column");
assert_eq!(
ts.text, "\"2023-11-14 22:13:20\"",
"the instant must be pinned"
);
assert_eq!(ts.source, "default_substituted");
assert_eq!(batch.rows[0].substituted.len(), 1);
assert_eq!(batch.rows[0].substituted[0].column, "ts");
assert_eq!(batch.rows[0].substituted[0].expr, "now()");
assert_eq!(batch.rows[0].substituted[0].text, "\"2023-11-14 22:13:20\"");
let t = batch
.transformed
.iter()
.find(|t| t.column == "ts")
.expect("ts transform");
assert_eq!(t.reason, reason::DEFAULT_MATERIALIZED);
assert!(!t.lossy(), "materializing a DEFAULT loses nothing");
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"a":1}"#,
&[
(SETTING_CLOCK_OFFSET_NANOS, "5000000000"),
(SETTING_MAX_CLOCK_SKEW_NANOS, "1"),
],
)
.expect("rows");
assert_eq!(batch.rows[0].outcome, Outcome::Unsupported);
assert_eq!(batch.rows[0].err_code, CODE_UNSUPPORTED);
let ts = batch.rows[0]
.values
.iter()
.find(|v| v.column == "ts")
.unwrap();
assert_eq!(ts.source, "default_volatile_unresolved");
assert!(batch.transformed.iter().all(|t| t.column != "ts"));
}
#[test]
fn a_ttl_expired_row_is_reported_and_absent_from_the_stored_view() {
let reg = registry!();
let lib = primary(reg);
let mut schema = lib.compile("ts DateTime, v UInt8").compile().unwrap();
schema.set_ttl("ts + INTERVAL 1 DAY").expect("set_ttl");
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"ts":"2020-01-01 00:00:00","v":9}"#,
&[(SETTING_NOW_EPOCH_NANOS, "1700000000000000000")],
)
.expect("rows");
assert_eq!(batch.engine_rows.as_deref(), Some(&[][..]));
let t = batch
.transformed
.iter()
.find(|t| t.reason == reason::TTL_EXPIRED)
.expect("ttl_expired must be folded into the batch");
assert_eq!(t.row, 0);
assert!(t.lossy());
let mut schema = lib.compile("ts DateTime, v UInt8").compile().unwrap();
let err = schema.set_ttl("now() + INTERVAL 1 DAY").unwrap_err();
assert!(err.is_unsupported(), "{err}");
assert_eq!(err.code(), Some(CODE_UNSUPPORTED));
}
#[test]
fn an_engines_insert_time_merge_is_the_stored_truth() {
let reg = registry!();
let lib = primary(reg);
let mut schema = lib
.compile("day Date, key UInt8, v UInt64")
.compile()
.unwrap();
schema
.set_engine("SummingMergeTree", "(day, key)", NO_SETTINGS)
.expect("set_engine");
let batch = stored(
&schema,
Format::JsonEachRow,
b"{\"day\":\"2026-01-01\",\"key\":1,\"v\":5}\n{\"day\":\"2026-01-01\",\"key\":1,\"v\":7}\n",
);
assert_eq!(batch.outcome, Outcome::Accepted);
assert_eq!(batch.rows.len(), 2, "the type layer still saw both inputs");
let engine_rows = batch.engine_rows.expect("engine_rows");
assert_eq!(engine_rows.len(), 1, "{engine_rows:?}");
let merged: serde_json::Value = serde_json::from_str(&engine_rows[0]).unwrap();
assert_eq!(merged["v"], 12, "{}", engine_rows[0]);
assert_eq!(merged["day"], "2026-01-01");
let mut schema = lib.compile("id UInt8, sign Int8").compile().unwrap();
schema
.set_engine("CollapsingMergeTree(sign)", "id", NO_SETTINGS)
.expect("set_engine");
let batch = stored(&schema, Format::JsonEachRow, br#"{"id":1,"sign":3}"#);
assert_eq!(batch.outcome, Outcome::Rejected);
assert_eq!(batch.err_code, 117, "{}", batch.err_msg);
}
#[test]
fn a_server_property_default_is_unsupported_not_rejected() {
let reg = registry!();
let lib = primary(reg);
let schema = lib
.compile("h String DEFAULT hostName()")
.compile()
.unwrap();
let batch = stored(&schema, Format::JsonEachRow, b"{}");
assert_eq!(batch.outcome, Outcome::Unsupported);
assert_eq!(batch.rows[0].outcome, Outcome::Unsupported);
assert!(
batch.rows[0]
.err_msg
.contains("property of the ClickHouse server"),
"{}",
batch.rows[0].err_msg
);
assert_eq!(batch.rows[0].values[0].source, "default_expr_unsupported");
assert_eq!(batch.rows[0].values[0].text, "null");
assert!(batch.rows[0].values[0].null);
assert!(batch.transformed.is_empty());
assert_ne!(batch.outcome, Outcome::Accepted);
assert_ne!(batch.outcome, Outcome::Rejected);
assert_eq!(CODE_UNSUPPORTED, -2);
let err = lib
.compile("b UInt8 DEFAULT range(400000000)[1]")
.compile()
.unwrap_err();
assert!(err.is_unsupported(), "{err}");
assert_eq!(err.code(), Some(CODE_UNSUPPORTED));
assert!(err.to_string().contains("budget"), "{err}");
}
#[test]
fn an_accepted_insert_that_destroys_the_value_is_reported_as_accepted() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("e Enum8('a'=1,'b'=2)").compile().unwrap();
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"e":null}"#,
&[("input_format_defaults_for_omitted_fields", "0")],
)
.expect("rows");
assert_eq!(batch.outcome, Outcome::AcceptedPoisoned);
assert_eq!(batch.err_code, 691);
assert_eq!(batch.rows[0].values[0].text, "");
assert!(!batch.rows[0].values[0].null);
let t = &batch.transformed[0];
assert_eq!(t.reason, reason::POISONED);
assert_eq!(t.stored, "<unreadable>");
assert!(t.lossy());
}
#[test]
fn a_schema_can_move_between_threads_and_answers_stably() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("x UInt8, s String").compile().unwrap();
let handle = std::thread::spawn(move || {
let mut last = String::new();
for _ in 0..2_000 {
let batch = schema
.rows(Format::JsonEachRow, br#"{"x":256,"s":"hi"}"#, NO_SETTINGS)
.expect("rows");
assert_eq!(batch.outcome, Outcome::Accepted);
let rendered = format!("{:?}", batch.rows[0].values);
if !last.is_empty() {
assert_eq!(rendered, last, "answers drifted across calls");
}
last = rendered;
}
last
});
let rendered = handle.join().expect("thread");
assert!(rendered.contains("hi"), "{rendered}");
assert!(!rendered.contains("overflow"), "{rendered}");
}
#[test]
fn the_row_and_batch_entry_points_agree_on_one_row() {
let reg = registry!();
let lib = primary(reg);
let schema = lib
.compile("x UInt8, m UInt8 MATERIALIZED x + 1, s String")
.compile()
.unwrap();
let batch = stored(&schema, Format::Csv, b"7,\"hey\"\n");
assert_eq!(batch.outcome, Outcome::Accepted, "{}", batch.err_msg);
let row = &batch.rows[0];
assert_eq!(row.values.len(), 2, "{:?}", row.values);
assert_eq!(row.values[0].text, "7");
assert_eq!(row.values[1].text, "\"hey\"");
assert_eq!(row.computed.len(), 1);
assert_eq!(row.computed[0].column, "m");
assert_eq!(row.computed[0].kind, "MATERIALIZED");
assert_eq!(row.computed[0].text, "8");
let json = stored(&schema, Format::JsonEachRow, br#"{"x":7,"s":"hey"}"#);
assert_eq!(json.rows[0].computed[0].text, "8");
assert_eq!(json.rows[0].values.len(), 2);
let single = schema.row(Format::Csv, b"7,\"hey\"").expect("row");
assert_eq!(single.outcome, row.outcome);
assert_eq!(single.values, row.values);
assert_eq!(single.computed, row.computed);
}
#[test]
fn an_empty_body_is_accepted_with_zero_rows() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("x UInt8").compile().unwrap();
let batch = stored(&schema, Format::JsonEachRow, b"");
assert_eq!(batch.outcome, Outcome::Accepted);
assert_eq!(batch.rows_read, 0);
assert!(batch.rows.is_empty());
assert!(batch.transformed.is_empty());
}
fn compile_settings_lib(reg: &Registry) -> Option<Arc<chtypes::Library>> {
let Ok(lib) = reg.for_version("25.8") else {
announce("\nSKIP: no 25.8 artifact loaded; the compile-settings tests target 25.8\n");
return None;
};
Some(lib)
}
#[test]
fn every_loaded_artifact_reports_compile_settings_support() {
let reg = registry!();
for lib in reg.libraries() {
assert!(
lib.has_compile_settings(),
"{} must report compile-settings support: chs_schema_compile is one of \
the four mandatory symbols",
lib.version()
);
}
}
#[test]
fn a_declared_flatten_nested_zero_compiles_one_nested_column() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let schema = lib
.compile("n Nested(a Int64, b String)")
.settings([("flatten_nested", "0")])
.compile()
.expect("compile with a settings profile");
assert_eq!(schema.columns().len(), 1, "{:?}", schema.columns());
assert_eq!(schema.columns()[0].name, "n");
assert!(
schema.columns()[0].ty.starts_with("Nested("),
"type = {}",
schema.columns()[0].ty
);
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"n":[{"a":1,"b":"x"}]}"#,
&[
("flatten_nested", "0"),
("input_format_import_nested_json", "1"),
],
)
.expect("rows");
assert_eq!(batch.outcome, Outcome::Accepted, "{}", batch.err_msg);
assert_eq!(batch.rows.len(), 1);
assert_eq!(batch.rows[0].values.len(), 1);
assert_eq!(batch.rows[0].values[0].text, r#"[{"a":1,"b":"x"}]"#);
let batch = schema
.rows(
Format::JsonEachRow,
br#"{"n.a":[1,2],"n.b":["x","y"]}"#,
&[
("flatten_nested", "0"),
("input_format_import_nested_json", "1"),
],
)
.expect("rows");
assert_eq!(batch.outcome, Outcome::Accepted, "{}", batch.err_msg);
assert_eq!(batch.rows.len(), 1);
assert_eq!(
batch.rows[0].unknown_fields.len(),
2,
"dotted keys: {:?}",
batch.rows[0].unknown_fields
);
}
#[test]
fn an_explicit_empty_settings_call_matches_the_bare_compile() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let ddl = "n Nested(a Int64, b String), x UInt8";
let bare = lib.compile(ddl).compile().expect("compile");
let explicit = lib
.compile(ddl)
.settings(Vec::<(&str, &str)>::new())
.compile()
.expect("compile with an explicit empty settings call");
assert_eq!(
bare.columns(),
explicit.columns(),
"an empty settings profile must produce the identical column list"
);
assert_eq!(bare.columns().len(), 3, "{:?}", bare.columns());
}
#[test]
fn the_compile_settings_error_channels_stay_distinct() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let err = lib
.compile("x UInt8")
.settings([("made_up_setting_xyz", "1")])
.compile()
.unwrap_err();
assert_eq!(err.code(), Some(115), "unknown setting: {err}");
assert!(!err.is_unsupported());
let err = lib
.compile("x UInt8")
.settings([(chtypes::SETTING_NOW_EPOCH_NANOS, "1")])
.compile()
.unwrap_err();
assert_eq!(err.code(), Some(115), "chtypes_* at compile: {err}");
}
#[test]
fn set_engine_validates_merge_tree_settings_and_refuses_non_defaults() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let mut schema = lib.compile("id UInt64").compile().unwrap();
schema
.set_engine("MergeTree", "tuple()", NO_SETTINGS)
.expect("empty settings");
schema
.set_engine("MergeTree", "tuple()", &[("allow_nullable_key", "0")])
.expect("declared-at-default");
let err = schema
.set_engine("MergeTree", "tuple()", &[("allow_nullable_key", "1")])
.unwrap_err();
assert!(err.is_unsupported(), "non-default MergeTree setting: {err}");
assert_eq!(err.code(), Some(CODE_UNSUPPORTED));
let err = schema
.set_engine("MergeTree", "tuple()", &[("totally_made_up_mt", "1")])
.unwrap_err();
assert_eq!(err.code(), Some(115), "unknown MergeTree name: {err}");
assert!(!err.is_unsupported());
}
#[test]
fn a_reconstructed_table_compiles_through_the_library_itself() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
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"}
{"name":"e","type":"UInt8","default_kind":"EPHEMERAL","default_expression":"","position":4}
{"name":"m","type":"UInt64","default_kind":"MATERIALIZED","default_expression":"id + 1","position":"5"}
"#;
let cols = chtypes::parse_columns_result(body).expect("parse_columns_result");
assert_eq!(cols.len(), 5);
assert_eq!(
cols.iter().map(|c| c.position).collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5]
);
let ddl = chtypes::reconstruct_ddl(&cols).expect("reconstruct_ddl");
let schema = lib
.compile(&ddl)
.compile()
.unwrap_or_else(|e| panic!("reconstructed DDL does not compile: {e}\n{ddl}"));
assert_eq!(schema.columns().len(), 5, "{:?}", schema.columns());
assert_eq!(schema.columns()[2].name, "n.a");
}
#[test]
fn a_discovered_profile_compiles_the_unflattened_shape() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let settings =
chtypes::parse_changed_settings_result(b"{\"name\":\"flatten_nested\",\"value\":\"0\"}\n")
.expect("parse_changed_settings_result");
let schema = lib
.compile("n Nested(a Int64, b String)")
.settings(settings)
.compile()
.expect("compile with the discovered profile");
assert_eq!(schema.columns().len(), 1, "{:?}", schema.columns());
assert_eq!(schema.columns()[0].name, "n");
assert!(schema.columns()[0].ty.starts_with("Nested("));
}
#[test]
fn rowbinary_is_read_in_clickhouses_storage_encoding() {
let reg = registry!();
let lib = primary(reg);
let schema = lib.compile("x UInt8").compile().unwrap();
let batch = stored(&schema, Format::RowBinary, &[0x00]);
assert_eq!(batch.outcome, Outcome::Accepted, "{}", batch.err_msg);
assert_eq!(batch.rows[0].values[0].text, "0");
assert_eq!(batch.rows[0].transformed, vec![]);
let batch = stored(
&schema,
Format::RowBinaryWithNamesAndTypesAndDefaults,
&[0x01],
);
assert_eq!(batch.outcome, Outcome::Rejected);
assert!(
batch.err_code == 73 || batch.err_code == 32,
"unexpected code {} on {}: {}",
batch.err_code,
lib.version(),
batch.err_msg
);
}
const ISO_Z_ROW: &[u8] = br#"{"ts":"2020-01-02T03:04:05Z"}"#;
#[test]
fn a_compile_profiles_settings_govern_row_calls() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let be = [("date_time_input_format", "best_effort")];
let basic = [("date_time_input_format", "basic")];
for (what, profile, per_call, want) in [
(
"profile best_effort, no per-call",
&be,
NO_SETTINGS,
Outcome::Accepted,
),
(
"profile basic, no per-call",
&basic,
NO_SETTINGS,
Outcome::Rejected,
),
(
"profile basic, per-call best_effort",
&basic,
&be[..],
Outcome::Accepted,
),
(
"profile best_effort, per-call basic",
&be,
&basic[..],
Outcome::Rejected,
),
] {
let schema = lib
.compile("ts DateTime")
.settings(profile.iter().copied())
.compile()
.unwrap_or_else(|e| panic!("{what}: compile: {e}"));
let batch = schema
.rows(Format::JsonEachRow, ISO_Z_ROW, per_call)
.unwrap_or_else(|e| panic!("{what}: rows: {e}"));
assert_eq!(
batch.outcome, want,
"{what}: code {:?} {}",
batch.err_code, batch.err_msg
);
}
}
#[test]
fn a_profile_less_handle_is_untouched_by_the_profile_channel() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let schema = lib.compile("ts DateTime").compile().expect("compile");
let forced = schema
.rows(
Format::JsonEachRow,
ISO_Z_ROW,
&[("date_time_input_format", "basic")],
)
.expect("rows");
assert_eq!(forced.outcome, Outcome::Rejected, "{}", forced.err_msg);
let bare = schema
.rows(Format::JsonEachRow, ISO_Z_ROW, NO_SETTINGS)
.expect("rows");
assert!(matches!(
bare.outcome,
Outcome::Accepted | Outcome::Rejected
));
}
#[test]
fn set_engine_tells_a_server_refusal_from_a_decline() {
let reg = registry!();
let Some(lib) = compile_settings_lib(reg) else {
return;
};
let mut schema = lib.compile("id UInt64, sign Int8").compile().unwrap();
let err = schema
.set_engine("MergeTree", "tuple()", &[("totally_made_up_mt", "1")])
.unwrap_err();
match &err {
chtypes::Error::Schema { code, message, .. } => {
assert_eq!(*code, 115, "unknown MergeTree name: {err}");
assert!(message.contains("totally_made_up_mt"), "{message}");
}
other => panic!("unknown MergeTree name must be Error::Schema, got {other:?}"),
}
assert!(!err.is_unsupported());
for (what, err) in [
(
"unmodeled engine",
schema
.set_engine("NotAnEngine", "tuple()", NO_SETTINGS)
.unwrap_err(),
),
(
"non-default MergeTree setting",
schema
.set_engine("MergeTree", "tuple()", &[("allow_nullable_key", "1")])
.unwrap_err(),
),
(
"sorting key not in the schema",
schema
.set_engine("MergeTree", "nosuchcolumn", NO_SETTINGS)
.unwrap_err(),
),
] {
assert!(
matches!(err, chtypes::Error::Unsupported { .. }),
"{what} must be Error::Unsupported, got {err:?}"
);
assert!(err.is_unsupported(), "{what}: {err}");
assert_eq!(err.code(), Some(CODE_UNSUPPORTED), "{what}");
}
}
#[test]
fn every_artifact_reports_a_compatible_abi_revision() {
let reg = registry!();
const { assert!(chtypes::ABI_REVISION > 0, "ABI_REVISION must be positive") };
let mut seen = 0usize;
let mut predates = Vec::new();
for lib in reg.libraries() {
let rev = lib.abi_revision();
assert!(
rev == 0 || rev == chtypes::ABI_REVISION,
"{}: Library exists with ABI revision {rev}, but the loader must \
refuse anything but {} or 0",
lib.version(),
chtypes::ABI_REVISION
);
if rev == 0 {
predates.push(lib.minor().to_string());
}
seen += 1;
}
assert!(seen > 0, "no artifact was probed for its ABI revision");
announce(&format!(
"\nABI revision {}: {} artifact(s) current, {} predate the probe {predates:?}\n",
chtypes::ABI_REVISION,
seen - predates.len(),
predates.len()
));
}
#[test]
fn one_image_means_one_lock_across_registries() {
let shared = registry!();
let second = match Registry::new(registry_dir()) {
Ok(r) => Arc::new(r),
Err(err) => {
announce(&format!("\nSKIP one_image_means_one_lock: {err}\n"));
return;
}
};
let version = shared.versions().last().expect("a version").to_string();
let body = br#"{"a": 1}"#;
let want = {
let lib = shared.for_version(&version).expect("library");
let schema = lib.compile("a UInt8").compile().expect("compile");
schema
.rows(Format::JsonEachRow, body, NO_SETTINGS)
.expect("rows")
.outcome
};
assert_eq!(want, Outcome::Accepted);
const WANT_READS: usize = 100;
const WANT_SWAPS: usize = 50;
const BOUND: std::time::Duration = std::time::Duration::from_secs(30);
let deadline = std::time::Instant::now() + BOUND;
let reads = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let swaps = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let ord = std::sync::atomic::Ordering::Relaxed;
let done = |r: usize, w: usize| {
(r > WANT_READS && w > WANT_SWAPS) || std::time::Instant::now() >= deadline
};
std::thread::scope(|scope| {
for _ in 0..4 {
let reg = Arc::clone(shared);
let reads = Arc::clone(&reads);
let swaps = Arc::clone(&swaps);
let version = version.clone();
scope.spawn(move || {
let lib = reg.for_version(&version).expect("library");
let schema = lib.compile("a UInt8").compile().expect("compile");
while !done(reads.load(ord), swaps.load(ord)) {
let got = schema
.rows(Format::JsonEachRow, body, NO_SETTINGS)
.expect("rows");
assert_eq!(got.outcome, want);
reads.fetch_add(1, ord);
}
});
}
let writer_reg = Arc::clone(&second);
let swaps_w = Arc::clone(&swaps);
let reads_w = Arc::clone(&reads);
let version_w = version.clone();
scope.spawn(move || {
let lib = writer_reg.for_version(&version_w).expect("library");
while !done(reads_w.load(ord), swaps_w.load(ord)) {
let n = swaps_w.fetch_add(1, ord);
let pairs: Vec<(&str, &str)> = if n % 2 == 0 {
vec![]
} else {
vec![("chtypes_default_eval_wall_nanos", "2000000000")]
};
lib.set_default_settings(&pairs).expect("seed");
std::thread::sleep(std::time::Duration::from_millis(1));
}
});
});
shared
.for_version(&version)
.expect("library")
.set_default_settings(NO_SETTINGS)
.expect("reset");
let (r, w) = (reads.load(ord), swaps.load(ord));
if r <= WANT_READS || w <= WANT_SWAPS {
announce(&format!(
"\nSKIP one_image_means_one_lock_across_registries: this machine did not \
produce contention within {BOUND:?} — {r} reads (wanted > {WANT_READS}) against \
{w} settings swaps (wanted > {WANT_SWAPS}). Every read that did run matched the \
uncontended answer; too few of them overlapped to prove the lock. Re-run on \
an idle machine.\n"
));
return;
}
announce(&format!(
"\none image, one lock: {r} batch reads across 4 threads on registry A \
against {w} settings swaps on registry B, 0 mismatches\n"
));
}