use std::collections::HashMap;
use snafu::prelude::*;
use crate::metadata::{
logical_schema::{LogicalDataType, LogicalField, LogicalSchema, LogicalSchemaError},
table_metadata::{TableMeta, TimeIndexSpec},
};
#[derive(Debug, Snafu)]
pub enum SchemaCompatibilityError {
#[snafu(display("Table has no logical_schema; v0.1 cannot append without a canonical schema"))]
MissingTableSchema,
#[snafu(display("Segment schema is missing required column {column}"))]
MissingColumn {
column: String,
},
#[snafu(display("Segment schema has extra column {column} not present in table schema"))]
ExtraColumn {
column: String,
},
#[snafu(display(
"Type mismatch for column {column}: table has {table_type}, segment has {segment_type}"
))]
TypeMismatch {
column: String,
table_type: LogicalDataType,
segment_type: LogicalDataType,
},
#[snafu(display(
"Time index column {column} has incompatible type: table has {table_type}, \
segment has {segment_type}"
))]
TimeIndexTypeMismatch {
column: String,
table_type: LogicalDataType,
segment_type: LogicalDataType,
},
#[snafu(display("Logical schema is invalid: {source}"))]
LogicalSchema {
#[snafu(source)]
source: LogicalSchemaError,
},
}
pub type SchemaResult<T> = Result<T, SchemaCompatibilityError>;
pub fn require_table_schema(meta: &TableMeta) -> SchemaResult<&LogicalSchema> {
match &meta.logical_schema {
Some(schema) => Ok(schema),
None => MissingTableSchemaSnafu.fail(),
}
}
fn columns_by_name(schema: &LogicalSchema) -> HashMap<&str, &LogicalField> {
schema
.columns()
.iter()
.map(|col| (col.name.as_str(), col))
.collect()
}
pub fn ensure_schema_exact_match(
table_schema: &LogicalSchema,
segment_schema: &LogicalSchema,
index: &TimeIndexSpec,
) -> SchemaResult<()> {
let time_col_name = index.timestamp_column.as_str();
let table_cols = columns_by_name(table_schema);
let seg_cols = columns_by_name(segment_schema);
for (name, table_field) in &table_cols {
let seg_field =
seg_cols
.get(name)
.ok_or_else(|| SchemaCompatibilityError::MissingColumn {
column: (*name).to_string(),
})?;
if table_field.data_type != seg_field.data_type
|| table_field.nullable != seg_field.nullable
{
let err = if *name == time_col_name {
SchemaCompatibilityError::TimeIndexTypeMismatch {
column: (*name).to_string(),
table_type: table_field.data_type.clone(),
segment_type: seg_field.data_type.clone(),
}
} else {
SchemaCompatibilityError::TypeMismatch {
column: (*name).to_string(),
table_type: table_field.data_type.clone(),
segment_type: seg_field.data_type.clone(),
}
};
return Err(err);
}
}
for name in seg_cols.keys() {
if !table_cols.contains_key(name) {
return Err(SchemaCompatibilityError::ExtraColumn {
column: (*name).to_string(),
});
}
}
Ok(())
}