use crate::raw::RawText;
use crate::transform::classify;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[repr(i32)]
pub enum Format {
JsonEachRow = 0,
Csv = 1,
Tsv = 2,
Values = 3,
JsonCompactEachRow = 4,
RowBinary = 5,
RowBinaryWithDefaults = 6,
RowBinaryWithNamesAndTypesAndDefaults = 7,
Native = 8,
Buffers = 9,
}
impl Format {
pub fn code(self) -> i32 {
self as i32
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub struct DocFlags(u32);
impl DocFlags {
pub const NONE: DocFlags = DocFlags(0);
pub const VALUES: DocFlags = DocFlags(0x1);
pub const TRANSFORMS: DocFlags = DocFlags(0x2);
pub const DEFAULTS: DocFlags = DocFlags(0x4);
pub const ALL: DocFlags = DocFlags(0x7);
pub fn bits(self) -> u32 {
self.0
}
pub fn from_bits(bits: u32) -> DocFlags {
DocFlags(bits)
}
}
impl std::ops::BitOr for DocFlags {
type Output = DocFlags;
fn bitor(self, rhs: DocFlags) -> DocFlags {
DocFlags(self.0 | rhs.0)
}
}
impl std::ops::BitOrAssign for DocFlags {
fn bitor_assign(&mut self, rhs: DocFlags) {
self.0 |= rhs.0;
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub struct Span {
pub off: usize,
pub len: usize,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub enum Verdict {
#[default]
Decline,
True,
False,
Error,
}
impl Verdict {
pub fn answered(self) -> bool {
matches!(self, Verdict::True | Verdict::False)
}
pub fn as_char(self) -> char {
match self {
Verdict::True => 't',
Verdict::False => 'f',
Verdict::Error => 'e',
Verdict::Decline => 'd',
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub enum FilterOutcome {
Ok,
Rejected,
#[default]
Unsupported,
}
impl FilterOutcome {
pub fn as_str(self) -> &'static str {
match self {
FilterOutcome::Ok => "ok",
FilterOutcome::Rejected => "rejected",
FilterOutcome::Unsupported => "unsupported",
}
}
}
impl std::fmt::Display for FilterOutcome {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct FilterRowError {
pub row: usize,
pub code: i32,
pub err: String,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct FilterResult {
pub outcome: FilterOutcome,
pub err_code: i32,
pub err_msg: String,
pub rows_read: usize,
pub unsupported_settings: Vec<String>,
pub verdicts: Vec<Verdict>,
pub errors: Vec<FilterRowError>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub enum Outcome {
Accepted,
Rejected,
AcceptedPoisoned,
#[default]
Unsupported,
Skipped,
}
impl Outcome {
pub fn as_str(self) -> &'static str {
match self {
Outcome::Accepted => "accepted",
Outcome::Rejected => "rejected",
Outcome::AcceptedPoisoned => "accepted_poisoned",
Outcome::Unsupported => "unsupported",
Outcome::Skipped => "skipped",
}
}
}
impl std::fmt::Display for Outcome {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
fn outcome_of(s: &str) -> Outcome {
match s {
"accepted" => Outcome::Accepted,
"rejected" => Outcome::Rejected,
"accepted_poisoned" => Outcome::AcceptedPoisoned,
"unsupported" => Outcome::Unsupported,
"skipped" => Outcome::Skipped,
_ => Outcome::Unsupported,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Value {
pub column: String,
pub text: RawText,
pub null: bool,
pub source: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Transform {
pub column: String,
pub input: String,
pub stored: String,
pub reason: String,
pub row: usize,
}
impl Transform {
pub fn lossy(&self) -> bool {
!matches!(
self.reason.as_str(),
crate::reason::REFORMAT
| crate::reason::DEFAULT_FILLED
| crate::reason::ZERO_FILLED
| crate::reason::DEFAULT_MATERIALIZED
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Substitution {
pub column: String,
pub expr: String,
pub text: RawText,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Computed {
pub column: String,
pub kind: String,
pub text: RawText,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RowResult {
pub values: Vec<Value>,
pub transformed: Vec<Transform>,
pub outcome: Outcome,
pub err_code: i32,
pub err_msg: String,
pub unknown_fields: Vec<String>,
pub unsupported_settings: Vec<String>,
pub substituted: Vec<Substitution>,
pub computed: Vec<Computed>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct BatchResult {
pub rows: Vec<RowResult>,
pub outcome: Outcome,
pub err_code: i32,
pub err_msg: String,
pub rows_read: usize,
pub rows_skipped: usize,
pub transformed: Vec<Transform>,
pub engine_rows: Option<Vec<String>>,
engine_rows_raw: Option<Vec<RawText>>,
pub payload: Option<Vec<u8>>,
pub spans: Option<Vec<Span>>,
pub export_declined: String,
}
impl BatchResult {
pub fn engine_rows_bytes(&self) -> Option<&[RawText]> {
self.engine_rows_raw.as_deref()
}
}
#[derive(Debug, Default)]
pub(crate) struct RowDoc {
pub outcome: String,
pub code: i32,
pub err: String,
pub unknown_fields: Vec<String>,
pub unsupported_settings: Vec<String>,
pub cols: Vec<ColDoc>,
pub computed: Vec<CompDoc>,
}
#[derive(Debug, Default)]
pub(crate) struct CompDoc {
pub name: String,
pub kind: String,
pub stored: Option<RawText>,
}
#[derive(Debug, Default)]
pub(crate) struct ColDoc {
pub name: String,
#[allow(dead_code)]
pub declared_type: String,
pub base: String,
pub src: String,
pub input: RawText,
pub ref_type: String,
pub stored: Option<RawText>,
pub reference: Option<RawText>,
pub nullable: bool,
pub poison: bool,
pub null_input: bool,
pub dup_dropped: bool,
pub ref_unclassified: bool,
}
impl ColDoc {
pub(crate) fn stored_raw(&self) -> RawText {
match (&self.stored, self.poison) {
(Some(v), false) => v.clone(),
_ => RawText::default(),
}
}
pub(crate) fn ref_raw(&self) -> RawText {
match &self.reference {
Some(v) if v != "null" => v.clone(),
_ => RawText::default(),
}
}
pub(crate) fn stored_is_json_null(&self) -> bool {
self.stored.as_ref().is_some_and(|v| v == "null")
}
}
#[derive(Debug, Default)]
pub(crate) struct BatchDoc {
pub outcome: String,
pub code: i32,
pub err: String,
pub rows_read: usize,
pub rows_skipped: usize,
pub rows: Vec<RowDoc>,
pub engine_rows: Option<Vec<RawText>>,
pub storage_transforms: Vec<StorageTransformDoc>,
pub row_spans: Option<Vec<Span>>,
pub export_declined: String,
}
#[derive(Debug, Default)]
pub(crate) struct FilterDoc {
pub outcome: String,
pub code: i32,
pub err: String,
pub rows_read: usize,
pub unsupported_settings: Vec<String>,
pub verdicts: String,
pub errors: Vec<FilterRowError>,
}
#[derive(Debug, Default)]
pub(crate) struct StorageTransformDoc {
pub row: usize,
pub column: String,
pub reason: String,
pub stored: Option<RawText>,
}
pub(crate) fn row_result_of(doc: RowDoc) -> RowResult {
let mut res = RowResult {
outcome: outcome_of(&doc.outcome),
err_code: doc.code,
err_msg: doc.err,
unknown_fields: doc.unknown_fields,
unsupported_settings: doc.unsupported_settings,
..Default::default()
};
if !res.unsupported_settings.is_empty() && res.outcome != Outcome::Rejected {
res.outcome = Outcome::Unsupported;
}
for c in &doc.cols {
if c.src == "skipped" {
continue;
}
res.values.push(Value {
column: c.name.clone(),
text: c.stored_raw(),
null: c.stored_is_json_null() && !c.poison,
source: c.src.clone(),
});
if c.src == "default_substituted" {
res.substituted.push(Substitution {
column: c.name.clone(),
expr: c.input.to_lossy().into_owned(),
text: c.stored_raw(),
});
}
res.transformed.extend(classify(c));
}
for m in doc.computed {
res.computed.push(Computed {
column: m.name,
kind: m.kind,
text: m.stored.unwrap_or_default(),
});
}
res
}
pub(crate) fn batch_result_of(doc: BatchDoc) -> BatchResult {
let mut res = BatchResult {
outcome: outcome_of(&doc.outcome),
err_code: doc.code,
err_msg: doc.err,
rows_read: doc.rows_read,
rows_skipped: doc.rows_skipped,
engine_rows: doc.engine_rows.as_ref().map(|rows| {
rows.iter()
.map(|r| r.to_lossy().into_owned())
.collect::<Vec<_>>()
}),
engine_rows_raw: doc.engine_rows,
spans: doc.row_spans,
export_declined: doc.export_declined,
..Default::default()
};
for (i, rd) in doc.rows.into_iter().enumerate() {
let mut rr = row_result_of(rd);
for t in &mut rr.transformed {
t.row = i;
}
res.transformed.extend(rr.transformed.iter().cloned());
res.rows.push(rr);
}
for st in doc.storage_transforms {
res.transformed.push(Transform {
column: st.column,
input: String::new(),
stored: st
.stored
.map(|v| v.to_lossy().into_owned())
.unwrap_or_default(),
reason: st.reason,
row: st.row,
});
}
res
}
pub(crate) fn filter_result_of(doc: FilterDoc) -> FilterResult {
let outcome = match doc.outcome.as_str() {
"ok" => FilterOutcome::Ok,
"rejected" => FilterOutcome::Rejected,
_ => FilterOutcome::Unsupported,
};
let mut res = FilterResult {
outcome,
err_code: doc.code,
err_msg: doc.err,
rows_read: doc.rows_read,
unsupported_settings: doc.unsupported_settings,
..Default::default()
};
if res.outcome != FilterOutcome::Ok {
return res;
}
res.verdicts = doc
.verdicts
.chars()
.map(|c| match c {
't' => Verdict::True,
'f' => Verdict::False,
'e' => Verdict::Error,
_ => Verdict::Decline,
})
.collect();
res.errors = doc.errors;
res
}
#[cfg(test)]
mod tests {
use super::*;
use crate::json::quote_bare_denormals;
fn parse_row(s: &str) -> RowResult {
let repaired = quote_bare_denormals(s.as_bytes());
row_result_of(crate::doc::row_doc(&repaired).unwrap())
}
#[test]
fn an_unknown_outcome_degrades_to_unsupported_never_rejected() {
let rr = parse_row(r#"{"outcome":"verdict_from_the_future","code":0,"err":"","cols":[]}"#);
assert_eq!(rr.outcome, Outcome::Unsupported);
assert_eq!(outcome_of("accepted"), Outcome::Accepted);
assert_eq!(outcome_of("rejected"), Outcome::Rejected);
assert_eq!(outcome_of("accepted_poisoned"), Outcome::AcceptedPoisoned);
assert_eq!(outcome_of("unsupported"), Outcome::Unsupported);
assert_eq!(outcome_of("skipped"), Outcome::Skipped);
assert_eq!(outcome_of(""), Outcome::Unsupported);
}
#[test]
fn a_skipped_row_parses_with_the_caught_error_and_no_values() {
let rr =
parse_row(r#"{"outcome":"skipped","code":27,"err":"Cannot parse input","cols":[]}"#);
assert_eq!(rr.outcome, Outcome::Skipped);
assert_eq!(rr.err_code, 27);
assert_eq!(rr.err_msg, "Cannot parse input");
assert!(rr.values.is_empty());
}
#[test]
fn a_rejected_document_omits_most_keys() {
let rr =
parse_row(r#"{"outcome":"rejected","code":27,"err":"Cannot parse input","cols":[]}"#);
assert_eq!(rr.outcome, Outcome::Rejected);
assert_eq!(rr.err_code, 27);
assert!(rr.values.is_empty());
assert!(rr.unknown_fields.is_empty());
}
#[test]
fn overflow_wrap_is_derived_from_the_reference_parse() {
let rr = parse_row(
r#"{"outcome":"accepted","code":0,"err":"","cols":[
{"name":"x","type":"UInt8","base":"UInt8","src":"input","input":"256",
"stored":0,"ref":256,"ref_type":"Int256",
"nullable":false,"poison":false,"null_input":false,"dup_dropped":false}]}"#,
);
assert_eq!(rr.outcome, Outcome::Accepted);
assert_eq!(rr.values[0].text, "0");
assert_eq!(rr.transformed.len(), 1);
assert_eq!(rr.transformed[0].reason, crate::reason::OVERFLOW_WRAP);
assert!(rr.transformed[0].lossy());
}
#[test]
fn declined_settings_promote_the_row_to_unsupported() {
let rr = parse_row(
r#"{"outcome":"accepted","unsupported_settings":["some_setting"],"cols":[]}"#,
);
assert_eq!(rr.outcome, Outcome::Unsupported);
}
#[test]
fn storage_transforms_are_folded_into_the_batch() {
let doc = crate::doc::batch_doc(
br#"{"outcome":"accepted","code":0,"err":"","engine_rows":[],
"storage_transforms":[{"row":0,"column":"","reason":"ttl_expired"}],
"rows_read":1,"rows_skipped":0,
"rows":[{"outcome":"accepted","cols":[]}]}"#,
)
.unwrap();
let br = batch_result_of(doc);
assert_eq!(br.outcome, Outcome::Accepted);
assert_eq!(br.engine_rows.as_deref(), Some(&[][..]));
assert_eq!(br.engine_rows_bytes().map(<[RawText]>::len), Some(0));
assert_eq!(br.rows.len(), 1);
assert_eq!(br.transformed.len(), 1);
assert_eq!(br.transformed[0].reason, crate::reason::TTL_EXPIRED);
assert!(br.transformed[0].lossy());
}
#[test]
fn engine_rows_keep_clickhouses_bytes_in_the_authoritative_view() {
let mut doc: Vec<u8> = br#"{"outcome":"accepted","engine_rows":[{"x":""#.to_vec();
doc.extend_from_slice(&[0xc3, b'(']);
doc.extend_from_slice(br#""}],"rows":[]}"#);
let br = batch_result_of(crate::doc::batch_doc(&doc).unwrap());
assert_eq!(
br.engine_rows_bytes().unwrap()[0].as_bytes(),
[b'{', b'"', b'x', b'"', b':', b'"', 0xc3, b'(', b'"', b'}'],
);
assert_eq!(br.engine_rows.as_ref().unwrap()[0], "{\"x\":\"\u{fffd}(\"}");
}
#[test]
fn a_stored_null_is_a_value_and_an_absent_one_is_not() {
let rr = parse_row(
r#"{"outcome":"accepted","cols":[
{"name":"a","type":"Nullable(UInt8)","base":"UInt8","src":"input",
"input":"null","stored":null,"nullable":true,"null_input":true},
{"name":"b","type":"Nullable(UInt8)","base":"UInt8","src":"input",
"input":"null","nullable":true,"null_input":true}]}"#,
);
assert_eq!(
(rr.values[0].text.as_str(), rr.values[0].null),
(Some("null"), true)
);
assert_eq!(
(rr.values[1].text.as_str(), rr.values[1].null),
(Some(""), false)
);
}
#[test]
fn denormals_survive_as_text() {
let rr = parse_row(
r#"{"outcome":"accepted","cols":[
{"name":"f","type":"Float64","base":"Float64","src":"input",
"input":"1e400","stored":inf,"ref":null,"ref_type":""}]}"#,
);
assert_eq!(rr.values[0].text, "\"inf\"");
assert_eq!(rr.transformed[0].reason, crate::reason::LOSSY_NUMERIC);
}
#[test]
fn format_codes_are_the_abi_numbers() {
assert_eq!(Format::JsonEachRow.code(), 0);
assert_eq!(Format::Csv.code(), 1);
assert_eq!(Format::Tsv.code(), 2);
assert_eq!(Format::Values.code(), 3);
assert_eq!(Format::JsonCompactEachRow.code(), 4);
assert_eq!(Format::RowBinary.code(), 5);
assert_eq!(Format::RowBinaryWithDefaults.code(), 6);
assert_eq!(Format::RowBinaryWithNamesAndTypesAndDefaults.code(), 7);
assert_eq!(Format::Native.code(), 8);
assert_eq!(Format::Buffers.code(), 9);
}
}