1use std::borrow::Cow;
27use std::cell::RefCell;
28use std::collections::hash_map::Entry;
29use std::collections::{HashMap, HashSet};
30use std::fmt;
31use std::rc::Rc;
32use std::str::FromStr;
33
34use digest::Digest;
35use itertools::Itertools;
36use regex::Regex;
37use serde::{
38 ser::{SerializeMap, SerializeSeq},
39 Serialize, Serializer,
40};
41use serde_json::{self, Map, Value};
42use tracing::{debug, warn};
43use types::{DecimalValue, Value as AvroValue};
44
45use crate::error::Error as AvroError;
46use crate::reader::SchemaResolver;
47use crate::types;
48use crate::types::AvroMap;
49use crate::util::MapHelper;
50
51pub fn resolve_schemas(
52 writer_schema: &Schema,
53 reader_schema: &Schema,
54) -> Result<Schema, AvroError> {
55 let r_indices = reader_schema.indices.clone();
56 let (reader_to_writer_names, writer_to_reader_names): (HashMap<_, _>, HashMap<_, _>) =
57 writer_schema
58 .indices
59 .iter()
60 .flat_map(|(name, widx)| {
61 r_indices
62 .get(name)
63 .map(|ridx| ((*ridx, *widx), (*widx, *ridx)))
64 })
65 .unzip();
66 let reader_fullnames = reader_schema
67 .indices
68 .iter()
69 .map(|(f, i)| (*i, f))
70 .collect::<HashMap<_, _>>();
71 let mut resolver = SchemaResolver {
72 named: Default::default(),
73 indices: Default::default(),
74 human_readable_field_path: Vec::new(),
75 current_human_readable_path_start: 0,
76 writer_to_reader_names,
77 reader_to_writer_names,
78 reader_to_resolved_names: Default::default(),
79 reader_fullnames,
80 reader_schema,
81 };
82 let writer_node = writer_schema.top_node_or_named();
83 let reader_node = reader_schema.top_node_or_named();
84 let inner = resolver.resolve(writer_node, reader_node)?;
85 let sch = Schema {
86 named: resolver.named.into_iter().map(Option::unwrap).collect(),
87 indices: resolver.indices,
88 top: inner,
89 };
90 Ok(sch)
91}
92
93#[derive(Clone, Debug, Eq, PartialEq)]
95pub struct ParseSchemaError(String);
96
97impl ParseSchemaError {
98 pub fn new<S>(msg: S) -> ParseSchemaError
99 where
100 S: Into<String>,
101 {
102 ParseSchemaError(msg.into())
103 }
104}
105
106impl fmt::Display for ParseSchemaError {
107 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
108 self.0.fmt(f)
109 }
110}
111
112impl std::error::Error for ParseSchemaError {}
113
114#[derive(Debug)]
118pub struct SchemaFingerprint {
119 pub bytes: Vec<u8>,
120}
121
122impl fmt::Display for SchemaFingerprint {
123 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
124 write!(
125 f,
126 "{}",
127 self.bytes
128 .iter()
129 .map(|byte| format!("{:02x}", byte))
130 .collect::<Vec<String>>()
131 .join("")
132 )
133 }
134}
135
136#[derive(Clone, Debug, PartialEq)]
137pub enum SchemaPieceOrNamed {
138 Piece(SchemaPiece),
139 Named(usize),
140}
141impl SchemaPieceOrNamed {
142 pub fn get_human_name(&self, root: &Schema) -> String {
143 match self {
144 Self::Piece(piece) => format!("{:?}", piece),
145 Self::Named(idx) => format!("{}", root.lookup(*idx).name),
146 }
147 }
148 #[inline(always)]
149 pub fn get_piece_and_name<'a>(
150 &'a self,
151 root: &'a Schema,
152 ) -> (&'a SchemaPiece, Option<&'a FullName>) {
153 self.as_ref().get_piece_and_name(root)
154 }
155
156 #[inline(always)]
157 pub fn as_ref(&self) -> SchemaPieceRefOrNamed {
158 match self {
159 SchemaPieceOrNamed::Piece(piece) => SchemaPieceRefOrNamed::Piece(piece),
160 SchemaPieceOrNamed::Named(index) => SchemaPieceRefOrNamed::Named(*index),
161 }
162 }
163}
164
165impl From<SchemaPiece> for SchemaPieceOrNamed {
166 #[inline(always)]
167 fn from(piece: SchemaPiece) -> Self {
168 Self::Piece(piece)
169 }
170}
171
172#[derive(Clone, Debug, PartialEq)]
173pub enum SchemaPiece {
174 Null,
176 Boolean,
178 Int,
180 Long,
182 Float,
184 Double,
186 Date,
188 TimestampMilli,
192 TimestampMicro,
196 Decimal {
202 precision: usize,
203 scale: usize,
204 fixed_size: Option<usize>,
205 },
206 Bytes,
209 String,
212 Json,
214 Uuid,
216 Array(Box<SchemaPieceOrNamed>),
219 Map(Box<SchemaPieceOrNamed>),
223 Union(UnionSchema),
225 ResolveIntTsMilli,
228 ResolveIntTsMicro,
231 ResolveDateTimestamp,
234 ResolveIntLong,
236 ResolveIntFloat,
238 ResolveIntDouble,
240 ResolveLongFloat,
242 ResolveLongDouble,
244 ResolveFloatDouble,
246 ResolveConcreteUnion {
249 index: usize,
251 inner: Box<SchemaPieceOrNamed>,
253 n_reader_variants: usize,
254 reader_null_variant: Option<usize>,
255 },
256 ResolveUnionUnion {
259 permutation: Vec<Result<(usize, SchemaPieceOrNamed), AvroError>>,
265 n_reader_variants: usize,
266 reader_null_variant: Option<usize>,
267 },
268 ResolveUnionConcrete {
270 index: usize,
271 inner: Box<SchemaPieceOrNamed>,
272 },
273 Record {
278 doc: Documentation,
279 fields: Vec<RecordField>,
280 lookup: HashMap<String, usize>,
281 },
282 Enum {
284 doc: Documentation,
285 symbols: Vec<String>,
286 default_idx: Option<usize>,
292 },
293 Fixed { size: usize },
295 ResolveRecord {
298 defaults: Vec<ResolvedDefaultValueField>,
301 fields: Vec<ResolvedRecordField>,
305 n_reader_fields: usize,
307 },
308 ResolveEnum {
311 doc: Documentation,
312 symbols: Vec<Result<(usize, String), String>>,
315 default: Option<(usize, String)>,
317 },
318}
319
320impl SchemaPiece {
321 pub fn is_underlying_int(&self) -> bool {
323 matches!(self, SchemaPiece::Int | SchemaPiece::Date)
324 }
325 pub fn is_underlying_long(&self) -> bool {
327 matches!(
328 self,
329 SchemaPiece::Long | SchemaPiece::TimestampMilli | SchemaPiece::TimestampMicro
330 )
331 }
332}
333
334#[derive(Clone, PartialEq)]
338pub struct Schema {
339 pub(crate) named: Vec<NamedSchemaPiece>,
340 pub(crate) indices: HashMap<FullName, usize>,
341 pub top: SchemaPieceOrNamed,
342}
343
344impl ToString for Schema {
345 fn to_string(&self) -> String {
346 let json = serde_json::to_value(self).unwrap();
347 json.to_string()
348 }
349}
350
351impl std::fmt::Debug for Schema {
352 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
353 if f.alternate() {
354 f.write_str(
355 &serde_json::to_string_pretty(self)
356 .unwrap_or_else(|e| format!("failed to serialize: {}", e)),
357 )
358 } else {
359 f.write_str(
360 &serde_json::to_string(self)
361 .unwrap_or_else(|e| format!("failed to serialize: {}", e)),
362 )
363 }
364 }
365}
366
367impl Schema {
368 pub fn top_node(&self) -> SchemaNode {
369 let (inner, name) = self.top.get_piece_and_name(self);
370 SchemaNode {
371 root: self,
372 inner,
373 name,
374 }
375 }
376 pub fn top_node_or_named(&self) -> SchemaNodeOrNamed {
377 SchemaNodeOrNamed {
378 root: self,
379 inner: self.top.as_ref(),
380 }
381 }
382 pub fn lookup(&self, idx: usize) -> &NamedSchemaPiece {
383 &self.named[idx]
384 }
385 pub fn try_lookup_name(&self, name: &FullName) -> Option<&NamedSchemaPiece> {
386 self.indices.get(name).map(|&idx| &self.named[idx])
387 }
388}
389
390#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
398pub enum SchemaKind {
399 Null,
401 Boolean,
402 Int,
403 Long,
404 Float,
405 Double,
406 Bytes,
408 String,
409 Array,
410 Map,
411 Union,
412 Record,
413 Enum,
414 Fixed,
415 Unknown,
418}
419
420impl SchemaKind {
421 pub fn name(self) -> &'static str {
422 match self {
423 SchemaKind::Null => "null",
424 SchemaKind::Boolean => "boolean",
425 SchemaKind::Int => "int",
426 SchemaKind::Long => "long",
427 SchemaKind::Float => "float",
428 SchemaKind::Double => "double",
429 SchemaKind::Bytes => "bytes",
430 SchemaKind::String => "string",
431 SchemaKind::Array => "array",
432 SchemaKind::Map => "map",
433 SchemaKind::Union => "union",
434 SchemaKind::Record => "record",
435 SchemaKind::Enum => "enum",
436 SchemaKind::Fixed => "fixed",
437 SchemaKind::Unknown => "unknown",
438 }
439 }
440}
441
442impl<'a> From<&'a SchemaPiece> for SchemaKind {
443 #[inline(always)]
444 fn from(piece: &'a SchemaPiece) -> SchemaKind {
445 match piece {
446 SchemaPiece::Null => SchemaKind::Null,
447 SchemaPiece::Boolean => SchemaKind::Boolean,
448 SchemaPiece::Int => SchemaKind::Int,
449 SchemaPiece::Long => SchemaKind::Long,
450 SchemaPiece::Float => SchemaKind::Float,
451 SchemaPiece::Double => SchemaKind::Double,
452 SchemaPiece::Date => SchemaKind::Int,
453 SchemaPiece::TimestampMilli
454 | SchemaPiece::TimestampMicro
455 | SchemaPiece::ResolveIntTsMilli
456 | SchemaPiece::ResolveDateTimestamp
457 | SchemaPiece::ResolveIntTsMicro => SchemaKind::Long,
458 SchemaPiece::Decimal {
459 fixed_size: None, ..
460 } => SchemaKind::Bytes,
461 SchemaPiece::Decimal {
462 fixed_size: Some(_),
463 ..
464 } => SchemaKind::Fixed,
465 SchemaPiece::Bytes => SchemaKind::Bytes,
466 SchemaPiece::String => SchemaKind::String,
467 SchemaPiece::Array(_) => SchemaKind::Array,
468 SchemaPiece::Map(_) => SchemaKind::Map,
469 SchemaPiece::Union(_) => SchemaKind::Union,
470 SchemaPiece::ResolveUnionUnion { .. } => SchemaKind::Union,
471 SchemaPiece::ResolveIntLong => SchemaKind::Long,
472 SchemaPiece::ResolveIntFloat => SchemaKind::Float,
473 SchemaPiece::ResolveIntDouble => SchemaKind::Double,
474 SchemaPiece::ResolveLongFloat => SchemaKind::Float,
475 SchemaPiece::ResolveLongDouble => SchemaKind::Double,
476 SchemaPiece::ResolveFloatDouble => SchemaKind::Double,
477 SchemaPiece::ResolveConcreteUnion { .. } => SchemaKind::Union,
478 SchemaPiece::ResolveUnionConcrete { inner: _, .. } => SchemaKind::Unknown,
479 SchemaPiece::Record { .. } => SchemaKind::Record,
480 SchemaPiece::Enum { .. } => SchemaKind::Enum,
481 SchemaPiece::Fixed { .. } => SchemaKind::Fixed,
482 SchemaPiece::ResolveRecord { .. } => SchemaKind::Record,
483 SchemaPiece::ResolveEnum { .. } => SchemaKind::Enum,
484 SchemaPiece::Json => SchemaKind::String,
485 SchemaPiece::Uuid => SchemaKind::String,
486 }
487 }
488}
489
490impl<'a> From<SchemaNode<'a>> for SchemaKind {
491 #[inline(always)]
492 fn from(schema: SchemaNode<'a>) -> SchemaKind {
493 SchemaKind::from(schema.inner)
494 }
495}
496
497impl<'a> From<&'a Schema> for SchemaKind {
498 #[inline(always)]
499 fn from(schema: &'a Schema) -> SchemaKind {
500 Self::from(schema.top_node())
501 }
502}
503
504#[derive(Clone, Debug, PartialEq)]
515pub struct Name {
516 pub name: String,
517 pub namespace: Option<String>,
518 pub aliases: Option<Vec<String>>,
519}
520
521#[derive(Clone, Debug, Hash, PartialEq, Eq)]
522pub struct FullName {
523 name: String,
524 namespace: String,
525}
526
527impl FullName {
528 pub fn from_parts(name: &str, namespace: Option<&str>, default_namespace: &str) -> FullName {
529 if let Some(ns) = namespace {
530 FullName {
531 name: name.to_owned(),
532 namespace: ns.to_owned(),
533 }
534 } else {
535 let mut split = name.rsplitn(2, '.');
536 let name = split.next().unwrap();
537 let namespace = split.next().unwrap_or(default_namespace);
538
539 FullName {
540 name: name.into(),
541 namespace: namespace.into(),
542 }
543 }
544 }
545 pub fn base_name(&self) -> &str {
546 &self.name
547 }
548 pub fn human_name(&self) -> String {
549 if self.namespace.is_empty() {
550 return self.name.clone();
551 }
552 return format!("{}.{}", self.namespace, self.name);
553 }
554}
555
556impl fmt::Display for FullName {
557 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
558 write!(f, "{}.{}", self.namespace, self.name)
559 }
560}
561
562pub type Documentation = Option<String>;
564
565impl Name {
566 pub fn new(name: &str) -> Name {
569 Name {
570 name: name.to_owned(),
571 namespace: None,
572 aliases: None,
573 }
574 }
575
576 fn parse(complex: &Map<String, Value>) -> Result<Self, AvroError> {
578 let name = complex
579 .name()
580 .ok_or_else(|| ParseSchemaError::new("No `name` field"))?;
581 if name.is_empty() {
582 return Err(ParseSchemaError::new(format!(
583 "Name cannot be the empty string: {:?}",
584 complex
585 ))
586 .into());
587 }
588
589 let (namespace, name) = if let Some(index) = name.rfind('.') {
590 let computed_namespace = name[..index].to_owned();
591 let computed_name = name[index + 1..].to_owned();
592 if let Some(provided_namespace) = complex.string("namespace") {
593 if provided_namespace != computed_namespace {
594 warn!(
595 "Found dots in name {}, updating to namespace {} and name {}",
596 name, computed_namespace, computed_name
597 );
598 }
599 }
600 (Some(computed_namespace), computed_name)
601 } else {
602 (complex.string("namespace"), name)
603 };
604
605 if !Regex::new(r"(^[A-Za-z_][A-Za-z0-9_]*)$")
606 .unwrap()
607 .is_match(&name)
608 {
609 return Err(ParseSchemaError::new(format!(
610 "Invalid name. Must start with [A-Za-z_] and subsequently only contain [A-Za-z0-9_]. Found: {}",
611 name
612 ))
613 .into());
614 }
615
616 let aliases: Option<Vec<String>> = complex
617 .get("aliases")
618 .and_then(|aliases| aliases.as_array())
619 .and_then(|aliases| {
620 aliases
621 .iter()
622 .map(|alias| alias.as_str())
623 .map(|alias| alias.map(|a| a.to_string()))
624 .collect::<Option<_>>()
625 });
626
627 Ok(Name {
628 name,
629 namespace,
630 aliases,
631 })
632 }
633
634 pub fn fullname(&self, default_namespace: &str) -> FullName {
639 FullName::from_parts(&self.name, self.namespace.as_deref(), default_namespace)
640 }
641}
642
643#[derive(Clone, Debug, PartialEq)]
644pub struct ResolvedDefaultValueField {
645 pub name: String,
646 pub doc: Documentation,
647 pub default: types::Value,
648 pub order: RecordFieldOrder,
649 pub position: usize,
650}
651
652#[derive(Clone, Debug, PartialEq)]
653pub enum ResolvedRecordField {
654 Absent(Schema),
655 Present(RecordField),
656}
657
658#[derive(Clone, Debug, PartialEq)]
660pub struct RecordField {
661 pub name: String,
663 pub doc: Documentation,
665 pub default: Option<Value>,
669 pub schema: SchemaPieceOrNamed,
671 pub order: RecordFieldOrder,
675 pub position: usize,
677}
678
679#[derive(Copy, Clone, Debug, PartialEq)]
681pub enum RecordFieldOrder {
682 Ascending,
683 Descending,
684 Ignore,
685}
686
687impl RecordField {}
688
689#[derive(Debug, Clone)]
690pub struct UnionSchema {
691 schemas: Vec<SchemaPieceOrNamed>,
692
693 anon_variant_index: HashMap<SchemaKind, usize>,
696
697 named_variant_index: HashMap<usize, usize>,
699}
700
701impl UnionSchema {
702 pub(crate) fn new(schemas: Vec<SchemaPieceOrNamed>) -> Result<Self, AvroError> {
703 let mut avindex = HashMap::new();
704 let mut nvindex = HashMap::new();
705 for (i, schema) in schemas.iter().enumerate() {
706 match schema {
707 SchemaPieceOrNamed::Piece(sp) => {
708 if let SchemaPiece::Union(_) = sp {
709 return Err(ParseSchemaError::new(
710 "Unions may not directly contain a union",
711 )
712 .into());
713 }
714 let kind = SchemaKind::from(sp);
715 if avindex.insert(kind, i).is_some() {
716 return Err(
717 ParseSchemaError::new("Unions cannot contain duplicate types").into(),
718 );
719 }
720 }
721 SchemaPieceOrNamed::Named(idx) => {
722 if nvindex.insert(*idx, i).is_some() {
723 return Err(
724 ParseSchemaError::new("Unions cannot contain duplicate types").into(),
725 );
726 }
727 }
728 }
729 }
730 Ok(UnionSchema {
731 schemas,
732 anon_variant_index: avindex,
733 named_variant_index: nvindex,
734 })
735 }
736
737 pub fn variants(&self) -> &[SchemaPieceOrNamed] {
739 &self.schemas
740 }
741
742 pub fn is_nullable(&self) -> bool {
744 !self.schemas.is_empty() && self.schemas[0] == SchemaPieceOrNamed::Piece(SchemaPiece::Null)
745 }
746
747 pub fn match_piece(&self, sp: &SchemaPiece) -> Option<(usize, &SchemaPieceOrNamed)> {
748 self.anon_variant_index
749 .get(&SchemaKind::from(sp))
750 .map(|idx| (*idx, &self.schemas[*idx]))
751 }
752
753 pub fn match_ref(
754 &self,
755 other: SchemaPieceRefOrNamed,
756 names_map: &HashMap<usize, usize>,
757 ) -> Option<(usize, &SchemaPieceOrNamed)> {
758 match other {
759 SchemaPieceRefOrNamed::Piece(sp) => self.match_piece(sp),
760 SchemaPieceRefOrNamed::Named(idx) => names_map
761 .get(&idx)
762 .and_then(|idx| self.named_variant_index.get(idx))
763 .map(|idx| (*idx, &self.schemas[*idx])),
764 }
765 }
766
767 #[inline(always)]
768 pub fn match_(
769 &self,
770 other: &SchemaPieceOrNamed,
771 names_map: &HashMap<usize, usize>,
772 ) -> Option<(usize, &SchemaPieceOrNamed)> {
773 self.match_ref(other.as_ref(), names_map)
774 }
775}
776
777impl PartialEq for UnionSchema {
779 fn eq(&self, other: &UnionSchema) -> bool {
780 self.schemas.eq(&other.schemas)
781 }
782}
783
784#[derive(Default)]
785struct SchemaParser {
786 named: Vec<Option<NamedSchemaPiece>>,
787 indices: HashMap<FullName, usize>,
788}
789
790impl SchemaParser {
791 fn parse(mut self, value: &Value) -> Result<Schema, AvroError> {
792 let top = self.parse_inner("", value)?;
793 let SchemaParser { named, indices } = self;
794 Ok(Schema {
795 named: named.into_iter().map(|o| o.unwrap()).collect(),
796 indices,
797 top,
798 })
799 }
800
801 fn parse_inner(
802 &mut self,
803 default_namespace: &str,
804 value: &Value,
805 ) -> Result<SchemaPieceOrNamed, AvroError> {
806 match *value {
807 Value::String(ref t) => {
808 let name = FullName::from_parts(t.as_str(), None, default_namespace);
809 if let Some(idx) = self.indices.get(&name) {
810 Ok(SchemaPieceOrNamed::Named(*idx))
811 } else {
812 Ok(SchemaPieceOrNamed::Piece(Schema::parse_primitive(
813 t.as_str(),
814 )?))
815 }
816 }
817 Value::Object(ref data) => self.parse_complex(default_namespace, data),
818 Value::Array(ref data) => Ok(SchemaPieceOrNamed::Piece(
819 self.parse_union(default_namespace, data)?,
820 )),
821 _ => Err(ParseSchemaError::new("Must be a JSON string, object or array").into()),
822 }
823 }
824
825 fn alloc_name(&mut self, fullname: FullName) -> Result<usize, AvroError> {
826 let idx = match self.indices.entry(fullname) {
827 Entry::Vacant(ve) => *ve.insert(self.named.len()),
828 Entry::Occupied(oe) => {
829 return Err(ParseSchemaError::new(format!(
830 "Sub-schema with name {} encountered multiple times",
831 oe.key()
832 ))
833 .into())
834 }
835 };
836 self.named.push(None);
837 Ok(idx)
838 }
839
840 fn insert(&mut self, index: usize, schema: NamedSchemaPiece) {
841 assert!(self.named[index].is_none());
842 self.named[index] = Some(schema);
843 }
844
845 fn parse_named_type(
846 &mut self,
847 type_name: &str,
848 default_namespace: &str,
849 complex: &Map<String, Value>,
850 ) -> Result<usize, AvroError> {
851 let name = Name::parse(complex)?;
852 match name.name.as_str() {
853 "null" | "boolean" | "int" | "long" | "float" | "double" | "bytes" | "string" => {
854 return Err(ParseSchemaError::new(format!(
855 "{} may not be used as a custom type name",
856 name.name
857 ))
858 .into())
859 }
860 _ => {}
861 };
862 let fullname = name.fullname(default_namespace);
863 let default_namespace = fullname.namespace.clone();
864 let idx = self.alloc_name(fullname.clone())?;
865 let piece = match type_name {
866 "record" => self.parse_record(&default_namespace, complex),
867 "enum" => self.parse_enum(complex),
868 "fixed" => self.parse_fixed(&default_namespace, complex),
869 _ => unreachable!("Unknown named type kind: {}", type_name),
870 }?;
871
872 self.insert(
873 idx,
874 NamedSchemaPiece {
875 name: fullname,
876 piece,
877 },
878 );
879
880 Ok(idx)
881 }
882
883 fn parse_complex(
889 &mut self,
890 default_namespace: &str,
891 complex: &Map<String, Value>,
892 ) -> Result<SchemaPieceOrNamed, AvroError> {
893 match complex.get("type") {
894 Some(&Value::String(ref t)) => Ok(match t.as_str() {
895 "record" | "enum" | "fixed" => SchemaPieceOrNamed::Named(self.parse_named_type(
896 t,
897 default_namespace,
898 complex,
899 )?),
900 "array" => SchemaPieceOrNamed::Piece(self.parse_array(default_namespace, complex)?),
901 "map" => SchemaPieceOrNamed::Piece(self.parse_map(default_namespace, complex)?),
902 "bytes" => SchemaPieceOrNamed::Piece(Self::parse_bytes(complex)?),
903 "int" => SchemaPieceOrNamed::Piece(Self::parse_int(complex)?),
904 "long" => SchemaPieceOrNamed::Piece(Self::parse_long(complex)?),
905 "string" => SchemaPieceOrNamed::Piece(Self::from_string(complex)),
906 other => {
907 let name = FullName {
908 name: other.into(),
909 namespace: default_namespace.into(),
910 };
911 if let Some(idx) = self.indices.get(&name) {
912 SchemaPieceOrNamed::Named(*idx)
913 } else {
914 SchemaPieceOrNamed::Piece(Schema::parse_primitive(t.as_str())?)
915 }
916 }
917 }),
918 Some(&Value::Object(ref data)) => match data.get("type") {
919 Some(ref value) => self.parse_inner(default_namespace, value),
920 None => Err(
921 ParseSchemaError::new(format!("Unknown complex type: {:?}", complex)).into(),
922 ),
923 },
924 _ => Err(ParseSchemaError::new("No `type` in complex type").into()),
925 }
926 }
927
928 fn parse_record(
931 &mut self,
932 default_namespace: &str,
933 complex: &Map<String, Value>,
934 ) -> Result<SchemaPiece, AvroError> {
935 let mut lookup = HashMap::new();
936
937 let fields: Vec<RecordField> = complex
938 .get("fields")
939 .and_then(|fields| fields.as_array())
940 .ok_or_else(|| ParseSchemaError::new("No `fields` in record").into())
941 .and_then(|fields| {
942 fields
943 .iter()
944 .filter_map(|field| field.as_object())
945 .enumerate()
946 .map(|(position, field)| {
947 self.parse_record_field(default_namespace, field, position)
948 })
949 .collect::<Result<_, _>>()
950 })?;
951
952 for field in &fields {
953 lookup.insert(field.name.clone(), field.position);
954 }
955
956 Ok(SchemaPiece::Record {
957 doc: complex.doc(),
958 fields,
959 lookup,
960 })
961 }
962
963 fn parse_record_field(
965 &mut self,
966 default_namespace: &str,
967 field: &Map<String, Value>,
968 position: usize,
969 ) -> Result<RecordField, AvroError> {
970 let name = field
971 .name()
972 .ok_or_else(|| ParseSchemaError::new("No `name` in record field"))?;
973
974 let schema = field
975 .get("type")
976 .ok_or_else(|| ParseSchemaError::new("No `type` in record field").into())
977 .and_then(|type_| self.parse_inner(default_namespace, type_))?;
978
979 let default = field.get("default").cloned();
980
981 let order = field
982 .get("order")
983 .and_then(|order| order.as_str())
984 .and_then(|order| match order {
985 "ascending" => Some(RecordFieldOrder::Ascending),
986 "descending" => Some(RecordFieldOrder::Descending),
987 "ignore" => Some(RecordFieldOrder::Ignore),
988 _ => None,
989 })
990 .unwrap_or(RecordFieldOrder::Ascending);
991
992 Ok(RecordField {
993 name,
994 doc: field.doc(),
995 default,
996 schema,
997 order,
998 position,
999 })
1000 }
1001
1002 fn parse_enum(&mut self, complex: &Map<String, Value>) -> Result<SchemaPiece, AvroError> {
1005 let symbols: Vec<String> = complex
1006 .get("symbols")
1007 .and_then(|v| v.as_array())
1008 .ok_or_else(|| ParseSchemaError::new("No `symbols` field in enum"))
1009 .and_then(|symbols| {
1010 symbols
1011 .iter()
1012 .map(|symbol| symbol.as_str().map(|s| s.to_string()))
1013 .collect::<Option<_>>()
1014 .ok_or_else(|| ParseSchemaError::new("Unable to parse `symbols` in enum"))
1015 })?;
1016
1017 let mut unique_symbols: HashSet<&String> = HashSet::new();
1018 for symbol in symbols.iter() {
1019 if unique_symbols.contains(symbol) {
1020 return Err(ParseSchemaError::new(format!(
1021 "Enum symbols must be unique, found multiple: {}",
1022 symbol
1023 ))
1024 .into());
1025 } else {
1026 unique_symbols.insert(symbol);
1027 }
1028 }
1029
1030 let default_idx = if let Some(default) = complex.get("default") {
1031 let default_str = default.as_str().ok_or_else(|| {
1032 ParseSchemaError::new(format!(
1033 "Enum default should be a string, got: {:?}",
1034 default
1035 ))
1036 })?;
1037 let default_idx = symbols
1038 .iter()
1039 .position(|x| x == default_str)
1040 .ok_or_else(|| {
1041 ParseSchemaError::new(format!(
1042 "Enum default not found in list of symbols: {}",
1043 default_str
1044 ))
1045 })?;
1046 Some(default_idx)
1047 } else {
1048 None
1049 };
1050
1051 Ok(SchemaPiece::Enum {
1052 doc: complex.doc(),
1053 symbols,
1054 default_idx,
1055 })
1056 }
1057
1058 fn parse_array(
1061 &mut self,
1062 default_namespace: &str,
1063 complex: &Map<String, Value>,
1064 ) -> Result<SchemaPiece, AvroError> {
1065 complex
1066 .get("items")
1067 .ok_or_else(|| ParseSchemaError::new("No `items` in array").into())
1068 .and_then(|items| self.parse_inner(default_namespace, items))
1069 .map(|schema| SchemaPiece::Array(Box::new(schema)))
1070 }
1071
1072 fn parse_map(
1075 &mut self,
1076 default_namespace: &str,
1077 complex: &Map<String, Value>,
1078 ) -> Result<SchemaPiece, AvroError> {
1079 complex
1080 .get("values")
1081 .ok_or_else(|| ParseSchemaError::new("No `values` in map").into())
1082 .and_then(|items| self.parse_inner(default_namespace, items))
1083 .map(|schema| SchemaPiece::Map(Box::new(schema)))
1084 }
1085
1086 fn parse_union(
1089 &mut self,
1090 default_namespace: &str,
1091 items: &[Value],
1092 ) -> Result<SchemaPiece, AvroError> {
1093 items
1094 .iter()
1095 .map(|value| self.parse_inner(default_namespace, value))
1096 .collect::<Result<Vec<_>, _>>()
1097 .and_then(|schemas| Ok(SchemaPiece::Union(UnionSchema::new(schemas)?)))
1098 }
1099
1100 fn parse_decimal(complex: &Map<String, Value>) -> Result<(usize, usize), AvroError> {
1103 let precision = complex
1104 .get("precision")
1105 .and_then(|v| v.as_i64())
1106 .ok_or_else(|| ParseSchemaError::new("No `precision` in decimal"))?;
1107
1108 let scale = complex.get("scale").and_then(|v| v.as_i64()).unwrap_or(0);
1109
1110 if scale < 0 {
1111 return Err(ParseSchemaError::new("Decimal scale must be greater than zero").into());
1112 }
1113
1114 if precision < 0 {
1115 return Err(
1116 ParseSchemaError::new("Decimal precision must be greater than zero").into(),
1117 );
1118 }
1119
1120 if scale > precision {
1121 return Err(ParseSchemaError::new("Decimal scale is greater than precision").into());
1122 }
1123
1124 Ok((precision as usize, scale as usize))
1125 }
1126
1127 fn parse_bytes(complex: &Map<String, Value>) -> Result<SchemaPiece, AvroError> {
1130 let logical_type = complex.get("logicalType").and_then(|v| v.as_str());
1131
1132 if let Some("decimal") = logical_type {
1133 match Self::parse_decimal(complex) {
1134 Ok((precision, scale)) => {
1135 return Ok(SchemaPiece::Decimal {
1136 precision,
1137 scale,
1138 fixed_size: None,
1139 })
1140 }
1141 Err(e) => warn!(
1142 "parsing decimal as regular bytes due to parse error: {:?}, {:?}",
1143 complex, e
1144 ),
1145 }
1146 }
1147
1148 Ok(SchemaPiece::Bytes)
1149 }
1150
1151 fn parse_int(complex: &Map<String, Value>) -> Result<SchemaPiece, AvroError> {
1159 const AVRO_DATE: &str = "date";
1160 const DEBEZIUM_DATE: &str = "io.debezium.time.Date";
1161 const KAFKA_DATE: &str = "org.apache.kafka.connect.data.Date";
1162 if let Some(name) = complex.get("connect.name") {
1163 if name == DEBEZIUM_DATE || name == KAFKA_DATE {
1164 if name == KAFKA_DATE {
1165 warn!("using deprecated debezium date format");
1166 }
1167 return Ok(SchemaPiece::Date);
1168 }
1169 }
1170 if let Some(name) = complex.get("logicalType") {
1174 if name == AVRO_DATE {
1175 return Ok(SchemaPiece::Date);
1176 }
1177 }
1178 if !complex.is_empty() {
1179 debug!("parsing complex type as regular int: {:?}", complex);
1180 }
1181 Ok(SchemaPiece::Int)
1182 }
1183
1184 fn parse_long(complex: &Map<String, Value>) -> Result<SchemaPiece, AvroError> {
1192 const AVRO_MILLI_TS: &str = "timestamp-millis";
1193 const AVRO_MICRO_TS: &str = "timestamp-micros";
1194
1195 const CONNECT_MILLI_TS: &[&str] = &[
1196 "io.debezium.time.Timestamp",
1197 "org.apache.kafka.connect.data.Timestamp",
1198 ];
1199 const CONNECT_MICRO_TS: &str = "io.debezium.time.MicroTimestamp";
1200
1201 if let Some(serde_json::Value::String(name)) = complex.get("connect.name") {
1202 if CONNECT_MILLI_TS.contains(&&**name) {
1203 return Ok(SchemaPiece::TimestampMilli);
1204 }
1205 if name == CONNECT_MICRO_TS {
1206 return Ok(SchemaPiece::TimestampMicro);
1207 }
1208 }
1209 if let Some(name) = complex.get("logicalType") {
1210 if name == AVRO_MILLI_TS {
1211 return Ok(SchemaPiece::TimestampMilli);
1212 }
1213 if name == AVRO_MICRO_TS {
1214 return Ok(SchemaPiece::TimestampMicro);
1215 }
1216 }
1217 if !complex.is_empty() {
1218 debug!("parsing complex type as regular long: {:?}", complex);
1219 }
1220 Ok(SchemaPiece::Long)
1221 }
1222
1223 fn from_string(complex: &Map<String, Value>) -> SchemaPiece {
1224 const CONNECT_JSON: &str = "io.debezium.data.Json";
1225
1226 if let Some(serde_json::Value::String(name)) = complex.get("connect.name") {
1227 if CONNECT_JSON == name.as_str() {
1228 return SchemaPiece::Json;
1229 }
1230 }
1231 if let Some(name) = complex.get("logicalType") {
1232 if name == "uuid" {
1233 return SchemaPiece::Uuid;
1234 }
1235 }
1236 debug!("parsing complex type as regular string: {:?}", complex);
1237 SchemaPiece::String
1238 }
1239
1240 fn parse_fixed(
1243 &mut self,
1244 _default_namespace: &str,
1245 complex: &Map<String, Value>,
1246 ) -> Result<SchemaPiece, AvroError> {
1247 let _name = Name::parse(complex)?;
1248
1249 let size = complex
1250 .get("size")
1251 .and_then(|v| v.as_i64())
1252 .ok_or_else(|| ParseSchemaError::new("No `size` in fixed"))?;
1253 if size <= 0 {
1254 return Err(ParseSchemaError::new(format!(
1255 "Fixed values require a positive size attribute, found: {}",
1256 size
1257 ))
1258 .into());
1259 }
1260
1261 let logical_type = complex.get("logicalType").and_then(|v| v.as_str());
1262
1263 if let Some("decimal") = logical_type {
1264 match Self::parse_decimal(complex) {
1265 Ok((precision, scale)) => {
1266 let max = ((2_usize.pow((8 * size - 1) as u32) - 1) as f64).log10() as usize;
1267 if precision > max {
1268 warn!("Decimal precision {} requires more than {} bytes of space, parsing as fixed", precision, size);
1269 } else {
1270 return Ok(SchemaPiece::Decimal {
1271 precision,
1272 scale,
1273 fixed_size: Some(size as usize),
1274 });
1275 }
1276 }
1277 Err(e) => warn!(
1278 "parsing decimal as fixed due to parse error: {:?}, {:?}",
1279 complex, e
1280 ),
1281 }
1282 }
1283
1284 Ok(SchemaPiece::Fixed {
1285 size: size as usize,
1286 })
1287 }
1288}
1289
1290impl Schema {
1291 pub fn parse(value: &Value) -> Result<Self, AvroError> {
1294 let p = SchemaParser {
1295 named: vec![],
1296 indices: Default::default(),
1297 };
1298 p.parse(value)
1299 }
1300
1301 pub fn canonical_form(&self) -> String {
1306 let json = serde_json::to_value(self).unwrap();
1307 parsing_canonical_form(&json)
1308 }
1309
1310 pub fn fingerprint<D: Digest>(&self) -> SchemaFingerprint {
1317 let mut d = D::new();
1318 d.update(self.canonical_form());
1319 SchemaFingerprint {
1320 bytes: d.finalize().to_vec(),
1321 }
1322 }
1323
1324 fn parse_primitive(primitive: &str) -> Result<SchemaPiece, AvroError> {
1327 match primitive {
1328 "null" => Ok(SchemaPiece::Null),
1329 "boolean" => Ok(SchemaPiece::Boolean),
1330 "int" => Ok(SchemaPiece::Int),
1331 "long" => Ok(SchemaPiece::Long),
1332 "double" => Ok(SchemaPiece::Double),
1333 "float" => Ok(SchemaPiece::Float),
1334 "bytes" => Ok(SchemaPiece::Bytes),
1335 "string" => Ok(SchemaPiece::String),
1336 other => Err(ParseSchemaError::new(format!("Unknown type: {}", other)).into()),
1337 }
1338 }
1339}
1340
1341impl FromStr for Schema {
1342 type Err = AvroError;
1343
1344 fn from_str(input: &str) -> Result<Self, AvroError> {
1346 let value = serde_json::from_str(input)
1347 .map_err(|e| ParseSchemaError::new(format!("Error parsing JSON: {}", e)))?;
1348 Self::parse(&value)
1349 }
1350}
1351
1352#[derive(Clone, Debug, PartialEq)]
1353pub struct NamedSchemaPiece {
1354 pub name: FullName,
1355 pub piece: SchemaPiece,
1356}
1357
1358#[derive(Copy, Clone, Debug)]
1359pub struct SchemaNode<'a> {
1360 pub root: &'a Schema,
1361 pub inner: &'a SchemaPiece,
1362 pub name: Option<&'a FullName>,
1363}
1364
1365#[derive(Copy, Clone, Debug)]
1366pub enum SchemaPieceRefOrNamed<'a> {
1367 Piece(&'a SchemaPiece),
1368 Named(usize),
1369}
1370
1371impl<'a> SchemaPieceRefOrNamed<'a> {
1372 pub fn get_human_name(&self, root: &Schema) -> String {
1373 match self {
1374 Self::Piece(piece) => format!("{:?}", piece),
1375 Self::Named(idx) => format!("{}", root.lookup(*idx).name),
1376 }
1377 }
1378
1379 #[inline(always)]
1380 pub fn get_piece_and_name(self, root: &'a Schema) -> (&'a SchemaPiece, Option<&'a FullName>) {
1381 match self {
1382 SchemaPieceRefOrNamed::Piece(sp) => (sp, None),
1383 SchemaPieceRefOrNamed::Named(index) => {
1384 let named_piece = root.lookup(index);
1385 (&named_piece.piece, Some(&named_piece.name))
1386 }
1387 }
1388 }
1389}
1390
1391#[derive(Copy, Clone, Debug)]
1392pub struct SchemaNodeOrNamed<'a> {
1393 pub root: &'a Schema,
1394 pub inner: SchemaPieceRefOrNamed<'a>,
1395}
1396
1397impl<'a> SchemaNodeOrNamed<'a> {
1398 #[inline(always)]
1399 pub fn lookup(self) -> SchemaNode<'a> {
1400 let (inner, name) = self.inner.get_piece_and_name(self.root);
1401 SchemaNode {
1402 root: self.root,
1403 inner,
1404 name,
1405 }
1406 }
1407 #[inline(always)]
1408 pub fn step(self, next: &'a SchemaPieceOrNamed) -> Self {
1409 self.step_ref(next.as_ref())
1410 }
1411 #[inline(always)]
1412 pub fn step_ref(self, next: SchemaPieceRefOrNamed<'a>) -> Self {
1413 Self {
1414 root: self.root,
1415 inner: match next {
1416 SchemaPieceRefOrNamed::Piece(piece) => SchemaPieceRefOrNamed::Piece(piece),
1417 SchemaPieceRefOrNamed::Named(index) => SchemaPieceRefOrNamed::Named(index),
1418 },
1419 }
1420 }
1421
1422 pub fn to_schema(self) -> Schema {
1423 let mut cloner = SchemaSubtreeDeepCloner {
1424 old_root: self.root,
1425 old_to_new_names: Default::default(),
1426 named: Default::default(),
1427 };
1428 let piece = cloner.clone_piece_or_named(self.inner);
1429 let named: Vec<NamedSchemaPiece> = cloner.named.into_iter().map(Option::unwrap).collect();
1430 let indices: HashMap<FullName, usize> = named
1431 .iter()
1432 .enumerate()
1433 .map(|(i, nsp)| (nsp.name.clone(), i))
1434 .collect();
1435 Schema {
1436 named,
1437 indices,
1438 top: piece,
1439 }
1440 }
1441}
1442
1443struct SchemaSubtreeDeepCloner<'a> {
1444 old_root: &'a Schema,
1445 old_to_new_names: HashMap<usize, usize>,
1446 named: Vec<Option<NamedSchemaPiece>>,
1447}
1448
1449impl<'a> SchemaSubtreeDeepCloner<'a> {
1450 fn clone_piece(&mut self, piece: &SchemaPiece) -> SchemaPiece {
1451 match piece {
1452 SchemaPiece::Null => SchemaPiece::Null,
1453 SchemaPiece::Boolean => SchemaPiece::Boolean,
1454 SchemaPiece::Int => SchemaPiece::Int,
1455 SchemaPiece::Long => SchemaPiece::Long,
1456 SchemaPiece::Float => SchemaPiece::Float,
1457 SchemaPiece::Double => SchemaPiece::Double,
1458 SchemaPiece::Date => SchemaPiece::Date,
1459 SchemaPiece::TimestampMilli => SchemaPiece::TimestampMilli,
1460 SchemaPiece::TimestampMicro => SchemaPiece::TimestampMicro,
1461 SchemaPiece::Json => SchemaPiece::Json,
1462 SchemaPiece::Decimal {
1463 scale,
1464 precision,
1465 fixed_size,
1466 } => SchemaPiece::Decimal {
1467 scale: *scale,
1468 precision: *precision,
1469 fixed_size: *fixed_size,
1470 },
1471 SchemaPiece::Bytes => SchemaPiece::Bytes,
1472 SchemaPiece::String => SchemaPiece::String,
1473 SchemaPiece::Uuid => SchemaPiece::Uuid,
1474 SchemaPiece::Array(inner) => {
1475 SchemaPiece::Array(Box::new(self.clone_piece_or_named(inner.as_ref().as_ref())))
1476 }
1477 SchemaPiece::Map(inner) => {
1478 SchemaPiece::Map(Box::new(self.clone_piece_or_named(inner.as_ref().as_ref())))
1479 }
1480 SchemaPiece::Union(us) => SchemaPiece::Union(UnionSchema {
1481 schemas: us
1482 .schemas
1483 .iter()
1484 .map(|s| self.clone_piece_or_named(s.as_ref()))
1485 .collect(),
1486 anon_variant_index: us.anon_variant_index.clone(),
1487 named_variant_index: us.named_variant_index.clone(),
1488 }),
1489 SchemaPiece::ResolveIntLong => SchemaPiece::ResolveIntLong,
1490 SchemaPiece::ResolveIntFloat => SchemaPiece::ResolveIntFloat,
1491 SchemaPiece::ResolveIntDouble => SchemaPiece::ResolveIntDouble,
1492 SchemaPiece::ResolveLongFloat => SchemaPiece::ResolveLongFloat,
1493 SchemaPiece::ResolveLongDouble => SchemaPiece::ResolveLongDouble,
1494 SchemaPiece::ResolveFloatDouble => SchemaPiece::ResolveFloatDouble,
1495 SchemaPiece::ResolveIntTsMilli => SchemaPiece::ResolveIntTsMilli,
1496 SchemaPiece::ResolveIntTsMicro => SchemaPiece::ResolveIntTsMicro,
1497 SchemaPiece::ResolveDateTimestamp => SchemaPiece::ResolveDateTimestamp,
1498 SchemaPiece::ResolveConcreteUnion {
1499 index,
1500 inner,
1501 n_reader_variants,
1502 reader_null_variant,
1503 } => SchemaPiece::ResolveConcreteUnion {
1504 index: *index,
1505 inner: Box::new(self.clone_piece_or_named(inner.as_ref().as_ref())),
1506 n_reader_variants: *n_reader_variants,
1507 reader_null_variant: *reader_null_variant,
1508 },
1509 SchemaPiece::ResolveUnionUnion {
1510 permutation,
1511 n_reader_variants,
1512 reader_null_variant,
1513 } => SchemaPiece::ResolveUnionUnion {
1514 permutation: permutation
1515 .clone()
1516 .into_iter()
1517 .map(|o| o.map(|(idx, piece)| (idx, self.clone_piece_or_named(piece.as_ref()))))
1518 .collect(),
1519 n_reader_variants: *n_reader_variants,
1520 reader_null_variant: *reader_null_variant,
1521 },
1522 SchemaPiece::ResolveUnionConcrete { index, inner } => {
1523 SchemaPiece::ResolveUnionConcrete {
1524 index: *index,
1525 inner: Box::new(self.clone_piece_or_named(inner.as_ref().as_ref())),
1526 }
1527 }
1528 SchemaPiece::Record {
1529 doc,
1530 fields,
1531 lookup,
1532 } => SchemaPiece::Record {
1533 doc: doc.clone(),
1534 fields: fields
1535 .iter()
1536 .map(|rf| RecordField {
1537 name: rf.name.clone(),
1538 doc: rf.doc.clone(),
1539 default: rf.default.clone(),
1540 schema: self.clone_piece_or_named(rf.schema.as_ref()),
1541 order: rf.order,
1542 position: rf.position,
1543 })
1544 .collect(),
1545 lookup: lookup.clone(),
1546 },
1547 SchemaPiece::Enum {
1548 doc,
1549 symbols,
1550 default_idx,
1551 } => SchemaPiece::Enum {
1552 doc: doc.clone(),
1553 symbols: symbols.clone(),
1554 default_idx: *default_idx,
1555 },
1556 SchemaPiece::Fixed { size } => SchemaPiece::Fixed { size: *size },
1557 SchemaPiece::ResolveRecord {
1558 defaults,
1559 fields,
1560 n_reader_fields,
1561 } => SchemaPiece::ResolveRecord {
1562 defaults: defaults.clone(),
1563 fields: fields
1564 .iter()
1565 .map(|rf| match rf {
1566 ResolvedRecordField::Present(rf) => {
1567 ResolvedRecordField::Present(RecordField {
1568 name: rf.name.clone(),
1569 doc: rf.doc.clone(),
1570 default: rf.default.clone(),
1571 schema: self.clone_piece_or_named(rf.schema.as_ref()),
1572 order: rf.order,
1573 position: rf.position,
1574 })
1575 }
1576 ResolvedRecordField::Absent(writer_schema) => {
1577 ResolvedRecordField::Absent(writer_schema.clone())
1578 }
1579 })
1580 .collect(),
1581 n_reader_fields: *n_reader_fields,
1582 },
1583 SchemaPiece::ResolveEnum {
1584 doc,
1585 symbols,
1586 default,
1587 } => SchemaPiece::ResolveEnum {
1588 doc: doc.clone(),
1589 symbols: symbols.clone(),
1590 default: default.clone(),
1591 },
1592 }
1593 }
1594 fn clone_piece_or_named(&mut self, piece: SchemaPieceRefOrNamed) -> SchemaPieceOrNamed {
1595 match piece {
1596 SchemaPieceRefOrNamed::Piece(piece) => self.clone_piece(piece).into(),
1597 SchemaPieceRefOrNamed::Named(index) => {
1598 let new_index = match self.old_to_new_names.entry(index) {
1599 Entry::Vacant(ve) => {
1600 let new_index = self.named.len();
1601 self.named.push(None);
1602 ve.insert(new_index);
1603 let old_named_piece = self.old_root.lookup(index);
1604 let new_named_piece = NamedSchemaPiece {
1605 name: old_named_piece.name.clone(),
1606 piece: self.clone_piece(&old_named_piece.piece),
1607 };
1608 self.named[new_index] = Some(new_named_piece);
1609 new_index
1610 }
1611 Entry::Occupied(oe) => *oe.get(),
1612 };
1613 SchemaPieceOrNamed::Named(new_index)
1614 }
1615 }
1616 }
1617}
1618
1619impl<'a> SchemaNode<'a> {
1620 #[inline(always)]
1621 pub fn step(self, next: &'a SchemaPieceOrNamed) -> Self {
1622 let (inner, name) = next.get_piece_and_name(self.root);
1623 Self {
1624 root: self.root,
1625 inner,
1626 name,
1627 }
1628 }
1629
1630 pub fn json_to_value(self, json: &serde_json::Value) -> Result<AvroValue, ParseSchemaError> {
1631 use serde_json::Value::*;
1632 let val = match (json, self.inner) {
1633 (json, SchemaPiece::Union(us)) => match us.schemas.first() {
1635 Some(variant) => AvroValue::Union {
1636 index: 0,
1637 inner: Box::new(self.step(variant).json_to_value(json)?),
1638 n_variants: us.schemas.len(),
1639 null_variant: us
1640 .schemas
1641 .iter()
1642 .position(|s| s == &SchemaPieceOrNamed::Piece(SchemaPiece::Null)),
1643 },
1644 None => return Err(ParseSchemaError("Union schema has no variants".to_owned())),
1645 },
1646 (Null, SchemaPiece::Null) => AvroValue::Null,
1647 (Bool(b), SchemaPiece::Boolean) => AvroValue::Boolean(*b),
1648 (Number(n), piece) => match piece {
1649 SchemaPiece::Int => {
1650 let i = n
1651 .as_i64()
1652 .and_then(|i| i32::try_from(i).ok())
1653 .ok_or_else(|| {
1654 ParseSchemaError(format!("{} is not a 32-bit integer", n))
1655 })?;
1656 AvroValue::Int(i)
1657 }
1658 SchemaPiece::Long => {
1659 let i = n.as_i64().ok_or_else(|| {
1660 ParseSchemaError(format!("{} is not a 64-bit integer", n))
1661 })?;
1662 AvroValue::Long(i)
1663 }
1664 SchemaPiece::Float => {
1665 let f = n
1666 .as_f64()
1667 .ok_or_else(|| ParseSchemaError(format!("{} is not a 32-bit float", n)))?;
1668 AvroValue::Float(f as f32)
1669 }
1670 SchemaPiece::Double => {
1671 let f = n
1672 .as_f64()
1673 .ok_or_else(|| ParseSchemaError(format!("{} is not a 64-bit float", n)))?;
1674 AvroValue::Double(f)
1675 }
1676 _ => {
1677 return Err(ParseSchemaError(format!(
1678 "Unexpected number in default: {}",
1679 n
1680 )))
1681 }
1682 },
1683 (String(s), piece)
1684 if s.eq_ignore_ascii_case("nan")
1685 && (piece == &SchemaPiece::Float || piece == &SchemaPiece::Double) =>
1686 {
1687 match piece {
1688 SchemaPiece::Float => AvroValue::Float(f32::NAN),
1689 SchemaPiece::Double => AvroValue::Double(f64::NAN),
1690 _ => unreachable!(),
1691 }
1692 }
1693 (String(s), piece)
1694 if s.eq_ignore_ascii_case("infinity")
1695 && (piece == &SchemaPiece::Float || piece == &SchemaPiece::Double) =>
1696 {
1697 match piece {
1698 SchemaPiece::Float => AvroValue::Float(f32::INFINITY),
1699 SchemaPiece::Double => AvroValue::Double(f64::INFINITY),
1700 _ => unreachable!(),
1701 }
1702 }
1703 (String(s), piece)
1704 if s.eq_ignore_ascii_case("-infinity")
1705 && (piece == &SchemaPiece::Float || piece == &SchemaPiece::Double) =>
1706 {
1707 match piece {
1708 SchemaPiece::Float => AvroValue::Float(f32::NEG_INFINITY),
1709 SchemaPiece::Double => AvroValue::Double(f64::NEG_INFINITY),
1710 _ => unreachable!(),
1711 }
1712 }
1713 (String(s), SchemaPiece::Bytes) => AvroValue::Bytes(s.clone().into_bytes()),
1714 (
1715 String(s),
1716 SchemaPiece::Decimal {
1717 precision, scale, ..
1718 },
1719 ) => AvroValue::Decimal(DecimalValue {
1720 precision: *precision,
1721 scale: *scale,
1722 unscaled: s.clone().into_bytes(),
1723 }),
1724 (String(s), SchemaPiece::String) => AvroValue::String(s.clone()),
1725 (Object(map), SchemaPiece::Record { fields, .. }) => {
1726 let field_values = fields
1727 .iter()
1728 .map(|rf| {
1729 let jval = map.get(&rf.name).ok_or_else(|| {
1730 ParseSchemaError(format!(
1731 "Field not found in default value: {}",
1732 rf.name
1733 ))
1734 })?;
1735 let value = self.step(&rf.schema).json_to_value(jval)?;
1736 Ok((rf.name.clone(), value))
1737 })
1738 .collect::<Result<Vec<(std::string::String, AvroValue)>, ParseSchemaError>>()?;
1739 AvroValue::Record(field_values)
1740 }
1741 (String(s), SchemaPiece::Enum { symbols, .. }) => {
1742 match symbols.iter().find_position(|sym| s == *sym) {
1743 Some((index, sym)) => AvroValue::Enum(index, sym.clone()),
1744 None => return Err(ParseSchemaError(format!("Enum variant not found: {}", s))),
1745 }
1746 }
1747 (Array(vals), SchemaPiece::Array(inner)) => {
1748 let node = self.step(&**inner);
1749 let vals = vals
1750 .iter()
1751 .map(|val| node.json_to_value(val))
1752 .collect::<Result<Vec<_>, ParseSchemaError>>()?;
1753 AvroValue::Array(vals)
1754 }
1755 (Object(map), SchemaPiece::Map(inner)) => {
1756 let node = self.step(&**inner);
1757 let map = map
1758 .iter()
1759 .map(|(k, v)| node.json_to_value(v).map(|v| (k.clone(), v)))
1760 .collect::<Result<HashMap<_, _>, ParseSchemaError>>()?;
1761 AvroValue::Map(AvroMap(map))
1762 }
1763 (String(s), SchemaPiece::Fixed { size }) if s.len() == *size => {
1764 AvroValue::Fixed(*size, s.clone().into_bytes())
1765 }
1766 _ => {
1767 return Err(ParseSchemaError(format!(
1768 "Json default value {} does not match schema",
1769 json
1770 )))
1771 }
1772 };
1773 Ok(val)
1774 }
1775}
1776
1777#[derive(Clone)]
1778struct SchemaSerContext<'a> {
1779 node: SchemaNodeOrNamed<'a>,
1780 seen_named: Rc<RefCell<HashMap<usize, String>>>,
1785}
1786
1787#[derive(Clone)]
1788struct RecordFieldSerContext<'a> {
1789 outer: &'a SchemaSerContext<'a>,
1790 inner: &'a RecordField,
1791}
1792
1793impl<'a> SchemaSerContext<'a> {
1794 fn step(&'a self, next: SchemaPieceRefOrNamed<'a>) -> Self {
1795 Self {
1796 node: self.node.step_ref(next),
1797 seen_named: Rc::clone(&self.seen_named),
1798 }
1799 }
1800}
1801
1802impl<'a> Serialize for SchemaSerContext<'a> {
1803 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1804 where
1805 S: Serializer,
1806 {
1807 match self.node.inner {
1808 SchemaPieceRefOrNamed::Piece(piece) => match piece {
1809 SchemaPiece::Null => serializer.serialize_str("null"),
1810 SchemaPiece::Boolean => serializer.serialize_str("boolean"),
1811 SchemaPiece::Int => serializer.serialize_str("int"),
1812 SchemaPiece::Long => serializer.serialize_str("long"),
1813 SchemaPiece::Float => serializer.serialize_str("float"),
1814 SchemaPiece::Double => serializer.serialize_str("double"),
1815 SchemaPiece::Date => {
1816 let mut map = serializer.serialize_map(Some(2))?;
1817 map.serialize_entry("type", "int")?;
1818 map.serialize_entry("logicalType", "date")?;
1819 map.end()
1820 }
1821 SchemaPiece::TimestampMilli | SchemaPiece::TimestampMicro => {
1822 let mut map = serializer.serialize_map(Some(2))?;
1823 map.serialize_entry("type", "long")?;
1824 if piece == &SchemaPiece::TimestampMilli {
1825 map.serialize_entry("logicalType", "timestamp-millis")?;
1826 } else {
1827 map.serialize_entry("logicalType", "timestamp-micros")?;
1828 }
1829 map.end()
1830 }
1831 SchemaPiece::Decimal {
1832 precision,
1833 scale,
1834 fixed_size: None,
1835 } => {
1836 let mut map = serializer.serialize_map(Some(4))?;
1837 map.serialize_entry("type", "bytes")?;
1838 map.serialize_entry("precision", precision)?;
1839 map.serialize_entry("scale", scale)?;
1840 map.serialize_entry("logicalType", "decimal")?;
1841 map.end()
1842 }
1843 SchemaPiece::Bytes => serializer.serialize_str("bytes"),
1844 SchemaPiece::String => serializer.serialize_str("string"),
1845 SchemaPiece::Array(inner) => {
1846 let mut map = serializer.serialize_map(Some(2))?;
1847 map.serialize_entry("type", "array")?;
1848 map.serialize_entry("items", &self.step(inner.as_ref().as_ref()))?;
1849 map.end()
1850 }
1851 SchemaPiece::Map(inner) => {
1852 let mut map = serializer.serialize_map(Some(2))?;
1853 map.serialize_entry("type", "map")?;
1854 map.serialize_entry("values", &self.step(inner.as_ref().as_ref()))?;
1855 map.end()
1856 }
1857 SchemaPiece::Union(inner) => {
1858 let variants = inner.variants();
1859 let mut seq = serializer.serialize_seq(Some(variants.len()))?;
1860 for v in variants {
1861 seq.serialize_element(&self.step(v.as_ref()))?;
1862 }
1863 seq.end()
1864 }
1865 SchemaPiece::Json => {
1866 let mut map = serializer.serialize_map(Some(2))?;
1867 map.serialize_entry("type", "string")?;
1868 map.serialize_entry("connect.name", "io.debezium.data.Json")?;
1869 map.end()
1870 }
1871 SchemaPiece::Uuid => {
1872 let mut map = serializer.serialize_map(Some(4))?;
1873 map.serialize_entry("type", "string")?;
1874 map.serialize_entry("logicalType", "uuid")?;
1875 map.end()
1876 }
1877 SchemaPiece::Record { .. }
1878 | SchemaPiece::Decimal {
1879 fixed_size: Some(_),
1880 ..
1881 }
1882 | SchemaPiece::Enum { .. }
1883 | SchemaPiece::Fixed { .. } => {
1884 unreachable!("Unexpected named schema piece in anonymous schema position")
1885 }
1886 SchemaPiece::ResolveIntLong
1887 | SchemaPiece::ResolveDateTimestamp
1888 | SchemaPiece::ResolveIntFloat
1889 | SchemaPiece::ResolveIntDouble
1890 | SchemaPiece::ResolveLongFloat
1891 | SchemaPiece::ResolveLongDouble
1892 | SchemaPiece::ResolveFloatDouble
1893 | SchemaPiece::ResolveConcreteUnion { .. }
1894 | SchemaPiece::ResolveUnionUnion { .. }
1895 | SchemaPiece::ResolveUnionConcrete { .. }
1896 | SchemaPiece::ResolveRecord { .. }
1897 | SchemaPiece::ResolveIntTsMicro
1898 | SchemaPiece::ResolveIntTsMilli
1899 | SchemaPiece::ResolveEnum { .. } => {
1900 panic!("Attempted to serialize resolved schema")
1901 }
1902 },
1903 SchemaPieceRefOrNamed::Named(index) => {
1904 let mut map = self.seen_named.borrow_mut();
1905 let named_piece = match map.get(&index) {
1906 Some(name) => {
1907 return serializer.serialize_str(name.as_str());
1908 }
1909 None => self.node.root.lookup(index),
1910 };
1911 let name = named_piece.name.to_string();
1912 map.insert(index, name.clone());
1913 std::mem::drop(map);
1914 match &named_piece.piece {
1915 SchemaPiece::Record { doc, fields, .. } => {
1916 let mut map = serializer.serialize_map(None)?;
1917 map.serialize_entry("type", "record")?;
1918 map.serialize_entry("name", &name)?;
1919 if let Some(ref docstr) = doc {
1920 map.serialize_entry("doc", docstr)?;
1921 }
1922 map.serialize_entry(
1924 "fields",
1925 &fields
1926 .iter()
1927 .map(|f| RecordFieldSerContext {
1928 outer: self,
1929 inner: f,
1930 })
1931 .collect::<Vec<_>>(),
1932 )?;
1933 map.end()
1934 }
1935 SchemaPiece::Enum {
1936 symbols,
1937 default_idx,
1938 ..
1939 } => {
1940 let mut map = serializer.serialize_map(None)?;
1941 map.serialize_entry("type", "enum")?;
1942 map.serialize_entry("name", &name)?;
1943 map.serialize_entry("symbols", symbols)?;
1944 if let Some(default_idx) = *default_idx {
1945 assert!(default_idx < symbols.len());
1946 map.serialize_entry("default", &symbols[default_idx])?;
1947 }
1948 map.end()
1949 }
1950 SchemaPiece::Fixed { size } => {
1951 let mut map = serializer.serialize_map(None)?;
1952 map.serialize_entry("type", "fixed")?;
1953 map.serialize_entry("name", &name)?;
1954 map.serialize_entry("size", size)?;
1955 map.end()
1956 }
1957 SchemaPiece::Decimal {
1958 scale,
1959 precision,
1960 fixed_size: Some(size),
1961 } => {
1962 let mut map = serializer.serialize_map(Some(6))?;
1963 map.serialize_entry("type", "fixed")?;
1964 map.serialize_entry("logicalType", "decimal")?;
1965 map.serialize_entry("name", &name)?;
1966 map.serialize_entry("size", size)?;
1967 map.serialize_entry("precision", precision)?;
1968 map.serialize_entry("scale", scale)?;
1969 map.end()
1970 }
1971 SchemaPiece::Null
1972 | SchemaPiece::Boolean
1973 | SchemaPiece::Int
1974 | SchemaPiece::Long
1975 | SchemaPiece::Float
1976 | SchemaPiece::Double
1977 | SchemaPiece::Date
1978 | SchemaPiece::TimestampMilli
1979 | SchemaPiece::TimestampMicro
1980 | SchemaPiece::Decimal {
1981 fixed_size: None, ..
1982 }
1983 | SchemaPiece::Bytes
1984 | SchemaPiece::String
1985 | SchemaPiece::Array(_)
1986 | SchemaPiece::Map(_)
1987 | SchemaPiece::Union(_)
1988 | SchemaPiece::Uuid
1989 | SchemaPiece::Json => {
1990 unreachable!("Unexpected anonymous schema piece in named schema position")
1991 }
1992 SchemaPiece::ResolveIntLong
1993 | SchemaPiece::ResolveDateTimestamp
1994 | SchemaPiece::ResolveIntFloat
1995 | SchemaPiece::ResolveIntDouble
1996 | SchemaPiece::ResolveLongFloat
1997 | SchemaPiece::ResolveLongDouble
1998 | SchemaPiece::ResolveFloatDouble
1999 | SchemaPiece::ResolveConcreteUnion { .. }
2000 | SchemaPiece::ResolveUnionUnion { .. }
2001 | SchemaPiece::ResolveUnionConcrete { .. }
2002 | SchemaPiece::ResolveRecord { .. }
2003 | SchemaPiece::ResolveIntTsMilli
2004 | SchemaPiece::ResolveIntTsMicro
2005 | SchemaPiece::ResolveEnum { .. } => {
2006 panic!("Attempted to serialize resolved schema")
2007 }
2008 }
2009 }
2010 }
2011 }
2012}
2013
2014impl Serialize for Schema {
2015 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
2016 where
2017 S: Serializer,
2018 {
2019 let ctx = SchemaSerContext {
2020 node: SchemaNodeOrNamed {
2021 root: self,
2022 inner: self.top.as_ref(),
2023 },
2024 seen_named: Rc::new(RefCell::new(Default::default())),
2025 };
2026 ctx.serialize(serializer)
2027 }
2028}
2029
2030impl<'a> Serialize for RecordFieldSerContext<'a> {
2031 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
2032 where
2033 S: Serializer,
2034 {
2035 let mut map = serializer.serialize_map(None)?;
2036 map.serialize_entry("name", &self.inner.name)?;
2037 map.serialize_entry("type", &self.outer.step(self.inner.schema.as_ref()))?;
2038 if let Some(default) = &self.inner.default {
2039 map.serialize_entry("default", default)?;
2040 }
2041 map.end()
2042 }
2043}
2044
2045fn parsing_canonical_form(schema: &serde_json::Value) -> String {
2048 match schema {
2049 serde_json::Value::Object(map) => pcf_map(map),
2050 serde_json::Value::String(s) => pcf_string(s),
2051 serde_json::Value::Array(v) => pcf_array(v),
2052 serde_json::Value::Number(n) => n.to_string(),
2053 _ => unreachable!("{:?} cannot yet be printed in canonical form", schema),
2054 }
2055}
2056
2057fn pcf_map(schema: &Map<String, serde_json::Value>) -> String {
2058 let ns = schema.get("namespace").and_then(|v| v.as_str());
2060 let mut fields = Vec::new();
2061 for (k, v) in schema {
2062 if schema.len() == 1 && k == "type" {
2064 if let serde_json::Value::String(s) = v {
2066 return pcf_string(s);
2067 }
2068 }
2069
2070 if field_ordering_position(k).is_none() {
2072 continue;
2073 }
2074
2075 if k == "name" {
2077 let name = v.as_str().unwrap();
2079 let n = match ns {
2080 Some(namespace) if !name.contains('.') => {
2081 Cow::Owned(format!("{}.{}", namespace, name))
2082 }
2083 _ => Cow::Borrowed(name),
2084 };
2085
2086 fields.push((k, format!("{}:{}", pcf_string(k), pcf_string(&*n))));
2087 continue;
2088 }
2089
2090 if k == "size" {
2092 let i = match v.as_str() {
2093 Some(s) => s.parse::<i64>().expect("Only valid schemas are accepted!"),
2094 None => v.as_i64().unwrap(),
2095 };
2096 fields.push((k, format!("{}:{}", pcf_string(k), i)));
2097 continue;
2098 }
2099
2100 fields.push((
2102 k,
2103 format!("{}:{}", pcf_string(k), parsing_canonical_form(v)),
2104 ));
2105 }
2106
2107 fields.sort_unstable_by_key(|(k, _)| field_ordering_position(k).unwrap());
2109 let inter = fields
2110 .into_iter()
2111 .map(|(_, v)| v)
2112 .collect::<Vec<_>>()
2113 .join(",");
2114 format!("{{{}}}", inter)
2115}
2116
2117fn pcf_array(arr: &[serde_json::Value]) -> String {
2118 let inter = arr
2119 .iter()
2120 .map(parsing_canonical_form)
2121 .collect::<Vec<String>>()
2122 .join(",");
2123 format!("[{}]", inter)
2124}
2125
2126fn pcf_string(s: &str) -> String {
2127 format!("\"{}\"", s)
2128}
2129
2130fn field_ordering_position(field: &str) -> Option<usize> {
2132 let v = match field {
2133 "name" => 1,
2134 "type" => 2,
2135 "fields" => 3,
2136 "symbols" => 4,
2137 "items" => 5,
2138 "values" => 6,
2139 "size" => 7,
2140 _ => return None,
2141 };
2142
2143 Some(v)
2144}
2145
2146#[cfg(test)]
2147mod tests {
2148 use types::Record;
2149 use types::ToAvro;
2150
2151 use super::*;
2152
2153 fn check_schema(schema: &str, expected: SchemaPiece) {
2154 let schema = Schema::from_str(schema).unwrap();
2155 assert_eq!(&expected, schema.top_node().inner);
2156
2157 let schema = serde_json::to_string(&schema).unwrap();
2159 let schema = Schema::from_str(&schema).unwrap();
2160 assert_eq!(&expected, schema.top_node().inner);
2161 }
2162
2163 #[test]
2164 fn test_primitive_schema() {
2165 check_schema("\"null\"", SchemaPiece::Null);
2166 check_schema("\"int\"", SchemaPiece::Int);
2167 check_schema("\"double\"", SchemaPiece::Double);
2168 }
2169
2170 #[test]
2171 fn test_array_schema() {
2172 check_schema(
2173 r#"{"type": "array", "items": "string"}"#,
2174 SchemaPiece::Array(Box::new(SchemaPieceOrNamed::Piece(SchemaPiece::String))),
2175 );
2176 }
2177
2178 #[test]
2179 fn test_map_schema() {
2180 check_schema(
2181 r#"{"type": "map", "values": "double"}"#,
2182 SchemaPiece::Map(Box::new(SchemaPieceOrNamed::Piece(SchemaPiece::Double))),
2183 );
2184 }
2185
2186 #[test]
2187 fn test_union_schema() {
2188 check_schema(
2189 r#"["null", "int"]"#,
2190 SchemaPiece::Union(
2191 UnionSchema::new(vec![
2192 SchemaPieceOrNamed::Piece(SchemaPiece::Null),
2193 SchemaPieceOrNamed::Piece(SchemaPiece::Int),
2194 ])
2195 .unwrap(),
2196 ),
2197 );
2198 }
2199
2200 #[test]
2201 fn test_multi_union_schema() {
2202 let schema = Schema::from_str(r#"["null", "int", "float", "string", "bytes"]"#);
2203 assert!(schema.is_ok());
2204 let schema = schema.unwrap();
2205 let node = schema.top_node();
2206 assert_eq!(SchemaKind::from(&schema), SchemaKind::Union);
2207 let union_schema = match node.inner {
2208 SchemaPiece::Union(u) => u,
2209 _ => unreachable!(),
2210 };
2211 assert_eq!(union_schema.variants().len(), 5);
2212 let mut variants = union_schema.variants().iter();
2213 assert_eq!(
2214 SchemaKind::from(node.step(variants.next().unwrap())),
2215 SchemaKind::Null
2216 );
2217 assert_eq!(
2218 SchemaKind::from(node.step(variants.next().unwrap())),
2219 SchemaKind::Int
2220 );
2221 assert_eq!(
2222 SchemaKind::from(node.step(variants.next().unwrap())),
2223 SchemaKind::Float
2224 );
2225 assert_eq!(
2226 SchemaKind::from(node.step(variants.next().unwrap())),
2227 SchemaKind::String
2228 );
2229 assert_eq!(
2230 SchemaKind::from(node.step(variants.next().unwrap())),
2231 SchemaKind::Bytes
2232 );
2233 assert_eq!(variants.next(), None);
2234 }
2235
2236 #[test]
2237 fn test_record_schema() {
2238 let schema = r#"
2239 {
2240 "type": "record",
2241 "name": "test",
2242 "fields": [
2243 {"name": "a", "type": "long", "default": 42},
2244 {"name": "b", "type": "string"}
2245 ]
2246 }
2247 "#;
2248
2249 let mut lookup = HashMap::new();
2250 lookup.insert("a".to_owned(), 0);
2251 lookup.insert("b".to_owned(), 1);
2252
2253 let expected = SchemaPiece::Record {
2254 doc: None,
2255 fields: vec![
2256 RecordField {
2257 name: "a".to_string(),
2258 doc: None,
2259 default: Some(Value::Number(42i64.into())),
2260 schema: SchemaPiece::Long.into(),
2261 order: RecordFieldOrder::Ascending,
2262 position: 0,
2263 },
2264 RecordField {
2265 name: "b".to_string(),
2266 doc: None,
2267 default: None,
2268 schema: SchemaPiece::String.into(),
2269 order: RecordFieldOrder::Ascending,
2270 position: 1,
2271 },
2272 ],
2273 lookup,
2274 };
2275
2276 check_schema(schema, expected);
2277 }
2278
2279 #[test]
2280 fn test_enum_schema() {
2281 let schema = r#"{"type": "enum", "name": "Suit", "symbols": ["diamonds", "spades", "jokers", "clubs", "hearts"], "default": "jokers"}"#;
2282
2283 let expected = SchemaPiece::Enum {
2284 doc: None,
2285 symbols: vec![
2286 "diamonds".to_owned(),
2287 "spades".to_owned(),
2288 "jokers".to_owned(),
2289 "clubs".to_owned(),
2290 "hearts".to_owned(),
2291 ],
2292 default_idx: Some(2),
2293 };
2294
2295 check_schema(schema, expected);
2296
2297 let bad_schema = Schema::from_str(
2298 r#"{"type": "enum", "name": "Suit", "symbols": ["diamonds", "spades", "jokers", "clubs", "hearts"], "default": "blah"}"#,
2299 );
2300
2301 assert!(bad_schema.is_err());
2302 }
2303
2304 #[test]
2305 fn test_fixed_schema() {
2306 let schema = r#"{"type": "fixed", "name": "test", "size": 16}"#;
2307
2308 let expected = SchemaPiece::Fixed { size: 16usize };
2309
2310 check_schema(schema, expected);
2311 }
2312
2313 #[test]
2314 fn test_date_schema() {
2315 let kinds = &[
2316 r#"{
2317 "type": "int",
2318 "name": "datish",
2319 "logicalType": "date"
2320 }"#,
2321 r#"{
2322 "type": "int",
2323 "name": "datish",
2324 "connect.name": "io.debezium.time.Date"
2325 }"#,
2326 r#"{
2327 "type": "int",
2328 "name": "datish",
2329 "connect.name": "org.apache.kafka.connect.data.Date"
2330 }"#,
2331 ];
2332 for kind in kinds {
2333 check_schema(*kind, SchemaPiece::Date);
2334
2335 let schema = Schema::from_str(*kind).unwrap();
2336 assert_eq!(
2337 serde_json::to_string(&schema).unwrap(),
2338 r#"{"type":"int","logicalType":"date"}"#
2339 );
2340 }
2341 }
2342
2343 #[test]
2344 fn new_field_in_middle() {
2345 let reader = r#"{
2346 "type": "record",
2347 "name": "MyRecord",
2348 "fields": [{"name": "f1", "type": "int"}, {"name": "f2", "type": "int"}]
2349 }"#;
2350 let writer = r#"{
2351 "type": "record",
2352 "name": "MyRecord",
2353 "fields": [{"name": "f1", "type": "int"}, {"name": "f_interposed", "type": "int"}, {"name": "f2", "type": "int"}]
2354 }"#;
2355 let reader = Schema::from_str(reader).unwrap();
2356 let writer = Schema::from_str(writer).unwrap();
2357
2358 let mut record = Record::new(writer.top_node()).unwrap();
2359 record.put("f1", 1);
2360 record.put("f2", 2);
2361 record.put("f_interposed", 42);
2362
2363 let value = record.avro();
2364
2365 let mut buf = vec![];
2366 crate::encode::encode(&value, &writer, &mut buf);
2367
2368 let resolved = resolve_schemas(&writer, &reader).unwrap();
2369
2370 let reader = &mut &buf[..];
2371 let reader_value = crate::decode::decode(resolved.top_node(), reader).unwrap();
2372 let expected = crate::types::Value::Record(vec![
2373 ("f1".to_string(), crate::types::Value::Int(1)),
2374 ("f2".to_string(), crate::types::Value::Int(2)),
2375 ]);
2376 assert_eq!(reader_value, expected);
2377 assert!(reader.is_empty()); }
2379
2380 #[test]
2381 fn new_field_at_end() {
2382 let reader = r#"{
2383 "type": "record",
2384 "name": "MyRecord",
2385 "fields": [{"name": "f1", "type": "int"}]
2386 }"#;
2387 let writer = r#"{
2388 "type": "record",
2389 "name": "MyRecord",
2390 "fields": [{"name": "f1", "type": "int"}, {"name": "f2", "type": "int"}]
2391 }"#;
2392 let reader = Schema::from_str(reader).unwrap();
2393 let writer = Schema::from_str(writer).unwrap();
2394
2395 let mut record = Record::new(writer.top_node()).unwrap();
2396 record.put("f1", 1);
2397 record.put("f2", 2);
2398
2399 let value = record.avro();
2400
2401 let mut buf = vec![];
2402 crate::encode::encode(&value, &writer, &mut buf);
2403
2404 let resolved = resolve_schemas(&writer, &reader).unwrap();
2405
2406 let reader = &mut &buf[..];
2407 let reader_value = crate::decode::decode(resolved.top_node(), reader).unwrap();
2408 let expected =
2409 crate::types::Value::Record(vec![("f1".to_string(), crate::types::Value::Int(1))]);
2410 assert_eq!(reader_value, expected);
2411 assert!(reader.is_empty()); }
2413
2414 #[test]
2415 fn default_non_nums() {
2416 let reader = r#"{
2417 "type": "record",
2418 "name": "MyRecord",
2419 "fields": [
2420 {"name": "f1", "type": "double", "default": "NaN"},
2421 {"name": "f2", "type": "double", "default": "Infinity"},
2422 {"name": "f3", "type": "double", "default": "-Infinity"}
2423 ]
2424 }
2425 "#;
2426 let writer = r#"{"type": "record", "name": "MyRecord", "fields": []}"#;
2427
2428 let writer_schema = Schema::from_str(writer).unwrap();
2429 let reader_schema = Schema::from_str(reader).unwrap();
2430 let resolved = resolve_schemas(&writer_schema, &reader_schema).unwrap();
2431
2432 let record = Record::new(writer_schema.top_node()).unwrap();
2433
2434 let value = record.avro();
2435 let mut buf = vec![];
2436 crate::encode::encode(&value, &writer_schema, &mut buf);
2437
2438 let reader = &mut &buf[..];
2439 let reader_value = crate::decode::decode(resolved.top_node(), reader).unwrap();
2440 let expected = crate::types::Value::Record(vec![
2441 ("f1".to_string(), crate::types::Value::Double(f64::NAN)),
2442 ("f2".to_string(), crate::types::Value::Double(f64::INFINITY)),
2443 (
2444 "f3".to_string(),
2445 crate::types::Value::Double(f64::NEG_INFINITY),
2446 ),
2447 ]);
2448
2449 #[derive(Debug)]
2450 struct NanEq(crate::types::Value);
2451 impl std::cmp::PartialEq for NanEq {
2452 fn eq(&self, other: &Self) -> bool {
2453 match (self, other) {
2454 (
2455 NanEq(crate::types::Value::Double(x)),
2456 NanEq(crate::types::Value::Double(y)),
2457 ) if x.is_nan() && y.is_nan() => true,
2458 (
2459 NanEq(crate::types::Value::Float(x)),
2460 NanEq(crate::types::Value::Float(y)),
2461 ) if x.is_nan() && y.is_nan() => true,
2462 (
2463 NanEq(crate::types::Value::Record(xs)),
2464 NanEq(crate::types::Value::Record(ys)),
2465 ) => {
2466 let xs = xs
2467 .iter()
2468 .cloned()
2469 .map(|(k, v)| (k, NanEq(v)))
2470 .collect::<Vec<_>>();
2471 let ys = ys
2472 .iter()
2473 .cloned()
2474 .map(|(k, v)| (k, NanEq(v)))
2475 .collect::<Vec<_>>();
2476
2477 xs == ys
2478 }
2479 (NanEq(x), NanEq(y)) => x == y,
2480 }
2481 }
2482 }
2483
2484 assert_eq!(NanEq(reader_value), NanEq(expected));
2485 assert!(reader.is_empty());
2486 }
2487
2488 #[test]
2489 fn test_decimal_schemas() {
2490 let schema = r#"{
2491 "type": "fixed",
2492 "name": "dec",
2493 "size": 8,
2494 "logicalType": "decimal",
2495 "precision": 12,
2496 "scale": 5
2497 }"#;
2498 let expected = SchemaPiece::Decimal {
2499 precision: 12,
2500 scale: 5,
2501 fixed_size: Some(8),
2502 };
2503 check_schema(schema, expected);
2504
2505 let schema = r#"{
2506 "type": "bytes",
2507 "logicalType": "decimal",
2508 "precision": 12,
2509 "scale": 5
2510 }"#;
2511 let expected = SchemaPiece::Decimal {
2512 precision: 12,
2513 scale: 5,
2514 fixed_size: None,
2515 };
2516 check_schema(schema, expected);
2517
2518 let res = Schema::from_str(
2519 r#"["bytes", {
2520 "type": "bytes",
2521 "logicalType": "decimal",
2522 "precision": 12,
2523 "scale": 5
2524 }]"#,
2525 );
2526 assert_eq!(
2527 res.unwrap_err().to_string(),
2528 "Schema parse error: Unions cannot contain duplicate types"
2529 );
2530
2531 let writer_schema = Schema::from_str(
2532 r#"["null", {
2533 "type": "bytes"
2534 }]"#,
2535 )
2536 .unwrap();
2537 let reader_schema = Schema::from_str(
2538 r#"["null", {
2539 "type": "bytes",
2540 "logicalType": "decimal",
2541 "precision": 12,
2542 "scale": 5
2543 }]"#,
2544 )
2545 .unwrap();
2546 let resolved = resolve_schemas(&writer_schema, &reader_schema).unwrap();
2547
2548 let expected = SchemaPiece::ResolveUnionUnion {
2549 permutation: vec![
2550 Ok((0, SchemaPieceOrNamed::Piece(SchemaPiece::Null))),
2551 Ok((
2552 1,
2553 SchemaPieceOrNamed::Piece(SchemaPiece::Decimal {
2554 precision: 12,
2555 scale: 5,
2556 fixed_size: None,
2557 }),
2558 )),
2559 ],
2560 n_reader_variants: 2,
2561 reader_null_variant: Some(0),
2562 };
2563 assert_eq!(resolved.top_node().inner, &expected);
2564 }
2565
2566 #[test]
2567 fn test_no_documentation() {
2568 let schema =
2569 Schema::from_str(r#"{"type": "enum", "name": "Coin", "symbols": ["heads", "tails"]}"#)
2570 .unwrap();
2571
2572 let doc = match schema.top_node().inner {
2573 SchemaPiece::Enum { doc, .. } => doc.clone(),
2574 _ => panic!(),
2575 };
2576
2577 assert!(doc.is_none());
2578 }
2579
2580 #[test]
2581 fn test_documentation() {
2582 let schema = Schema::from_str(
2583 r#"{"type": "enum", "name": "Coin", "doc": "Some documentation", "symbols": ["heads", "tails"]}"#
2584 ).unwrap();
2585
2586 let doc = match schema.top_node().inner {
2587 SchemaPiece::Enum { doc, .. } => doc.clone(),
2588 _ => None,
2589 };
2590
2591 assert_eq!("Some documentation".to_owned(), doc.unwrap());
2592 }
2593
2594 #[test]
2595 fn test_namespaces_and_names() {
2596 let schema = Schema::from_str(
2598 r#"{"type": "fixed", "namespace": "namespace", "name": "name", "size": 1}"#,
2599 )
2600 .unwrap();
2601 assert_eq!(schema.named.len(), 1);
2602 assert_eq!(
2603 schema.named[0].name,
2604 FullName {
2605 name: "name".into(),
2606 namespace: "namespace".into()
2607 }
2608 );
2609
2610 let schema =
2612 Schema::from_str(r#"{"type": "enum", "name": "name.has.dots", "symbols": ["A", "B"]}"#)
2613 .unwrap();
2614 assert_eq!(schema.named.len(), 1);
2615 assert_eq!(
2616 schema.named[0].name,
2617 FullName {
2618 name: "dots".into(),
2619 namespace: "name.has".into()
2620 }
2621 );
2622
2623 let schema = Schema::from_str(
2625 r#"{"type": "enum", "namespace": "namespace",
2626 "name": "name.has.dots", "symbols": ["A", "B"]}"#,
2627 )
2628 .unwrap();
2629 assert_eq!(schema.named.len(), 1);
2630 assert_eq!(
2631 schema.named[0].name,
2632 FullName {
2633 name: "dots".into(),
2634 namespace: "name.has".into()
2635 }
2636 );
2637
2638 let schema = Schema::from_str(
2641 r#"{"type": "record", "name": "TestDoc", "doc": "Doc string",
2642 "fields": [{"name": "name", "type": "string"}]}"#,
2643 )
2644 .unwrap();
2645 assert_eq!(schema.named.len(), 1);
2646 assert_eq!(
2647 schema.named[0].name,
2648 FullName {
2649 name: "TestDoc".into(),
2650 namespace: "".into()
2651 }
2652 );
2653
2654 let schema = Schema::from_str(
2656 r#"{"type": "record", "namespace": "", "name": "TestDoc", "doc": "Doc string",
2657 "fields": [{"name": "name", "type": "string"}]}"#,
2658 )
2659 .unwrap();
2660 assert_eq!(schema.named.len(), 1);
2661 assert_eq!(
2662 schema.named[0].name,
2663 FullName {
2664 name: "TestDoc".into(),
2665 namespace: "".into()
2666 }
2667 );
2668
2669 let first = Schema::from_str(
2671 r#"{"type": "fixed", "namespace": "namespace",
2672 "name": "name", "size": 1}"#,
2673 )
2674 .unwrap();
2675 let second = Schema::from_str(
2676 r#"{"type": "fixed", "name": "namespace.name",
2677 "size": 1}"#,
2678 )
2679 .unwrap();
2680 assert_eq!(first.named[0].name, second.named[0].name);
2681
2682 let first = Schema::from_str(
2683 r#"{"type": "fixed", "namespace": "namespace",
2684 "name": "name", "size": 1}"#,
2685 )
2686 .unwrap();
2687 let second = Schema::from_str(
2688 r#"{"type": "fixed", "name": "namespace.Name",
2689 "size": 1}"#,
2690 )
2691 .unwrap();
2692 assert_ne!(first.named[0].name, second.named[0].name);
2693
2694 let first = Schema::from_str(
2695 r#"{"type": "fixed", "namespace": "Namespace",
2696 "name": "name", "size": 1}"#,
2697 )
2698 .unwrap();
2699 let second = Schema::from_str(
2700 r#"{"type": "fixed", "namespace": "namespace",
2701 "name": "name", "size": 1}"#,
2702 )
2703 .unwrap();
2704 assert_ne!(first.named[0].name, second.named[0].name);
2705
2706 assert!(Schema::from_str(
2709 r#"{"type": "record", "name": "99 problems but a name aint one",
2710 "fields": [{"name": "name", "type": "string"}]}"#
2711 )
2712 .is_err());
2713
2714 assert!(Schema::from_str(
2715 r#"{"type": "record", "name": "!!!",
2716 "fields": [{"name": "name", "type": "string"}]}"#
2717 )
2718 .is_err());
2719
2720 assert!(Schema::from_str(
2721 r#"{"type": "record", "name": "_valid_until_©",
2722 "fields": [{"name": "name", "type": "string"}]}"#
2723 )
2724 .is_err());
2725
2726 let schema = Schema::from_str(r#"{"type": "record", "name": "org.apache.avro.tests.Hello", "fields": [
2728 {"name": "f1", "type": {"type": "enum", "name": "MyEnum", "symbols": ["Foo", "Bar", "Baz"]}},
2729 {"name": "f2", "type": "org.apache.avro.tests.MyEnum"},
2730 {"name": "f3", "type": "MyEnum"},
2731 {"name": "f4", "type": {"type": "enum", "name": "other.namespace.OtherEnum", "symbols": ["one", "two", "three"]}},
2732 {"name": "f5", "type": "other.namespace.OtherEnum"},
2733 {"name": "f6", "type": {"type": "enum", "name": "ThirdEnum", "namespace": "some.other", "symbols": ["Alice", "Bob"]}},
2734 {"name": "f7", "type": "some.other.ThirdEnum"}
2735 ]}"#).unwrap();
2736 assert_eq!(schema.named.len(), 4);
2737
2738 if let SchemaPiece::Record { fields, .. } = schema.named[0].clone().piece {
2739 assert_eq!(fields[0].schema, SchemaPieceOrNamed::Named(1)); assert_eq!(fields[1].schema, SchemaPieceOrNamed::Named(1)); assert_eq!(fields[2].schema, SchemaPieceOrNamed::Named(1)); assert_eq!(fields[3].schema, SchemaPieceOrNamed::Named(2)); assert_eq!(fields[4].schema, SchemaPieceOrNamed::Named(2)); assert_eq!(fields[5].schema, SchemaPieceOrNamed::Named(3)); assert_eq!(fields[6].schema, SchemaPieceOrNamed::Named(3)); } else {
2747 panic!("Expected SchemaPiece::Record, found something else");
2748 }
2749
2750 let schema = Schema::from_str(
2751 r#"{"type": "record", "name": "x.Y", "fields": [
2752 {"name": "e", "type":
2753 {"type": "record", "name": "Z", "fields": [
2754 {"name": "f", "type": "x.Y"},
2755 {"name": "g", "type": "x.Z"}
2756 ]}
2757 }
2758 ]}"#,
2759 )
2760 .unwrap();
2761 assert_eq!(schema.named.len(), 2);
2762
2763 if let SchemaPiece::Record { fields, .. } = schema.named[0].clone().piece {
2764 assert_eq!(fields[0].schema, SchemaPieceOrNamed::Named(1)); } else {
2766 panic!("Expected SchemaPiece::Record, found something else");
2767 }
2768
2769 if let SchemaPiece::Record { fields, .. } = schema.named[1].clone().piece {
2770 assert_eq!(fields[0].schema, SchemaPieceOrNamed::Named(0)); assert_eq!(fields[1].schema, SchemaPieceOrNamed::Named(1)); } else {
2773 panic!("Expected SchemaPiece::Record, found something else");
2774 }
2775
2776 let schema = Schema::from_str(
2777 r#"{"type": "record", "name": "R", "fields": [
2778 {"name": "s", "type": {"type": "record", "namespace": "x", "name": "Y", "fields": [
2779 {"name": "e", "type": {"type": "enum", "namespace": "", "name": "Z",
2780 "symbols": ["Foo", "Bar"]}
2781 }
2782 ]}},
2783 {"name": "t", "type": "Z"}
2784 ]}"#,
2785 )
2786 .unwrap();
2787 assert_eq!(schema.named.len(), 3);
2788
2789 if let SchemaPiece::Record { fields, .. } = schema.named[0].clone().piece {
2790 assert_eq!(fields[0].schema, SchemaPieceOrNamed::Named(1)); assert_eq!(fields[1].schema, SchemaPieceOrNamed::Named(2)); } else {
2793 panic!("Expected SchemaPiece::Record, found something else");
2794 }
2795 }
2796
2797 #[test]
2800 fn test_schema_is_send() {
2801 fn send<S: Send>(_s: S) {}
2802
2803 let schema = Schema {
2804 named: vec![],
2805 indices: Default::default(),
2806 top: SchemaPiece::Null.into(),
2807 };
2808 send(schema);
2809 }
2810
2811 #[test]
2812 fn test_schema_is_sync() {
2813 fn sync<S: Sync>(_s: S) {}
2814
2815 let schema = Schema {
2816 named: vec![],
2817 indices: Default::default(),
2818 top: SchemaPiece::Null.into(),
2819 };
2820 sync(&schema);
2821 sync(schema);
2822 }
2823
2824 #[test]
2825 fn test_schema_fingerprint() {
2826 use md5::Md5;
2827 use sha2::Sha256;
2828
2829 let raw_schema = r#"
2830 {
2831 "type": "record",
2832 "name": "test",
2833 "fields": [
2834 {"name": "a", "type": "long", "default": 42},
2835 {"name": "b", "type": "string"}
2836 ]
2837 }
2838 "#;
2839
2840 let schema = Schema::from_str(raw_schema).unwrap();
2841 assert_eq!(
2842 "5ecb2d1f0eaa647d409e6adbd5d70cd274d85802aa9167f5fe3b73ba70b32c76",
2843 format!("{}", schema.fingerprint::<Sha256>())
2844 );
2845
2846 assert_eq!(
2847 "a2c99a3f40ea2eea32593d63b483e962",
2848 format!("{}", schema.fingerprint::<Md5>())
2849 );
2850 }
2851}