use crate::{
error::CompatibilityError,
schema::{
ArraySchema, DecimalSchema, EnumSchema, InnerDecimalSchema, MapSchema, RecordSchema,
Schema, UuidSchema,
},
};
use std::{
collections::{HashMap, hash_map::DefaultHasher},
hash::Hasher,
iter::once,
ops::BitAndAssign,
ptr,
};
pub struct SchemaCompatibility;
impl SchemaCompatibility {
pub fn can_read(
writers_schema: &Schema,
readers_schema: &Schema,
) -> Result<Compatibility, CompatibilityError> {
let mut c = Checker::new();
c.can_read(writers_schema, readers_schema)
}
pub fn mutual_read(
schema_a: &Schema,
schema_b: &Schema,
) -> Result<Compatibility, CompatibilityError> {
let mut c = SchemaCompatibility::can_read(schema_a, schema_b)?;
c &= SchemaCompatibility::can_read(schema_b, schema_a)?;
Ok(c)
}
}
#[derive(Debug, Copy, Clone, Eq, PartialEq)]
pub enum Compatibility {
Full,
Partial,
}
impl BitAndAssign for Compatibility {
fn bitand_assign(&mut self, rhs: Self) {
match (*self, rhs) {
(Self::Full, Self::Full) => *self = Self::Full,
_ => *self = Self::Partial,
}
}
}
struct Checker {
recursion: HashMap<(u64, u64), Compatibility>,
}
impl Checker {
pub(crate) fn new() -> Self {
Self {
recursion: HashMap::new(),
}
}
pub(crate) fn full_match_schemas(
&mut self,
writers_schema: &Schema,
readers_schema: &Schema,
) -> Result<Compatibility, CompatibilityError> {
let key = (
Self::pointer_hash(writers_schema),
Self::pointer_hash(readers_schema),
);
if let Some(c) = self.recursion.get(&key).copied() {
Ok(c)
} else {
let c = self.inner_full_match_schemas(writers_schema, readers_schema)?;
self.recursion.insert(key, c);
Ok(c)
}
}
fn pointer_hash(schema: &Schema) -> u64 {
let mut hasher = DefaultHasher::new();
ptr::hash(schema, &mut hasher);
hasher.finish()
}
#[rustfmt::skip]
fn inner_full_match_schemas(
&mut self,
writers_schema: &Schema,
readers_schema: &Schema,
) -> Result<Compatibility, CompatibilityError> {
if let Some(w_name) = writers_schema.name()
&& let Some(r_name) = readers_schema.name()
&& w_name.name() != r_name.name()
{
return Err(CompatibilityError::NameMismatch {
writer_name: w_name.name().into(),
reader_name: r_name.name().into(),
});
}
match (writers_schema, readers_schema) {
(Schema::Ref { name: w_name }, Schema::Ref { name: r_name }) => {
if r_name == w_name {
Ok(Compatibility::Full)
} else {
Err(CompatibilityError::NameMismatch {
writer_name: w_name.fullname(None),
reader_name: r_name.fullname(None),
})
}
}
(Schema::Union(writer), Schema::Union(reader)) => {
let mut any = false;
let mut all = true;
for writer in &writer.schemas {
let mut local_any = false;
all &= reader.schemas.iter().any(|reader| {
match self.full_match_schemas(writer, reader) {
Ok(Compatibility::Full) => {
local_any = true;
true
}
Ok(Compatibility::Partial) => {
local_any = true;
false
}
Err(_) => false,
}
});
any |= local_any;
}
if all {
Ok(Compatibility::Full)
} else if any {
Ok(Compatibility::Partial)
} else {
Err(CompatibilityError::MissingUnionElements)
}
}
(Schema::Union(writer), _) => {
let mut any = false;
let mut all = true;
for writer in &writer.schemas {
match self.full_match_schemas(writer, readers_schema) {
Ok(Compatibility::Full) => any = true,
Ok(Compatibility::Partial) => {
any = true;
all = false;
}
Err(_) => {
all = false;
}
}
}
if all {
Ok(Compatibility::Full)
} else if any {
Ok(Compatibility::Partial)
} else {
Err(CompatibilityError::SchemaMismatchAllUnionElements)
}
}
(_, Schema::Union(reader)) => {
let mut partial = false;
if reader.schemas.iter().any(|reader| {
match self.full_match_schemas(writers_schema, reader) {
Ok(Compatibility::Full) => true,
Ok(Compatibility::Partial) => {
partial = true;
false
}
Err(_) => false,
}
}) {
Ok(Compatibility::Full)
} else if partial {
Ok(Compatibility::Partial)
} else {
Err(CompatibilityError::SchemaMismatchAllUnionElements)
}
}
(Schema::Null, Schema::Null) => Ok(Compatibility::Full),
(Schema::Boolean, Schema::Boolean) => Ok(Compatibility::Full),
(
Schema::Int | Schema::Date | Schema::TimeMillis,
Schema::Int | Schema::Long | Schema::Float | Schema::Double | Schema::Date
| Schema::TimeMillis | Schema::TimeMicros | Schema::TimestampMillis
| Schema::TimestampMicros | Schema::TimestampNanos | Schema::LocalTimestampMillis
| Schema::LocalTimestampMicros | Schema::LocalTimestampNanos,
) => Ok(Compatibility::Full),
(
Schema::Long | Schema::TimeMicros | Schema::TimestampMillis
| Schema::TimestampMicros | Schema::TimestampNanos | Schema::LocalTimestampMillis
| Schema::LocalTimestampMicros | Schema::LocalTimestampNanos,
Schema::Long | Schema::Float | Schema::Double | Schema::TimeMicros
| Schema::TimestampMillis | Schema::TimestampMicros | Schema::TimestampNanos
| Schema::LocalTimestampMillis | Schema::LocalTimestampMicros
| Schema::LocalTimestampNanos,
) => Ok(Compatibility::Full),
(Schema::Float, Schema::Float | Schema::Double) => Ok(Compatibility::Full),
(Schema::Double, Schema::Double) => Ok(Compatibility::Full),
(
Schema::Decimal(DecimalSchema { precision: w_precision, scale: w_scale, .. }),
Schema::Decimal(DecimalSchema { precision: r_precision, scale: r_scale, .. }),
) => {
if r_precision == w_precision && r_scale == w_scale {
Ok(Compatibility::Full)
} else {
Err(CompatibilityError::DecimalMismatch {
r_precision: *r_precision,
r_scale: *r_scale,
w_precision: *w_precision,
w_scale: *w_scale
})
}
}
(
Schema::Bytes | Schema::String | Schema::BigDecimal
| Schema::Uuid(UuidSchema::String | UuidSchema::Bytes)
| Schema::Decimal(DecimalSchema { inner: InnerDecimalSchema::Bytes, .. }),
Schema::Bytes | Schema::String | Schema::BigDecimal
| Schema::Uuid(UuidSchema::String | UuidSchema::Bytes)
| Schema::Decimal(DecimalSchema { inner: InnerDecimalSchema::Bytes, .. }),
) => Ok(Compatibility::Full),
(Schema::Uuid(_), Schema::Uuid(_)) => Ok(Compatibility::Full),
(
Schema::Fixed(w_fixed) | Schema::Uuid(UuidSchema::Fixed(w_fixed))
| Schema::Decimal(DecimalSchema { inner: InnerDecimalSchema::Fixed(w_fixed), .. })
| Schema::Duration(w_fixed),
Schema::Fixed(r_fixed) | Schema::Uuid(UuidSchema::Fixed(r_fixed))
| Schema::Decimal(DecimalSchema { inner: InnerDecimalSchema::Fixed(r_fixed), .. })
| Schema::Duration(r_fixed),
) => {
if r_fixed.size == w_fixed.size {
Ok(Compatibility::Full)
} else {
Err(CompatibilityError::FixedMismatch)
}
}
(
Schema::Array(ArraySchema { items: w_items, .. }),
Schema::Array(ArraySchema { items: r_items, .. }),
) => {
self.full_match_schemas(w_items, r_items)
}
(
Schema::Map(MapSchema { types: w_types, .. }),
Schema::Map(MapSchema { types: r_types, .. }),
) => {
self.full_match_schemas(w_types, r_types)
}
(
Schema::Enum(EnumSchema { symbols: w_symbols, .. }),
Schema::Enum(EnumSchema { symbols: r_symbols, default: r_default, .. }),
) => {
if r_default.is_some() {
Ok(Compatibility::Full)
} else {
let mut any = false;
let mut all = true;
w_symbols.iter().for_each(|s| {
let res = r_symbols.contains(s);
any |= res;
all &= res;
});
if all {
Ok(Compatibility::Full)
} else if any {
Ok(Compatibility::Partial)
} else {
Err(CompatibilityError::MissingSymbols)
}
}
}
(
Schema::Record(RecordSchema { fields: w_fields, .. }),
Schema::Record(RecordSchema { fields: r_fields, .. }),
) => {
let mut compatibility = Compatibility::Full;
for r_field in r_fields {
if let Some(w_field) = once(&r_field.name)
.chain(r_field.aliases.iter())
.find_map(|ra| w_fields.iter().find(|wf| &wf.name == ra))
{
match self.full_match_schemas(&w_field.schema, &r_field.schema) {
Ok(c) => compatibility &= c,
Err(err) => {
return Err(CompatibilityError::FieldTypeMismatch(
r_field.name.clone(),
Box::new(err),
));
}
}
} else if r_field.default.is_none() {
return Err(CompatibilityError::MissingDefaultValue(
r_field.name.clone(),
));
}
}
Ok(compatibility)
}
(_, _) => Err(CompatibilityError::WrongType {
writer_schema_type: format!("{writers_schema:#?}"),
reader_schema_type: format!("{readers_schema:#?}"),
}),
}
}
pub(crate) fn can_read(
&mut self,
writers_schema: &Schema,
readers_schema: &Schema,
) -> Result<Compatibility, CompatibilityError> {
self.full_match_schemas(writers_schema, readers_schema)
}
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use super::*;
use crate::{
Codec, Decimal, Reader, Writer,
schema::{FixedSchema, Name, UuidSchema},
types::{Record, Value},
};
use apache_avro_test_helper::TestResult;
use rstest::*;
fn int_array_schema() -> Schema {
Schema::parse_str(r#"{"type":"array", "items":"int"}"#).unwrap()
}
fn long_array_schema() -> Schema {
Schema::parse_str(r#"{"type":"array", "items":"long"}"#).unwrap()
}
fn string_array_schema() -> Schema {
Schema::parse_str(r#"{"type":"array", "items":"string"}"#).unwrap()
}
fn int_map_schema() -> Schema {
Schema::parse_str(r#"{"type":"map", "values":"int"}"#).unwrap()
}
fn long_map_schema() -> Schema {
Schema::parse_str(r#"{"type":"map", "values":"long"}"#).unwrap()
}
fn string_map_schema() -> Schema {
Schema::parse_str(r#"{"type":"map", "values":"string"}"#).unwrap()
}
fn enum1_ab_schema() -> Schema {
Schema::parse_str(r#"{"type":"enum", "name":"Enum1", "symbols":["A","B"]}"#).unwrap()
}
fn enum1_abc_schema() -> Schema {
Schema::parse_str(r#"{"type":"enum", "name":"Enum1", "symbols":["A","B","C"]}"#).unwrap()
}
fn enum1_bc_schema() -> Schema {
Schema::parse_str(r#"{"type":"enum", "name":"Enum1", "symbols":["B","C"]}"#).unwrap()
}
fn enum2_ab_schema() -> Schema {
Schema::parse_str(r#"{"type":"enum", "name":"Enum2", "symbols":["A","B"]}"#).unwrap()
}
fn empty_record1_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record1", "fields":[]}"#).unwrap()
}
fn empty_record2_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record2", "fields": []}"#).unwrap()
}
fn a_int_record1_schema() -> Schema {
Schema::parse_str(
r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"int"}]}"#,
)
.unwrap()
}
fn a_long_record1_schema() -> Schema {
Schema::parse_str(
r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"long"}]}"#,
)
.unwrap()
}
fn a_int_b_int_record1_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"int"}, {"name":"b", "type":"int"}]}"#).unwrap()
}
fn a_dint_record1_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"int", "default":0}]}"#).unwrap()
}
fn a_int_b_dint_record1_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"int"}, {"name":"b", "type":"int", "default":0}]}"#).unwrap()
}
fn a_dint_b_dint_record1_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"Record1", "fields":[{"name":"a", "type":"int", "default":0}, {"name":"b", "type":"int", "default":0}]}"#).unwrap()
}
fn nested_record() -> Schema {
Schema::parse_str(r#"{"type":"record","name":"parent","fields":[{"name":"attribute","type":{"type":"record","name":"child","fields":[{"name":"id","type":"string"}]}}]}"#).unwrap()
}
fn nested_optional_record() -> Schema {
Schema::parse_str(r#"{"type":"record","name":"parent","fields":[{"name":"attribute","type":["null",{"type":"record","name":"child","fields":[{"name":"id","type":"string"}]}],"default":null}]}"#).unwrap()
}
fn int_list_record_schema() -> Schema {
Schema::parse_str(r#"{"type":"record", "name":"List", "fields": [{"name": "head", "type": "int"},{"name": "tail", "type": {"type": "array", "items": "int"}}]}"#).unwrap()
}
fn long_list_record_schema() -> Schema {
Schema::parse_str(
r#"
{
"type":"record", "name":"List", "fields": [
{"name": "head", "type": "long"},
{"name": "tail", "type": {"type": "array", "items": "long"}}
]}
"#,
)
.unwrap()
}
fn union_schema(schemas: Vec<Schema>) -> Schema {
let schema_string = schemas
.into_iter()
.map(|s| s.canonical_form())
.collect::<Vec<String>>()
.join(",");
Schema::parse_str(&format!("[{schema_string}]")).unwrap()
}
fn empty_union_schema() -> Schema {
union_schema(vec![])
}
fn int_union_schema() -> Schema {
union_schema(vec![Schema::Int])
}
fn long_union_schema() -> Schema {
union_schema(vec![Schema::Long])
}
fn string_union_schema() -> Schema {
union_schema(vec![Schema::String])
}
fn int_string_union_schema() -> Schema {
union_schema(vec![Schema::Int, Schema::String])
}
fn string_int_union_schema() -> Schema {
union_schema(vec![Schema::String, Schema::Int])
}
#[test]
fn test_broken() {
assert_eq!(
Compatibility::Partial,
SchemaCompatibility::can_read(&int_string_union_schema(), &int_union_schema()).unwrap(),
"Only compatible if writer writes an int"
);
}
#[test]
fn test_incompatible_reader_writer_pairs() {
let incompatible_schemas = vec![
(Schema::Null, Schema::Int),
(Schema::Null, Schema::Long),
(Schema::Boolean, Schema::Int),
(Schema::Int, Schema::Null),
(Schema::Int, Schema::Boolean),
(Schema::Int, Schema::Long),
(Schema::Int, Schema::Float),
(Schema::Int, Schema::Double),
(Schema::Long, Schema::Float),
(Schema::Long, Schema::Double),
(Schema::Float, Schema::Double),
(Schema::String, Schema::Boolean),
(Schema::String, Schema::Int),
(Schema::Bytes, Schema::Null),
(Schema::Bytes, Schema::Int),
(Schema::TimeMicros, Schema::Int),
(Schema::TimestampMillis, Schema::Int),
(Schema::TimestampMicros, Schema::Int),
(Schema::TimestampNanos, Schema::Int),
(Schema::LocalTimestampMillis, Schema::Int),
(Schema::LocalTimestampMicros, Schema::Int),
(Schema::LocalTimestampNanos, Schema::Int),
(Schema::Date, Schema::Long),
(Schema::TimeMillis, Schema::Long),
(int_array_schema(), long_array_schema()),
(int_map_schema(), int_array_schema()),
(int_array_schema(), int_map_schema()),
(int_map_schema(), long_map_schema()),
(enum1_ab_schema(), enum1_abc_schema()),
(enum1_bc_schema(), enum1_abc_schema()),
(enum1_ab_schema(), enum2_ab_schema()),
(Schema::Int, enum2_ab_schema()),
(enum2_ab_schema(), Schema::Int),
(int_union_schema(), int_string_union_schema()),
(string_union_schema(), int_string_union_schema()),
(empty_record2_schema(), empty_record1_schema()),
(a_int_record1_schema(), empty_record1_schema()),
(a_int_b_dint_record1_schema(), empty_record1_schema()),
(int_list_record_schema(), long_list_record_schema()),
(nested_record(), nested_optional_record()),
];
assert!(
incompatible_schemas
.iter()
.any(|(reader, writer)| SchemaCompatibility::can_read(writer, reader).is_err())
);
}
#[rstest]
#[case(
r#"{"type": "record", "name": "record_a", "fields": [{"type": "long", "name": "date"}]}"#,
r#"{"type": "record", "name": "record_a", "fields": [{"type": "long", "name": "date", "default": 18181}]}"#
)]
#[case(
r#"{"type": "fixed", "name": "EmployeeId", "size": 16}"#,
r#"{"type": "fixed", "name": "EmployeeId", "size": 16, "default": "u00ffffffffffffx"}"#
)]
#[case(
r#"{"type": "enum", "name":"Enum1", "symbols": ["A","B"]}"#,
r#"{"type": "enum", "name":"Enum1", "symbols": ["A","B", "C"], "default": "C"}"#
)]
#[case(
r#"{"type": "map", "values": "int"}"#,
r#"{"type": "map", "values": "long"}"#
)]
#[case(r#"{"type": "int"}"#, r#"{"type": "int", "logicalType": "date"}"#)]
#[case(
r#"{"type": "int"}"#,
r#"{"type": "int", "logicalType": "time-millis"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "time-micros"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "timestamp-nanos"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "timestamp-millis"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "timestamp-micros"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-millis"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-micros"}"#
)]
#[case(
r#"{"type": "long"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-nanos"}"#
)]
#[case(
r#"{"type": "array", "items": "int"}"#,
r#"{"type": "array", "items": "long"}"#
)]
fn test_avro_3950_match_schemas_ok(
#[case] writer_schema_str: &str,
#[case] reader_schema_str: &str,
) {
let writer_schema = Schema::parse_str(writer_schema_str).unwrap();
let reader_schema = Schema::parse_str(reader_schema_str).unwrap();
assert!(SchemaCompatibility::can_read(&writer_schema, &reader_schema).is_ok());
}
#[rstest]
#[case(
r#"{"type": "record", "name":"record_a", "fields": [{"type": "long", "name": "date"}]}"#,
r#"{"type": "record", "name":"record_b", "fields": [{"type": "long", "name": "date"}]}"#,
CompatibilityError::NameMismatch{writer_name: String::from("record_a"), reader_name: String::from("record_b")}
)]
#[case(
r#"{"type": "fixed", "name": "EmployeeId", "size": 16}"#,
r#"{"type": "fixed", "name": "EmployeeId", "size": 20}"#,
CompatibilityError::FixedMismatch
)]
#[case(
r#"{"type": "enum", "name": "Enum1", "symbols": ["A","B"]}"#,
r#"{"type": "enum", "name": "Enum2", "symbols": ["A","B"]}"#,
CompatibilityError::NameMismatch{writer_name: String::from("Enum1"), reader_name: String::from("Enum2")}
)]
#[case(
r#"{"type":"map", "values": "long"}"#,
r#"{"type":"map", "values": "int"}"#,
CompatibilityError::WrongType { writer_schema_type: "Long".to_string(), reader_schema_type: "Int".to_string() }
)]
#[case(
r#"{"type": "array", "items": "long"}"#,
r#"{"type": "array", "items": "int"}"#,
CompatibilityError::WrongType { writer_schema_type: "Long".to_string(), reader_schema_type: "Int".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "int", "logicalType": "date"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "Date".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "int", "logicalType": "time-millis"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "TimeMillis".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "time-micros"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "TimeMicros".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "timestamp-nanos"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "TimestampNanos".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "timestamp-millis"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "TimestampMillis".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "timestamp-micros"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "TimestampMicros".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-millis"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "LocalTimestampMillis".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-micros"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "LocalTimestampMicros".to_string() }
)]
#[case(
r#"{"type": "string"}"#,
r#"{"type": "long", "logicalType": "local-timestamp-nanos"}"#,
CompatibilityError::WrongType { writer_schema_type: "String".to_string(), reader_schema_type: "LocalTimestampNanos".to_string() }
)]
#[case(
r#"{"type": "record", "name":"record_b", "fields": [{"type": "long", "name": "date"}]}"#,
r#"{"type": "fixed", "name": "EmployeeId", "size": 16}"#,
CompatibilityError::NameMismatch { writer_name: "record_b".to_string(), reader_name: "EmployeeId".to_string() }
)]
fn test_avro_3950_match_schemas_error(
#[case] writer_schema_str: &str,
#[case] reader_schema_str: &str,
#[case] expected_error: CompatibilityError,
) {
let writer_schema = Schema::parse_str(writer_schema_str).unwrap();
let reader_schema = Schema::parse_str(reader_schema_str).unwrap();
assert_eq!(
expected_error,
SchemaCompatibility::can_read(&writer_schema, &reader_schema).unwrap_err()
);
}
#[test]
fn test_compatible_reader_writer_pairs() {
let uuid_fixed = FixedSchema {
name: Name::new("uuid_fixed").unwrap(),
aliases: None,
doc: None,
size: 16,
attributes: Default::default(),
};
let compatible_schemas = vec![
(Schema::Null, Schema::Null),
(Schema::Long, Schema::Int),
(Schema::Float, Schema::Int),
(Schema::Float, Schema::Long),
(Schema::Double, Schema::Long),
(Schema::Double, Schema::Int),
(Schema::Double, Schema::Float),
(Schema::String, Schema::Bytes),
(Schema::Bytes, Schema::String),
(
Schema::Uuid(UuidSchema::String),
Schema::Uuid(UuidSchema::String),
),
(Schema::Uuid(UuidSchema::String), Schema::String),
(Schema::String, Schema::Uuid(UuidSchema::String)),
(
Schema::Uuid(UuidSchema::Bytes),
Schema::Uuid(UuidSchema::Bytes),
),
(Schema::Uuid(UuidSchema::Bytes), Schema::Bytes),
(Schema::Bytes, Schema::Uuid(UuidSchema::Bytes)),
(
Schema::Uuid(UuidSchema::Fixed(uuid_fixed.clone())),
Schema::Uuid(UuidSchema::Fixed(uuid_fixed.clone())),
),
(
Schema::Uuid(UuidSchema::Fixed(uuid_fixed.clone())),
Schema::Fixed(uuid_fixed.clone()),
),
(
Schema::Fixed(uuid_fixed.clone()),
Schema::Uuid(UuidSchema::Fixed(uuid_fixed.clone())),
),
(Schema::Date, Schema::Int),
(Schema::TimeMillis, Schema::Int),
(Schema::TimeMicros, Schema::Long),
(Schema::TimestampMillis, Schema::Long),
(Schema::TimestampMicros, Schema::Long),
(Schema::TimestampNanos, Schema::Long),
(Schema::LocalTimestampMillis, Schema::Long),
(Schema::LocalTimestampMicros, Schema::Long),
(Schema::LocalTimestampNanos, Schema::Long),
(Schema::Int, Schema::Date),
(Schema::Int, Schema::TimeMillis),
(Schema::Long, Schema::TimeMicros),
(Schema::Long, Schema::TimestampMillis),
(Schema::Long, Schema::TimestampMicros),
(Schema::Long, Schema::TimestampNanos),
(Schema::Long, Schema::LocalTimestampMillis),
(Schema::Long, Schema::LocalTimestampMicros),
(Schema::Long, Schema::LocalTimestampNanos),
(int_array_schema(), int_array_schema()),
(long_array_schema(), int_array_schema()),
(int_map_schema(), int_map_schema()),
(long_map_schema(), int_map_schema()),
(enum1_ab_schema(), enum1_ab_schema()),
(enum1_abc_schema(), enum1_ab_schema()),
(empty_union_schema(), empty_union_schema()),
(int_union_schema(), int_union_schema()),
(int_string_union_schema(), string_int_union_schema()),
(int_union_schema(), empty_union_schema()),
(long_union_schema(), int_union_schema()),
(int_union_schema(), Schema::Int),
(Schema::Int, int_union_schema()),
(empty_record1_schema(), empty_record1_schema()),
(empty_record1_schema(), a_int_record1_schema()),
(a_int_record1_schema(), a_int_record1_schema()),
(a_dint_record1_schema(), a_int_record1_schema()),
(a_dint_record1_schema(), a_dint_record1_schema()),
(a_int_record1_schema(), a_dint_record1_schema()),
(a_long_record1_schema(), a_int_record1_schema()),
(a_int_record1_schema(), a_int_b_int_record1_schema()),
(a_dint_record1_schema(), a_int_b_int_record1_schema()),
(a_int_b_dint_record1_schema(), a_int_record1_schema()),
(a_dint_b_dint_record1_schema(), empty_record1_schema()),
(a_dint_b_dint_record1_schema(), a_int_record1_schema()),
(a_int_b_int_record1_schema(), a_dint_b_dint_record1_schema()),
(int_list_record_schema(), int_list_record_schema()),
(long_list_record_schema(), long_list_record_schema()),
(long_list_record_schema(), int_list_record_schema()),
(nested_optional_record(), nested_record()),
];
for (reader, writer) in compatible_schemas {
SchemaCompatibility::can_read(&writer, &reader).unwrap();
}
}
fn writer_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield1", "type":"int"},
{"name":"oldfield2", "type":"string"}
]}
"#,
)
.unwrap()
}
#[test]
fn test_missing_field() -> TestResult {
let reader_schema = Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield1", "type":"int"}
]}
"#,
)?;
assert!(SchemaCompatibility::can_read(&writer_schema(), &reader_schema,).is_ok());
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("oldfield2")),
SchemaCompatibility::can_read(&reader_schema, &writer_schema()).unwrap_err()
);
Ok(())
}
#[test]
fn test_missing_second_field() -> TestResult {
let reader_schema = Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield2", "type":"string"}
]}
"#,
)?;
assert!(SchemaCompatibility::can_read(&writer_schema(), &reader_schema).is_ok());
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("oldfield1")),
SchemaCompatibility::can_read(&reader_schema, &writer_schema()).unwrap_err()
);
Ok(())
}
#[test]
fn test_all_fields() -> TestResult {
let reader_schema = Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield1", "type":"int"},
{"name":"oldfield2", "type":"string"}
]}
"#,
)?;
assert!(SchemaCompatibility::can_read(&writer_schema(), &reader_schema).is_ok());
assert!(SchemaCompatibility::can_read(&reader_schema, &writer_schema()).is_ok());
Ok(())
}
#[test]
fn test_new_field_with_default() -> TestResult {
let reader_schema = Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield1", "type":"int"},
{"name":"newfield1", "type":"int", "default":42}
]}
"#,
)?;
assert!(SchemaCompatibility::can_read(&writer_schema(), &reader_schema).is_ok());
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("oldfield2")),
SchemaCompatibility::can_read(&reader_schema, &writer_schema()).unwrap_err()
);
Ok(())
}
#[test]
fn test_new_field() -> TestResult {
let reader_schema = Schema::parse_str(
r#"
{"type":"record", "name":"Record", "fields":[
{"name":"oldfield1", "type":"int"},
{"name":"newfield1", "type":"int"}
]}
"#,
)?;
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("newfield1")),
SchemaCompatibility::can_read(&writer_schema(), &reader_schema).unwrap_err()
);
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("oldfield2")),
SchemaCompatibility::can_read(&reader_schema, &writer_schema()).unwrap_err()
);
Ok(())
}
#[test]
fn test_array_writer_schema() {
let valid_reader = string_array_schema();
let invalid_reader = string_map_schema();
assert_eq!(
Compatibility::Full,
SchemaCompatibility::can_read(&string_array_schema(), &valid_reader).unwrap()
);
assert!(matches!(
SchemaCompatibility::can_read(&string_array_schema(), &invalid_reader),
Err(CompatibilityError::WrongType { .. }),
));
}
#[test]
fn test_primitive_writer_schema() {
let valid_reader = Schema::String;
assert!(SchemaCompatibility::can_read(&Schema::String, &valid_reader).is_ok());
assert_eq!(
CompatibilityError::WrongType {
writer_schema_type: "Int".to_string(),
reader_schema_type: "String".to_string()
},
SchemaCompatibility::can_read(&Schema::Int, &Schema::String).unwrap_err()
);
}
#[test]
fn test_union_reader_writer_subset_incompatibility() {
let union_writer = union_schema(vec![Schema::Int, Schema::String]);
let union_reader = union_schema(vec![Schema::String]);
assert_eq!(
Compatibility::Partial,
SchemaCompatibility::can_read(&union_writer, &union_reader).unwrap()
);
assert_eq!(
Compatibility::Full,
SchemaCompatibility::can_read(&union_reader, &union_writer).unwrap()
);
}
#[test]
fn test_incompatible_record_field() -> TestResult {
let string_schema = Schema::parse_str(
r#"
{"type":"record", "name":"MyRecord", "namespace":"ns", "fields": [
{"name":"field1", "type":"string"}
]}
"#,
)?;
let int_schema = Schema::parse_str(
r#"
{"type":"record", "name":"MyRecord", "namespace":"ns", "fields": [
{"name":"field1", "type":"int"}
]}
"#,
)?;
assert_eq!(
CompatibilityError::FieldTypeMismatch(
"field1".to_owned(),
Box::new(CompatibilityError::WrongType {
writer_schema_type: "String".to_string(),
reader_schema_type: "Int".to_string()
})
),
SchemaCompatibility::can_read(&string_schema, &int_schema).unwrap_err()
);
Ok(())
}
#[test]
fn test_enum_symbols() -> TestResult {
let enum_schema1 = Schema::parse_str(
r#"
{"type":"enum", "name":"MyEnum", "symbols":["A","B"]}
"#,
)?;
let enum_schema2 =
Schema::parse_str(r#"{"type":"enum", "name":"MyEnum", "symbols":["A","B","C"]}"#)?;
assert_eq!(
Compatibility::Partial,
SchemaCompatibility::can_read(&enum_schema2, &enum_schema1)?
);
assert_eq!(
Compatibility::Full,
SchemaCompatibility::can_read(&enum_schema1, &enum_schema2)?
);
Ok(())
}
fn point_2d_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Point2D", "fields":[
{"name":"x", "type":"double"},
{"name":"y", "type":"double"}
]}
"#,
)
.unwrap()
}
fn point_2d_fullname_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Point", "namespace":"written", "fields":[
{"name":"x", "type":"double"},
{"name":"y", "type":"double"}
]}
"#,
)
.unwrap()
}
fn point_3d_no_default_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Point", "fields":[
{"name":"x", "type":"double"},
{"name":"y", "type":"double"},
{"name":"z", "type":"double"}
]}
"#,
)
.unwrap()
}
fn point_3d_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Point3D", "fields":[
{"name":"x", "type":"double"},
{"name":"y", "type":"double"},
{"name":"z", "type":"double", "default": 0.0}
]}
"#,
)
.unwrap()
}
fn point_3d_match_name_schema() -> Schema {
Schema::parse_str(
r#"
{"type":"record", "name":"Point", "fields":[
{"name":"x", "type":"double"},
{"name":"y", "type":"double"},
{"name":"z", "type":"double", "default": 0.0}
]}
"#,
)
.unwrap()
}
#[test]
fn test_union_resolution_no_structure_match() {
let read_schema = union_schema(vec![Schema::Null, point_3d_no_default_schema()]);
assert_eq!(
CompatibilityError::SchemaMismatchAllUnionElements,
SchemaCompatibility::can_read(&point_2d_fullname_schema(), &read_schema).unwrap_err()
);
}
#[test]
fn test_union_resolution_first_structure_match_2d() {
let read_schema = union_schema(vec![
Schema::Null,
point_3d_no_default_schema(),
point_2d_schema(),
point_3d_schema(),
]);
assert_eq!(
CompatibilityError::SchemaMismatchAllUnionElements,
SchemaCompatibility::can_read(&point_2d_fullname_schema(), &read_schema).unwrap_err()
);
}
#[test]
fn test_union_resolution_first_structure_match_3d() {
let read_schema = union_schema(vec![
Schema::Null,
point_3d_no_default_schema(),
point_3d_schema(),
point_2d_schema(),
]);
assert_eq!(
CompatibilityError::SchemaMismatchAllUnionElements,
SchemaCompatibility::can_read(&point_2d_fullname_schema(), &read_schema).unwrap_err()
);
}
#[test]
fn test_union_resolution_named_structure_match() {
let read_schema = union_schema(vec![
Schema::Null,
point_2d_schema(),
point_3d_match_name_schema(),
point_3d_schema(),
]);
assert_eq!(
CompatibilityError::SchemaMismatchAllUnionElements,
SchemaCompatibility::can_read(&point_2d_fullname_schema(), &read_schema).unwrap_err()
);
}
#[test]
fn test_union_resolution_full_name_match() {
let read_schema = union_schema(vec![
Schema::Null,
point_2d_schema(),
point_3d_match_name_schema(),
point_3d_schema(),
point_2d_fullname_schema(),
]);
assert!(SchemaCompatibility::can_read(&point_2d_fullname_schema(), &read_schema).is_ok());
}
#[test]
fn test_avro_3772_enum_default() -> TestResult {
let writer_raw_schema = r#"
{
"type": "record",
"name": "test",
"fields": [
{"name": "a", "type": "long", "default": 42},
{"name": "b", "type": "string"},
{
"name": "c",
"type": {
"type": "enum",
"name": "suit",
"symbols": ["diamonds", "spades", "clubs", "hearts"],
"default": "spades"
}
}
]
}
"#;
let reader_raw_schema = r#"
{
"type": "record",
"name": "test",
"fields": [
{"name": "a", "type": "long", "default": 42},
{"name": "b", "type": "string"},
{
"name": "c",
"type": {
"type": "enum",
"name": "suit",
"symbols": ["diamonds", "spades", "ninja", "hearts"],
"default": "spades"
}
}
]
}
"#;
let writer_schema = Schema::parse_str(writer_raw_schema)?;
let reader_schema = Schema::parse_str(reader_raw_schema)?;
let mut writer = Writer::with_codec(&writer_schema, Vec::new(), Codec::Null)?;
let mut record = Record::new(writer.schema()).unwrap();
record.put("a", 27i64);
record.put("b", "foo");
record.put("c", "clubs");
writer.append_value(record).unwrap();
let input = writer.into_inner()?;
let mut reader = Reader::builder(&input[..])
.reader_schema(&reader_schema)
.build()?;
assert_eq!(
reader.next().unwrap().unwrap(),
Value::Record(vec![
("a".to_string(), Value::Long(27)),
("b".to_string(), Value::String("foo".to_string())),
("c".to_string(), Value::Enum(1, "spades".to_string())),
])
);
assert!(reader.next().is_none());
Ok(())
}
#[test]
fn test_avro_3772_enum_default_less_symbols() -> TestResult {
let writer_raw_schema = r#"
{
"type": "record",
"name": "test",
"fields": [
{"name": "a", "type": "long", "default": 42},
{"name": "b", "type": "string"},
{
"name": "c",
"type": {
"type": "enum",
"name": "suit",
"symbols": ["diamonds", "spades", "clubs", "hearts"],
"default": "spades"
}
}
]
}
"#;
let reader_raw_schema = r#"
{
"type": "record",
"name": "test",
"fields": [
{"name": "a", "type": "long", "default": 42},
{"name": "b", "type": "string"},
{
"name": "c",
"type": {
"type": "enum",
"name": "suit",
"symbols": ["hearts", "spades"],
"default": "spades"
}
}
]
}
"#;
let writer_schema = Schema::parse_str(writer_raw_schema)?;
let reader_schema = Schema::parse_str(reader_raw_schema)?;
let mut writer = Writer::with_codec(&writer_schema, Vec::new(), Codec::Null)?;
let mut record = Record::new(writer.schema()).unwrap();
record.put("a", 27i64);
record.put("b", "foo");
record.put("c", "hearts");
writer.append_value(record).unwrap();
let input = writer.into_inner()?;
let mut reader = Reader::builder(&input[..])
.reader_schema(&reader_schema)
.build()?;
assert_eq!(
reader.next().unwrap().unwrap(),
Value::Record(vec![
("a".to_string(), Value::Long(27)),
("b".to_string(), Value::String("foo".to_string())),
("c".to_string(), Value::Enum(0, "hearts".to_string())),
])
);
assert!(reader.next().is_none());
Ok(())
}
#[test]
fn avro_3894_take_aliases_into_account_when_serializing_for_schema_compatibility() -> TestResult
{
let schema_v1 = Schema::parse_str(
r#"
{
"type": "record",
"name": "Conference",
"namespace": "advdaba",
"fields": [
{"type": "string", "name": "name"},
{"type": "long", "name": "date"}
]
}"#,
)?;
let schema_v2 = Schema::parse_str(
r#"
{
"type": "record",
"name": "Conference",
"namespace": "advdaba",
"fields": [
{"type": "string", "name": "name"},
{"type": "long", "name": "date", "aliases" : [ "time" ]}
]
}"#,
)?;
assert!(SchemaCompatibility::mutual_read(&schema_v1, &schema_v2).is_ok());
Ok(())
}
#[test]
fn avro_3917_take_aliases_into_account_for_schema_compatibility() -> TestResult {
let schema_v1 = Schema::parse_str(
r#"
{
"type": "record",
"name": "Conference",
"namespace": "advdaba",
"fields": [
{"type": "string", "name": "name"},
{"type": "long", "name": "date", "aliases" : [ "time" ]}
]
}"#,
)?;
let schema_v2 = Schema::parse_str(
r#"
{
"type": "record",
"name": "Conference",
"namespace": "advdaba",
"fields": [
{"type": "string", "name": "name"},
{"type": "long", "name": "time"}
]
}"#,
)?;
assert_eq!(
Compatibility::Full,
SchemaCompatibility::can_read(&schema_v2, &schema_v1)?
);
assert_eq!(
CompatibilityError::MissingDefaultValue(String::from("time")),
SchemaCompatibility::can_read(&schema_v1, &schema_v2).unwrap_err()
);
Ok(())
}
#[test]
fn test_avro_3898_record_schemas_match_by_unqualified_name() -> TestResult {
let schemas = [
(
Schema::parse_str(
r#"{
"type": "record",
"name": "Statistics",
"fields": [
{ "name": "success", "type": "int" },
{ "name": "fail", "type": "int" },
{ "name": "time", "type": "string" },
{ "name": "max", "type": "int", "default": 0 }
]
}"#,
)?,
Schema::parse_str(
r#"{
"type": "record",
"name": "Statistics",
"namespace": "my.namespace",
"fields": [
{ "name": "success", "type": "int" },
{ "name": "fail", "type": "int" },
{ "name": "time", "type": "string" },
{ "name": "average", "type": "int", "default": 0}
]
}"#,
)?,
),
(
Schema::parse_str(
r#"{
"type": "enum",
"name": "Suit",
"symbols": ["diamonds", "spades", "clubs"]
}"#,
)?,
Schema::parse_str(
r#"{
"type": "enum",
"name": "Suit",
"namespace": "my.namespace",
"symbols": ["diamonds", "spades", "clubs", "hearts"]
}"#,
)?,
),
(
Schema::parse_str(
r#"{
"type": "fixed",
"name": "EmployeeId",
"size": 16
}"#,
)?,
Schema::parse_str(
r#"{
"type": "fixed",
"name": "EmployeeId",
"namespace": "my.namespace",
"size": 16
}"#,
)?,
),
];
for (schema_1, schema_2) in schemas {
assert!(SchemaCompatibility::can_read(&schema_1, &schema_2).is_ok());
}
Ok(())
}
#[test]
fn test_can_read_compatibility_errors() -> TestResult {
let schemas = [
(
Schema::parse_str(
r#"{
"type": "record",
"name": "StatisticsMap",
"fields": [
{"name": "average", "type": "int", "default": 0},
{"name": "success", "type": {"type": "map", "values": "int"}}
]
}"#,
)?,
Schema::parse_str(
r#"{
"type": "record",
"name": "StatisticsMap",
"fields": [
{"name": "average", "type": "int", "default": 0},
{"name": "success", "type": ["null", {"type": "map", "values": "int"}], "default": null}
]
}"#,
)?,
),
(
Schema::parse_str(
r#"{
"type": "record",
"name": "StatisticsArray",
"fields": [
{"name": "max_values", "type": {"type": "array", "items": "int"}}
]
}"#,
)?,
Schema::parse_str(
r#"{
"type": "record",
"name": "StatisticsArray",
"fields": [
{"name": "max_values", "type": ["null", {"type": "array", "items": "int"}], "default": null}
]
}"#,
)?,
),
];
for (schema_1, schema_2) in schemas {
assert_eq!(
Compatibility::Full,
SchemaCompatibility::can_read(&schema_1, &schema_2).unwrap()
);
assert_eq!(
Compatibility::Partial,
SchemaCompatibility::can_read(&schema_2, &schema_1).unwrap()
);
}
Ok(())
}
#[test]
fn avro_3974_can_read_schema_references() -> TestResult {
let schema_strs = vec![
r#"{
"type": "record",
"name": "Child",
"namespace": "avro",
"fields": [
{
"name": "val",
"type": "int"
}
]
}
"#,
r#"{
"type": "record",
"name": "Parent",
"namespace": "avro",
"fields": [
{
"name": "child",
"type": "avro.Child"
}
]
}
"#,
];
let schemas = Schema::parse_list(schema_strs).unwrap();
SchemaCompatibility::can_read(&schemas[1], &schemas[1])?;
Ok(())
}
#[test]
fn duration_and_fixed_of_different_size() -> TestResult {
let schema_strs = vec![
r#"{
"type": "fixed",
"name": "Fixed25",
"size": 25
}
"#,
r#"{
"type": "fixed",
"logicalType": "duration",
"name": "Duration",
"size": 12
}
"#,
];
let schemas = Schema::parse_list(schema_strs).unwrap();
assert!(SchemaCompatibility::can_read(&schemas[0], &schemas[1]).is_err());
assert!(SchemaCompatibility::can_read(&schemas[1], &schemas[0]).is_err());
SchemaCompatibility::can_read(&schemas[1], &schemas[1])?;
SchemaCompatibility::can_read(&schemas[0], &schemas[0])?;
Ok(())
}
#[test]
fn avro_rs_342_decimal_fixed_and_bytes() -> TestResult {
let bytes = Schema::Decimal(DecimalSchema {
precision: 20,
scale: 0,
inner: InnerDecimalSchema::Bytes,
});
let fixed = Schema::Decimal(DecimalSchema {
precision: 20,
scale: 0,
inner: InnerDecimalSchema::Fixed(FixedSchema {
name: Name::new("DecimalFixed")?,
aliases: None,
doc: None,
size: 20,
attributes: BTreeMap::default(),
}),
});
assert_eq!(
Compatibility::Full,
SchemaCompatibility::mutual_read(&bytes, &fixed)?
);
let value = Value::Decimal(Decimal::from(vec![1; 10]));
let fixed_value = value.clone().resolve(&fixed)?;
let bytes_value = value.resolve(&bytes)?;
assert_eq!(fixed_value, bytes_value);
Ok(())
}
}