1#[cfg(feature = "canonical_extension_types")]
21use arrow_schema::extension::ExtensionType;
22use arrow_schema::{
23 ArrowError, DataType, Field as ArrowField, IntervalUnit, Metadata, Schema as ArrowSchema,
24 TimeUnit, UnionMode,
25};
26use serde::{Deserialize, Serialize};
27use serde_json::{Map as JsonMap, Value, json};
28#[cfg(feature = "sha256")]
29use sha2::{Digest, Sha256};
30use std::borrow::Cow;
31use std::cmp::PartialEq;
32use std::collections::hash_map::Entry;
33use std::collections::{HashMap, HashSet};
34use strum_macros::AsRefStr;
35
36pub const SINGLE_OBJECT_MAGIC: [u8; 2] = [0xC3, 0x01];
38
39pub const CONFLUENT_MAGIC: [u8; 1] = [0x00];
41
42pub const MAX_PREFIX_LEN: usize = 34;
45
46pub const SCHEMA_METADATA_KEY: &str = "avro.schema";
48
49pub const AVRO_ENUM_SYMBOLS_METADATA_KEY: &str = "avro.enum.symbols";
51
52pub const AVRO_FIELD_DEFAULT_METADATA_KEY: &str = "avro.field.default";
54
55pub const AVRO_NAME_METADATA_KEY: &str = "avro.name";
57
58pub const AVRO_NAMESPACE_METADATA_KEY: &str = "avro.namespace";
60
61pub const AVRO_DOC_METADATA_KEY: &str = "avro.doc";
63
64pub const AVRO_ROOT_RECORD_DEFAULT_NAME: &str = "topLevelRecord";
66
67#[derive(Debug, Copy, Clone, PartialEq, Default)]
73pub(crate) enum Nullability {
74 #[default]
76 NullFirst,
77 NullSecond,
79}
80
81impl Nullability {
82 pub(crate) fn non_null_index(&self) -> usize {
84 match self {
85 Nullability::NullFirst => 1,
86 Nullability::NullSecond => 0,
87 }
88 }
89}
90
91#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
95#[serde(untagged)]
96pub(crate) enum TypeName<'a> {
100 Primitive(PrimitiveType),
102 Ref(&'a str),
104}
105
106#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, AsRefStr)]
110#[serde(rename_all = "camelCase")]
111#[strum(serialize_all = "lowercase")]
112pub(crate) enum PrimitiveType {
113 Null,
115 Boolean,
117 Int,
119 Long,
121 Float,
123 Double,
125 Bytes,
127 String,
129}
130
131#[derive(Debug, Clone, PartialEq, Eq, Default, Deserialize, Serialize)]
135#[serde(rename_all = "camelCase")]
136pub(crate) struct Attributes<'a> {
137 #[serde(default)]
141 pub(crate) logical_type: Option<&'a str>,
142
143 #[serde(flatten)]
145 pub(crate) additional: HashMap<&'a str, Value>,
146}
147
148impl Attributes<'_> {
149 pub(crate) fn field_metadata(&self) -> HashMap<String, String> {
151 self.additional
152 .iter()
153 .map(|(k, v)| (k.to_string(), v.to_string()))
154 .collect()
155 }
156}
157
158#[derive(Debug, Clone, PartialEq, Eq, Deserialize, Serialize)]
160#[serde(rename_all = "camelCase")]
161pub(crate) struct Type<'a> {
162 #[serde(borrow)]
164 pub(crate) r#type: TypeName<'a>,
165 #[serde(flatten)]
167 pub(crate) attributes: Attributes<'a>,
168}
169
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
175#[serde(untagged)]
176pub(crate) enum Schema<'a> {
177 #[serde(borrow)]
179 TypeName(TypeName<'a>),
180 #[serde(borrow)]
182 Union(Vec<Schema<'a>>),
183 #[serde(borrow)]
185 Complex(ComplexType<'a>),
186 #[serde(borrow)]
188 Type(Type<'a>),
189}
190
191#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
195#[serde(tag = "type", rename_all = "camelCase")]
196pub(crate) enum ComplexType<'a> {
197 #[serde(borrow)]
199 Record(Record<'a>),
200 #[serde(borrow)]
202 Enum(Enum<'a>),
203 #[serde(borrow)]
205 Array(Array<'a>),
206 #[serde(borrow)]
208 Map(Map<'a>),
209 #[serde(borrow)]
211 Fixed(Fixed<'a>),
212}
213
214#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
218pub(crate) struct Record<'a> {
219 #[serde(borrow)]
221 pub(crate) name: &'a str,
222 #[serde(borrow, default)]
224 pub(crate) namespace: Option<&'a str>,
225 #[serde(borrow, default)]
227 pub(crate) doc: Option<Cow<'a, str>>,
228 #[serde(borrow, default)]
230 pub(crate) aliases: Vec<&'a str>,
231 #[serde(borrow)]
233 pub(crate) fields: Vec<Field<'a>>,
234 #[serde(flatten)]
236 pub(crate) attributes: Attributes<'a>,
237}
238
239fn deserialize_default<'de, D>(deserializer: D) -> Result<Option<Value>, D::Error>
240where
241 D: serde::Deserializer<'de>,
242{
243 Value::deserialize(deserializer).map(Some)
244}
245
246#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
248pub(crate) struct Field<'a> {
249 #[serde(borrow)]
251 pub(crate) name: &'a str,
252 #[serde(borrow, default)]
254 pub(crate) doc: Option<Cow<'a, str>>,
255 #[serde(borrow)]
257 pub(crate) r#type: Schema<'a>,
258 #[serde(deserialize_with = "deserialize_default", default)]
260 pub(crate) default: Option<Value>,
261 #[serde(borrow, default)]
264 pub(crate) aliases: Vec<&'a str>,
265}
266
267#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
271pub(crate) struct Enum<'a> {
272 #[serde(borrow)]
274 pub(crate) name: &'a str,
275 #[serde(borrow, default)]
277 pub(crate) namespace: Option<&'a str>,
278 #[serde(borrow, default)]
280 pub(crate) doc: Option<Cow<'a, str>>,
281 #[serde(borrow, default)]
283 pub(crate) aliases: Vec<&'a str>,
284 #[serde(borrow)]
286 pub(crate) symbols: Vec<&'a str>,
287 #[serde(borrow, default)]
289 pub(crate) default: Option<&'a str>,
290 #[serde(flatten)]
292 pub(crate) attributes: Attributes<'a>,
293}
294
295#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
299pub(crate) struct Array<'a> {
300 #[serde(borrow)]
302 pub(crate) items: Box<Schema<'a>>,
303 #[serde(flatten)]
305 pub(crate) attributes: Attributes<'a>,
306}
307
308#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
312pub(crate) struct Map<'a> {
313 #[serde(borrow)]
315 pub(crate) values: Box<Schema<'a>>,
316 #[serde(flatten)]
318 pub(crate) attributes: Attributes<'a>,
319}
320
321#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
325pub(crate) struct Fixed<'a> {
326 #[serde(borrow)]
328 pub(crate) name: &'a str,
329 #[serde(borrow, default)]
331 pub(crate) namespace: Option<&'a str>,
332 #[serde(borrow, default)]
334 pub(crate) aliases: Vec<&'a str>,
335 pub(crate) size: usize,
337 #[serde(flatten)]
339 pub(crate) attributes: Attributes<'a>,
340}
341
342#[derive(Debug, Copy, Clone, PartialEq, Default)]
343pub(crate) struct AvroSchemaOptions {
344 pub(crate) null_order: Option<Nullability>,
345 pub(crate) strip_metadata: bool,
346}
347
348#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
350pub struct AvroSchema {
351 pub json_string: String,
353}
354
355impl TryFrom<&ArrowSchema> for AvroSchema {
356 type Error = ArrowError;
357
358 fn try_from(schema: &ArrowSchema) -> Result<Self, Self::Error> {
362 AvroSchema::from_arrow_with_options(schema, None)
363 }
364}
365
366impl AvroSchema {
367 pub fn new(json_string: String) -> Self {
369 Self { json_string }
370 }
371
372 pub(crate) fn schema(&self) -> Result<Schema<'_>, ArrowError> {
373 serde_json::from_str(self.json_string.as_str())
374 .map_err(|e| ArrowError::ParseError(format!("Invalid Avro schema JSON: {e}")))
375 }
376
377 pub fn fingerprint(&self, hash_type: FingerprintAlgorithm) -> Result<Fingerprint, ArrowError> {
405 Self::generate_fingerprint(&self.schema()?, hash_type)
406 }
407
408 pub(crate) fn project(&self, projection: &[usize]) -> Result<Self, ArrowError> {
409 let mut value: Value = serde_json::from_str(&self.json_string)
410 .map_err(|e| ArrowError::AvroError(format!("Invalid Avro schema JSON: {e}")))?;
411 let obj = value.as_object_mut().ok_or_else(|| {
412 ArrowError::AvroError(
413 "Projected schema must be a JSON object Avro record schema".to_string(),
414 )
415 })?;
416 match obj.get("type").and_then(|v| v.as_str()) {
417 Some("record") => {}
418 Some(other) => {
419 return Err(ArrowError::AvroError(format!(
420 "Projected schema must be an Avro record, found type '{other}'"
421 )));
422 }
423 None => {
424 return Err(ArrowError::AvroError(
425 "Projected schema missing required 'type' field".to_string(),
426 ));
427 }
428 }
429 let fields_val = obj.get_mut("fields").ok_or_else(|| {
430 ArrowError::AvroError("Avro record schema missing required 'fields'".to_string())
431 })?;
432 let projected_fields = {
433 let mut original_fields = match fields_val {
434 Value::Array(arr) => std::mem::take(arr),
435 _ => {
436 return Err(ArrowError::AvroError(
437 "Avro record schema 'fields' must be an array".to_string(),
438 ));
439 }
440 };
441 let len = original_fields.len();
442 let mut seen: HashSet<usize> = HashSet::with_capacity(projection.len());
443 let mut out: Vec<Value> = Vec::with_capacity(projection.len());
444 for &i in projection {
445 if i >= len {
446 return Err(ArrowError::AvroError(format!(
447 "Projection index {i} out of bounds for record with {len} fields"
448 )));
449 }
450 if !seen.insert(i) {
451 return Err(ArrowError::AvroError(format!(
452 "Duplicate projection index {i}"
453 )));
454 }
455 out.push(std::mem::replace(&mut original_fields[i], Value::Null));
456 }
457 out
458 };
459 *fields_val = Value::Array(projected_fields);
460 let json_string = serde_json::to_string(&value).map_err(|e| {
461 ArrowError::AvroError(format!(
462 "Failed to serialize projected Avro schema JSON: {e}"
463 ))
464 })?;
465 Ok(Self::new(json_string))
466 }
467
468 pub(crate) fn generate_fingerprint(
469 schema: &Schema,
470 hash_type: FingerprintAlgorithm,
471 ) -> Result<Fingerprint, ArrowError> {
472 let canonical = Self::generate_canonical_form(schema).map_err(|e| {
473 ArrowError::ComputeError(format!("Failed to generate canonical form for schema: {e}"))
474 })?;
475 match hash_type {
476 FingerprintAlgorithm::Rabin => {
477 Ok(Fingerprint::Rabin(compute_fingerprint_rabin(&canonical)))
478 }
479 FingerprintAlgorithm::Id | FingerprintAlgorithm::Id64 => Err(ArrowError::SchemaError(
480 "FingerprintAlgorithm of Id or Id64 cannot be used to generate a fingerprint; \
481 if using Fingerprint::Id, pass the registry ID in instead using the set method."
482 .to_string(),
483 )),
484 #[cfg(feature = "md5")]
485 FingerprintAlgorithm::MD5 => Ok(Fingerprint::MD5(compute_fingerprint_md5(&canonical))),
486 #[cfg(feature = "sha256")]
487 FingerprintAlgorithm::SHA256 => {
488 Ok(Fingerprint::SHA256(compute_fingerprint_sha256(&canonical)))
489 }
490 }
491 }
492
493 pub(crate) fn generate_canonical_form(schema: &Schema) -> Result<String, ArrowError> {
504 build_canonical(schema, None)
505 }
506
507 pub(crate) fn from_arrow_with_options(
514 schema: &ArrowSchema,
515 options: Option<AvroSchemaOptions>,
516 ) -> Result<AvroSchema, ArrowError> {
517 let opts = options.unwrap_or_default();
518 let order = opts.null_order.unwrap_or_default();
519 let strip = opts.strip_metadata;
520 if !strip && let Some(json) = schema.metadata.get(SCHEMA_METADATA_KEY) {
521 return Ok(AvroSchema::new(json.clone()));
522 }
523 let mut name_gen = NameGenerator::default();
524 let fields_json = schema
525 .fields()
526 .iter()
527 .map(|f| arrow_field_to_avro(f, &mut name_gen, order, strip))
528 .collect::<Result<Vec<_>, _>>()?;
529 let record_name = schema
530 .metadata
531 .get(AVRO_NAME_METADATA_KEY)
532 .map_or(AVRO_ROOT_RECORD_DEFAULT_NAME, |s| s.as_str());
533 let mut record = JsonMap::with_capacity(schema.metadata.len() + 4);
534 record.insert("type".into(), Value::String("record".into()));
535 record.insert(
536 "name".into(),
537 Value::String(sanitise_avro_name(record_name)),
538 );
539 if let Some(ns) = schema.metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
540 record.insert("namespace".into(), Value::String(ns.clone()));
541 }
542 if let Some(doc) = schema.metadata.get(AVRO_DOC_METADATA_KEY) {
543 record.insert("doc".into(), Value::String(doc.clone()));
544 }
545 record.insert("fields".into(), Value::Array(fields_json));
546 extend_with_passthrough_metadata(&mut record, &schema.metadata);
547 let json_string = serde_json::to_string(&Value::Object(record))
548 .map_err(|e| ArrowError::SchemaError(format!("Serializing Avro JSON failed: {e}")))?;
549 Ok(AvroSchema::new(json_string))
550 }
551}
552
553#[derive(Debug, Copy, Clone)]
555pub(crate) struct Prefix {
556 buf: [u8; MAX_PREFIX_LEN],
557 len: u8,
558}
559
560impl Prefix {
561 #[inline]
562 pub(crate) fn as_slice(&self) -> &[u8] {
563 &self.buf[..self.len as usize]
564 }
565}
566
567#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
569pub enum FingerprintStrategy {
570 #[default]
572 Rabin,
573 Id(u32),
575 Id64(u64),
577 #[cfg(feature = "md5")]
578 MD5,
580 #[cfg(feature = "sha256")]
581 SHA256,
583}
584
585impl From<Fingerprint> for FingerprintStrategy {
586 fn from(f: Fingerprint) -> Self {
587 Self::from(&f)
588 }
589}
590
591impl From<FingerprintAlgorithm> for FingerprintStrategy {
592 fn from(f: FingerprintAlgorithm) -> Self {
593 match f {
594 FingerprintAlgorithm::Rabin => FingerprintStrategy::Rabin,
595 FingerprintAlgorithm::Id => FingerprintStrategy::Id(0),
596 FingerprintAlgorithm::Id64 => FingerprintStrategy::Id64(0),
597 #[cfg(feature = "md5")]
598 FingerprintAlgorithm::MD5 => FingerprintStrategy::MD5,
599 #[cfg(feature = "sha256")]
600 FingerprintAlgorithm::SHA256 => FingerprintStrategy::SHA256,
601 }
602 }
603}
604
605impl From<&Fingerprint> for FingerprintStrategy {
606 fn from(f: &Fingerprint) -> Self {
607 match f {
608 Fingerprint::Rabin(_) => FingerprintStrategy::Rabin,
609 Fingerprint::Id(_) => FingerprintStrategy::Id(0),
610 Fingerprint::Id64(_) => FingerprintStrategy::Id64(0),
611 #[cfg(feature = "md5")]
612 Fingerprint::MD5(_) => FingerprintStrategy::MD5,
613 #[cfg(feature = "sha256")]
614 Fingerprint::SHA256(_) => FingerprintStrategy::SHA256,
615 }
616 }
617}
618
619#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Default)]
622pub enum FingerprintAlgorithm {
623 #[default]
625 Rabin,
626 Id,
628 Id64,
630 #[cfg(feature = "md5")]
631 MD5,
633 #[cfg(feature = "sha256")]
634 SHA256,
636}
637
638impl From<&Fingerprint> for FingerprintAlgorithm {
640 fn from(fp: &Fingerprint) -> Self {
641 match fp {
642 Fingerprint::Rabin(_) => FingerprintAlgorithm::Rabin,
643 Fingerprint::Id(_) => FingerprintAlgorithm::Id,
644 Fingerprint::Id64(_) => FingerprintAlgorithm::Id64,
645 #[cfg(feature = "md5")]
646 Fingerprint::MD5(_) => FingerprintAlgorithm::MD5,
647 #[cfg(feature = "sha256")]
648 Fingerprint::SHA256(_) => FingerprintAlgorithm::SHA256,
649 }
650 }
651}
652
653impl From<FingerprintStrategy> for FingerprintAlgorithm {
654 fn from(s: FingerprintStrategy) -> Self {
655 Self::from(&s)
656 }
657}
658
659impl From<&FingerprintStrategy> for FingerprintAlgorithm {
660 fn from(s: &FingerprintStrategy) -> Self {
661 match s {
662 FingerprintStrategy::Rabin => FingerprintAlgorithm::Rabin,
663 FingerprintStrategy::Id(_) => FingerprintAlgorithm::Id,
664 FingerprintStrategy::Id64(_) => FingerprintAlgorithm::Id64,
665 #[cfg(feature = "md5")]
666 FingerprintStrategy::MD5 => FingerprintAlgorithm::MD5,
667 #[cfg(feature = "sha256")]
668 FingerprintStrategy::SHA256 => FingerprintAlgorithm::SHA256,
669 }
670 }
671}
672
673#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
682pub enum Fingerprint {
683 Rabin(u64),
685 Id(u32),
687 Id64(u64),
689 #[cfg(feature = "md5")]
690 MD5([u8; 16]),
692 #[cfg(feature = "sha256")]
693 SHA256([u8; 32]),
695}
696
697impl From<FingerprintStrategy> for Fingerprint {
698 fn from(s: FingerprintStrategy) -> Self {
699 Self::from(&s)
700 }
701}
702
703impl From<&FingerprintStrategy> for Fingerprint {
704 fn from(s: &FingerprintStrategy) -> Self {
705 match s {
706 FingerprintStrategy::Rabin => Fingerprint::Rabin(0),
707 FingerprintStrategy::Id(id) => Fingerprint::Id(*id),
708 FingerprintStrategy::Id64(id) => Fingerprint::Id64(*id),
709 #[cfg(feature = "md5")]
710 FingerprintStrategy::MD5 => Fingerprint::MD5([0; 16]),
711 #[cfg(feature = "sha256")]
712 FingerprintStrategy::SHA256 => Fingerprint::SHA256([0; 32]),
713 }
714 }
715}
716
717impl From<FingerprintAlgorithm> for Fingerprint {
718 fn from(s: FingerprintAlgorithm) -> Self {
719 match s {
720 FingerprintAlgorithm::Rabin => Fingerprint::Rabin(0),
721 FingerprintAlgorithm::Id => Fingerprint::Id(0),
722 FingerprintAlgorithm::Id64 => Fingerprint::Id64(0),
723 #[cfg(feature = "md5")]
724 FingerprintAlgorithm::MD5 => Fingerprint::MD5([0; 16]),
725 #[cfg(feature = "sha256")]
726 FingerprintAlgorithm::SHA256 => Fingerprint::SHA256([0; 32]),
727 }
728 }
729}
730
731impl Fingerprint {
732 pub fn load_fingerprint_id(id: u32) -> Self {
740 Fingerprint::Id(u32::from_be(id))
741 }
742
743 pub fn load_fingerprint_id64(id: u64) -> Self {
751 Fingerprint::Id64(u64::from_be(id))
752 }
753
754 pub(crate) fn make_prefix(&self) -> Prefix {
778 let mut buf = [0u8; MAX_PREFIX_LEN];
779 let len = match self {
780 Self::Id(val) => write_prefix(&mut buf, &CONFLUENT_MAGIC, &val.to_be_bytes()),
781 Self::Id64(val) => write_prefix(&mut buf, &CONFLUENT_MAGIC, &val.to_be_bytes()),
782 Self::Rabin(val) => write_prefix(&mut buf, &SINGLE_OBJECT_MAGIC, &val.to_le_bytes()),
783 #[cfg(feature = "md5")]
784 Self::MD5(val) => write_prefix(&mut buf, &SINGLE_OBJECT_MAGIC, val),
785 #[cfg(feature = "sha256")]
786 Self::SHA256(val) => write_prefix(&mut buf, &SINGLE_OBJECT_MAGIC, val),
787 };
788 Prefix { buf, len }
789 }
790}
791
792fn write_prefix<const MAGIC_LEN: usize, const PAYLOAD_LEN: usize>(
793 buf: &mut [u8; MAX_PREFIX_LEN],
794 magic: &[u8; MAGIC_LEN],
795 payload: &[u8; PAYLOAD_LEN],
796) -> u8 {
797 debug_assert!(MAGIC_LEN + PAYLOAD_LEN <= MAX_PREFIX_LEN);
798 let total = MAGIC_LEN + PAYLOAD_LEN;
799 let prefix_slice = &mut buf[..total];
800 prefix_slice[..MAGIC_LEN].copy_from_slice(magic);
801 prefix_slice[MAGIC_LEN..total].copy_from_slice(payload);
802 total as u8
803}
804
805#[derive(Debug, Clone, Default)]
832pub struct SchemaStore {
833 fingerprint_algorithm: FingerprintAlgorithm,
835 schemas: HashMap<Fingerprint, AvroSchema>,
837}
838
839impl TryFrom<HashMap<Fingerprint, AvroSchema>> for SchemaStore {
840 type Error = ArrowError;
841
842 fn try_from(schemas: HashMap<Fingerprint, AvroSchema>) -> Result<Self, Self::Error> {
845 Ok(Self {
846 schemas,
847 ..Self::default()
848 })
849 }
850}
851
852impl SchemaStore {
853 pub fn new() -> Self {
855 Self::default()
856 }
857
858 pub fn new_with_type(fingerprint_algorithm: FingerprintAlgorithm) -> Self {
860 Self {
861 fingerprint_algorithm,
862 ..Self::default()
863 }
864 }
865
866 pub fn set(
883 &mut self,
884 fingerprint: Fingerprint,
885 schema: AvroSchema,
886 ) -> Result<Fingerprint, ArrowError> {
887 match self.schemas.entry(fingerprint) {
888 Entry::Occupied(entry) => {
889 if entry.get() != &schema {
890 return Err(ArrowError::ComputeError(format!(
891 "Schema fingerprint collision detected for fingerprint {fingerprint:?}"
892 )));
893 }
894 }
895 Entry::Vacant(entry) => {
896 entry.insert(schema);
897 }
898 }
899 Ok(fingerprint)
900 }
901
902 pub fn register(&mut self, schema: AvroSchema) -> Result<Fingerprint, ArrowError> {
920 if self.fingerprint_algorithm == FingerprintAlgorithm::Id
921 || self.fingerprint_algorithm == FingerprintAlgorithm::Id64
922 {
923 return Err(ArrowError::SchemaError(
924 "Invalid FingerprintAlgorithm; unable to generate fingerprint. \
925 Use the set method directly instead, providing a valid fingerprint"
926 .to_string(),
927 ));
928 }
929 let fingerprint =
930 AvroSchema::generate_fingerprint(&schema.schema()?, self.fingerprint_algorithm)?;
931 self.set(fingerprint, schema)?;
932 Ok(fingerprint)
933 }
934
935 pub fn lookup(&self, fingerprint: &Fingerprint) -> Option<&AvroSchema> {
945 self.schemas.get(fingerprint)
946 }
947
948 pub fn fingerprints(&self) -> Vec<Fingerprint> {
954 self.schemas.keys().copied().collect()
955 }
956
957 pub(crate) fn fingerprint_algorithm(&self) -> FingerprintAlgorithm {
959 self.fingerprint_algorithm
960 }
961}
962
963fn quote(s: &str) -> Result<String, ArrowError> {
964 serde_json::to_string(s)
965 .map_err(|e| ArrowError::ComputeError(format!("Failed to quote string: {e}")))
966}
967
968pub(crate) fn make_full_name(
985 name: &str,
986 namespace_attr: Option<&str>,
987 enclosing_ns: Option<&str>,
988) -> (String, Option<String>) {
989 if let Some((ns, _)) = name.rsplit_once('.') {
991 return (name.to_string(), Some(ns.to_string()));
992 }
993 match namespace_attr.or(enclosing_ns) {
994 Some(ns) => (format!("{ns}.{name}"), Some(ns.to_string())),
995 None => (name.to_string(), None),
996 }
997}
998
999fn build_canonical(schema: &Schema, enclosing_ns: Option<&str>) -> Result<String, ArrowError> {
1000 Ok(match schema {
1001 Schema::TypeName(tn) | Schema::Type(Type { r#type: tn, .. }) => match tn {
1002 TypeName::Primitive(pt) => quote(pt.as_ref())?,
1003 TypeName::Ref(name) => {
1004 let (full_name, _) = make_full_name(name, None, enclosing_ns);
1005 quote(&full_name)?
1006 }
1007 },
1008 Schema::Union(branches) => format!(
1009 "[{}]",
1010 branches
1011 .iter()
1012 .map(|b| build_canonical(b, enclosing_ns))
1013 .collect::<Result<Vec<_>, _>>()?
1014 .join(",")
1015 ),
1016 Schema::Complex(ct) => match ct {
1017 ComplexType::Record(r) => {
1018 let (full_name, child_ns) = make_full_name(r.name, r.namespace, enclosing_ns);
1019 let fields = r
1020 .fields
1021 .iter()
1022 .map(|f| {
1023 let field_type =
1028 build_canonical(&f.r#type, child_ns.as_deref().or(enclosing_ns))?;
1029 Ok(format!(
1030 r#"{{"name":{},"type":{}}}"#,
1031 quote(f.name)?,
1032 field_type
1033 ))
1034 })
1035 .collect::<Result<Vec<_>, ArrowError>>()?
1036 .join(",");
1037 format!(
1038 r#"{{"name":{},"type":"record","fields":[{fields}]}}"#,
1039 quote(&full_name)?,
1040 )
1041 }
1042 ComplexType::Enum(e) => {
1043 let (full_name, _) = make_full_name(e.name, e.namespace, enclosing_ns);
1044 let symbols = e
1045 .symbols
1046 .iter()
1047 .map(|s| quote(s))
1048 .collect::<Result<Vec<_>, _>>()?
1049 .join(",");
1050 format!(
1051 r#"{{"name":{},"type":"enum","symbols":[{symbols}]}}"#,
1052 quote(&full_name)?
1053 )
1054 }
1055 ComplexType::Array(arr) => format!(
1056 r#"{{"type":"array","items":{}}}"#,
1057 build_canonical(&arr.items, enclosing_ns)?
1058 ),
1059 ComplexType::Map(map) => format!(
1060 r#"{{"type":"map","values":{}}}"#,
1061 build_canonical(&map.values, enclosing_ns)?
1062 ),
1063 ComplexType::Fixed(f) => {
1064 let (full_name, _) = make_full_name(f.name, f.namespace, enclosing_ns);
1065 format!(
1066 r#"{{"name":{},"type":"fixed","size":{}}}"#,
1067 quote(&full_name)?,
1068 f.size
1069 )
1070 }
1071 },
1072 })
1073}
1074
1075const EMPTY: u64 = 0xc15d_213a_a4d7_a795;
1077
1078const fn one_entry(i: usize) -> u64 {
1085 let mut fp = i as u64;
1086 let mut j = 0;
1087 while j < 8 {
1088 fp = (fp >> 1) ^ (EMPTY & (0u64.wrapping_sub(fp & 1)));
1089 j += 1;
1090 }
1091 fp
1092}
1093
1094const fn build_table() -> [u64; 256] {
1101 let mut table = [0u64; 256];
1102 let mut i = 0;
1103 while i < 256 {
1104 table[i] = one_entry(i);
1105 i += 1;
1106 }
1107 table
1108}
1109
1110static FINGERPRINT_TABLE: [u64; 256] = build_table();
1112
1113pub(crate) fn compute_fingerprint_rabin(canonical_form: &str) -> u64 {
1116 let mut fp = EMPTY;
1117 for &byte in canonical_form.as_bytes() {
1118 let idx = ((fp as u8) ^ byte) as usize;
1119 fp = (fp >> 8) ^ FINGERPRINT_TABLE[idx];
1120 }
1121 fp
1122}
1123
1124#[cfg(feature = "md5")]
1125#[inline]
1130pub(crate) fn compute_fingerprint_md5(canonical_form: &str) -> [u8; 16] {
1131 let digest = md5::compute(canonical_form.as_bytes());
1132 digest.0
1133}
1134
1135#[cfg(feature = "sha256")]
1136#[inline]
1140pub(crate) fn compute_fingerprint_sha256(canonical_form: &str) -> [u8; 32] {
1141 let mut hasher = Sha256::new();
1142 hasher.update(canonical_form.as_bytes());
1143 let digest = hasher.finalize();
1144 digest.into()
1145}
1146
1147#[inline]
1148fn is_internal_arrow_key(key: &str) -> bool {
1149 key.starts_with("ARROW:") || key == SCHEMA_METADATA_KEY
1150}
1151
1152fn extend_with_passthrough_metadata(target: &mut JsonMap<String, Value>, metadata: &Metadata) {
1157 for (meta_key, meta_val) in metadata {
1158 if meta_key.starts_with("avro.") || is_internal_arrow_key(meta_key) {
1159 continue;
1160 }
1161 let json_val =
1162 serde_json::from_str(meta_val).unwrap_or_else(|_| Value::String(meta_val.clone()));
1163 target.insert(meta_key.clone(), json_val);
1164 }
1165}
1166
1167fn sanitise_avro_name(base_name: &str) -> String {
1169 if base_name.is_empty() {
1170 return "_".to_owned();
1171 }
1172 let mut out: String = base_name
1173 .chars()
1174 .map(|char| {
1175 if char.is_ascii_alphanumeric() || char == '_' {
1176 char
1177 } else {
1178 '_'
1179 }
1180 })
1181 .collect();
1182 if out.as_bytes()[0].is_ascii_digit() {
1183 out.insert(0, '_');
1184 }
1185 out
1186}
1187
1188#[derive(Default)]
1189struct NameGenerator {
1190 used: HashSet<String>,
1191 counters: HashMap<String, usize>,
1192}
1193
1194impl NameGenerator {
1195 fn make_unique(&mut self, field_name: &str) -> String {
1196 let field_name = sanitise_avro_name(field_name);
1197 if self.used.insert(field_name.clone()) {
1198 self.counters.insert(field_name.clone(), 1);
1199 return field_name;
1200 }
1201 let counter = self.counters.entry(field_name.clone()).or_insert(1);
1202 loop {
1203 let candidate = format!("{field_name}_{}", *counter);
1204 if self.used.insert(candidate.clone()) {
1205 return candidate;
1206 }
1207 *counter += 1;
1208 }
1209 }
1210}
1211
1212fn merge_extras(schema: Value, extras: JsonMap<String, Value>) -> Value {
1213 if extras.is_empty() {
1214 return schema;
1215 }
1216 match schema {
1217 Value::Object(mut map) => {
1218 map.extend(extras);
1219 Value::Object(map)
1220 }
1221 Value::Array(mut union) => {
1222 if let Some(non_null) = union.iter_mut().find(|val| val.as_str() != Some("null")) {
1225 let original = std::mem::take(non_null);
1226 *non_null = merge_extras(original, extras);
1227 }
1228 Value::Array(union)
1229 }
1230 primitive => {
1231 let mut map = JsonMap::with_capacity(extras.len() + 1);
1232 map.insert("type".into(), primitive);
1233 map.extend(extras);
1234 Value::Object(map)
1235 }
1236 }
1237}
1238
1239#[inline]
1240fn is_avro_json_null(v: &Value) -> bool {
1241 matches!(v, Value::String(s) if s == "null")
1242}
1243
1244fn wrap_nullable(inner: Value, null_order: Nullability) -> Value {
1245 let null = Value::String("null".into());
1246 match inner {
1247 Value::Array(mut union) => {
1248 if union.iter().any(is_avro_json_null) {
1253 return Value::Array(union);
1254 }
1255 match null_order {
1257 Nullability::NullFirst => union.insert(0, null),
1258 Nullability::NullSecond => union.push(null),
1259 }
1260 Value::Array(union)
1261 }
1262 other => match null_order {
1263 Nullability::NullFirst => Value::Array(vec![null, other]),
1264 Nullability::NullSecond => Value::Array(vec![other, null]),
1265 },
1266 }
1267}
1268
1269fn min_fixed_bytes_for_precision(p: usize) -> usize {
1270 const MAX_P: [usize; 32] = [
1273 2, 4, 6, 9, 11, 14, 16, 18, 21, 23, 26, 28, 31, 33, 35, 38, 40, 43, 45, 47, 50, 52, 55, 57,
1274 59, 62, 64, 67, 69, 71, 74, 76,
1275 ];
1276 for (i, &max_p) in MAX_P.iter().enumerate() {
1277 if p <= max_p {
1278 return i + 1;
1279 }
1280 }
1281 32 }
1283
1284fn union_branch_signature(branch: &Value) -> Result<String, ArrowError> {
1285 match branch {
1286 Value::String(t) => Ok(format!("P:{t}")),
1287 Value::Object(map) => {
1288 let t = map.get("type").and_then(|v| v.as_str()).ok_or_else(|| {
1289 ArrowError::SchemaError("Union branch object missing string 'type'".into())
1290 })?;
1291 match t {
1292 "record" | "enum" | "fixed" => {
1293 let name = map.get("name").and_then(|v| v.as_str()).ok_or_else(|| {
1294 ArrowError::SchemaError(format!(
1295 "Union branch '{t}' missing required 'name'"
1296 ))
1297 })?;
1298 Ok(format!("N:{t}:{name}"))
1299 }
1300 "array" | "map" => Ok(format!("C:{t}")),
1301 other => Ok(format!("P:{other}")),
1302 }
1303 }
1304 Value::Array(_) => Err(ArrowError::SchemaError(
1305 "Avro union may not immediately contain another union".into(),
1306 )),
1307 _ => Err(ArrowError::SchemaError(
1308 "Invalid JSON for Avro union branch".into(),
1309 )),
1310 }
1311}
1312
1313fn datatype_to_avro(
1314 dt: &DataType,
1315 field_name: &str,
1316 metadata: &Metadata,
1317 name_gen: &mut NameGenerator,
1318 null_order: Nullability,
1319 strip: bool,
1320) -> Result<(Value, JsonMap<String, Value>), ArrowError> {
1321 let mut extras = JsonMap::new();
1322 let mut handle_decimal = |precision: &u8, scale: &i8| -> Result<Value, ArrowError> {
1323 if *scale < 0 {
1324 return Err(ArrowError::SchemaError(format!(
1325 "Invalid Avro decimal for field '{field_name}': scale ({scale}) must be >= 0"
1326 )));
1327 }
1328 if (*scale as usize) > (*precision as usize) {
1329 return Err(ArrowError::SchemaError(format!(
1330 "Invalid Avro decimal for field '{field_name}': scale ({scale}) \
1331 must be <= precision ({precision})"
1332 )));
1333 }
1334 let mut meta = JsonMap::from_iter([
1335 ("logicalType".into(), json!("decimal")),
1336 ("precision".into(), json!(*precision)),
1337 ("scale".into(), json!(*scale)),
1338 ]);
1339 let mut fixed_size = metadata.get("size").and_then(|v| v.parse::<usize>().ok());
1340 let carries_name = metadata.contains_key(AVRO_NAME_METADATA_KEY)
1341 || metadata.contains_key(AVRO_NAMESPACE_METADATA_KEY);
1342 if fixed_size.is_none() && carries_name {
1343 fixed_size = Some(min_fixed_bytes_for_precision(*precision as usize));
1344 }
1345 if let Some(size) = fixed_size {
1346 meta.insert("type".into(), json!("fixed"));
1347 meta.insert("size".into(), json!(size));
1348 let chosen_name = metadata
1349 .get(AVRO_NAME_METADATA_KEY)
1350 .map(|s| sanitise_avro_name(s))
1351 .unwrap_or_else(|| name_gen.make_unique(field_name));
1352 meta.insert("name".into(), json!(chosen_name));
1353 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1354 meta.insert("namespace".into(), json!(ns));
1355 }
1356 } else {
1357 meta.insert("type".into(), json!("bytes"));
1359 }
1360 Ok(Value::Object(meta))
1361 };
1362 let val = match dt {
1363 DataType::Null => Value::String("null".into()),
1364 DataType::Boolean => Value::String("boolean".into()),
1365 #[cfg(not(feature = "avro_custom_types"))]
1366 DataType::Int8 | DataType::Int16 | DataType::UInt8 | DataType::UInt16 => {
1367 Value::String("int".into())
1368 }
1369 DataType::Int32 => Value::String("int".into()),
1370 #[cfg(feature = "avro_custom_types")]
1371 DataType::Int8 => json!({ "type": "int", "logicalType": "arrow.int8" }),
1372 #[cfg(feature = "avro_custom_types")]
1373 DataType::Int16 => json!({ "type": "int", "logicalType": "arrow.int16" }),
1374 #[cfg(feature = "avro_custom_types")]
1375 DataType::UInt8 => json!({ "type": "int", "logicalType": "arrow.uint8" }),
1376 #[cfg(feature = "avro_custom_types")]
1377 DataType::UInt16 => json!({ "type": "int", "logicalType": "arrow.uint16" }),
1378 #[cfg(not(feature = "avro_custom_types"))]
1379 DataType::UInt32 => Value::String("long".into()),
1380 #[cfg(feature = "avro_custom_types")]
1381 DataType::UInt32 => json!({ "type": "long", "logicalType": "arrow.uint32" }),
1382 DataType::Int64 => Value::String("long".into()),
1383 #[cfg(not(feature = "avro_custom_types"))]
1384 DataType::UInt64 => Value::String("long".into()),
1385 #[cfg(feature = "avro_custom_types")]
1386 DataType::UInt64 => {
1387 let chosen_name = metadata
1389 .get(AVRO_NAME_METADATA_KEY)
1390 .map(|s| sanitise_avro_name(s))
1391 .unwrap_or_else(|| name_gen.make_unique(field_name));
1392 let mut obj = JsonMap::from_iter([
1393 ("type".into(), json!("fixed")),
1394 ("name".into(), json!(chosen_name)),
1395 ("size".into(), json!(8)),
1396 ("logicalType".into(), json!("arrow.uint64")),
1397 ]);
1398 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1399 obj.insert("namespace".into(), json!(ns));
1400 }
1401 json!(obj)
1402 }
1403 #[cfg(not(feature = "avro_custom_types"))]
1404 DataType::Float16 => Value::String("float".into()),
1405 #[cfg(feature = "avro_custom_types")]
1406 DataType::Float16 => {
1407 let chosen_name = metadata
1409 .get(AVRO_NAME_METADATA_KEY)
1410 .map(|s| sanitise_avro_name(s))
1411 .unwrap_or_else(|| name_gen.make_unique(field_name));
1412 let mut obj = JsonMap::from_iter([
1413 ("type".into(), json!("fixed")),
1414 ("name".into(), json!(chosen_name)),
1415 ("size".into(), json!(2)),
1416 ("logicalType".into(), json!("arrow.float16")),
1417 ]);
1418 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1419 obj.insert("namespace".into(), json!(ns));
1420 }
1421 json!(obj)
1422 }
1423 DataType::Float32 => Value::String("float".into()),
1424 DataType::Float64 => Value::String("double".into()),
1425 DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => Value::String("string".into()),
1426 DataType::Binary | DataType::LargeBinary => Value::String("bytes".into()),
1427 DataType::BinaryView => {
1428 if !strip {
1429 extras.insert("arrowBinaryView".into(), Value::Bool(true));
1430 }
1431 Value::String("bytes".into())
1432 }
1433 DataType::FixedSizeBinary(len) => {
1434 let md_is_uuid = metadata
1435 .get("logicalType")
1436 .map(|s| s.trim_matches('"') == "uuid")
1437 .unwrap_or(false);
1438 #[cfg(feature = "canonical_extension_types")]
1439 let ext_is_uuid = metadata
1440 .get(arrow_schema::extension::EXTENSION_TYPE_NAME_KEY)
1441 .map(|v| v == arrow_schema::extension::Uuid::NAME || v == "uuid")
1442 .unwrap_or(false);
1443 #[cfg(not(feature = "canonical_extension_types"))]
1444 let ext_is_uuid = false;
1445 let is_uuid = (*len == 16) && (md_is_uuid || ext_is_uuid);
1446 if is_uuid {
1447 json!({ "type": "string", "logicalType": "uuid" })
1448 } else {
1449 let chosen_name = metadata
1450 .get(AVRO_NAME_METADATA_KEY)
1451 .map(|s| sanitise_avro_name(s))
1452 .unwrap_or_else(|| name_gen.make_unique(field_name));
1453 let mut obj = JsonMap::from_iter([
1454 ("type".into(), json!("fixed")),
1455 ("name".into(), json!(chosen_name)),
1456 ("size".into(), json!(len)),
1457 ]);
1458 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1459 obj.insert("namespace".into(), json!(ns));
1460 }
1461 Value::Object(obj)
1462 }
1463 }
1464 #[cfg(feature = "small_decimals")]
1465 DataType::Decimal32(precision, scale) | DataType::Decimal64(precision, scale) => {
1466 handle_decimal(precision, scale)?
1467 }
1468 DataType::Decimal128(precision, scale) | DataType::Decimal256(precision, scale) => {
1469 handle_decimal(precision, scale)?
1470 }
1471 DataType::Date32 => json!({ "type": "int", "logicalType": "date" }),
1472 #[cfg(not(feature = "avro_custom_types"))]
1473 DataType::Date64 => json!({ "type": "long", "logicalType": "local-timestamp-millis" }),
1474 #[cfg(feature = "avro_custom_types")]
1475 DataType::Date64 => json!({ "type": "long", "logicalType": "arrow.date64" }),
1476 DataType::Time32(unit) => match unit {
1477 TimeUnit::Millisecond => json!({ "type": "int", "logicalType": "time-millis" }),
1478 #[cfg(not(feature = "avro_custom_types"))]
1479 TimeUnit::Second => {
1480 if !strip {
1482 extras.insert("arrowTimeUnit".into(), Value::String("second".into()));
1483 }
1484 json!({ "type": "int", "logicalType": "time-millis" })
1485 }
1486 #[cfg(feature = "avro_custom_types")]
1487 TimeUnit::Second => {
1488 json!({ "type": "int", "logicalType": "arrow.time32-second" })
1489 }
1490 _ => Value::String("int".into()),
1491 },
1492 DataType::Time64(unit) => match unit {
1493 TimeUnit::Microsecond => json!({ "type": "long", "logicalType": "time-micros" }),
1494 #[cfg(not(feature = "avro_custom_types"))]
1495 TimeUnit::Nanosecond => {
1496 if !strip {
1498 extras.insert("arrowTimeUnit".into(), Value::String("nanosecond".into()));
1499 }
1500 json!({ "type": "long", "logicalType": "time-micros" })
1501 }
1502 #[cfg(feature = "avro_custom_types")]
1503 TimeUnit::Nanosecond => {
1504 json!({ "type": "long", "logicalType": "arrow.time64-nanosecond" })
1505 }
1506 _ => Value::String("long".into()),
1507 },
1508 DataType::Timestamp(unit, tz) => {
1509 #[cfg(feature = "avro_custom_types")]
1510 if matches!(unit, TimeUnit::Second) {
1511 let logical_type = if tz.is_some() {
1512 "arrow.timestamp-second"
1513 } else {
1514 "arrow.local-timestamp-second"
1515 };
1516 return Ok((
1517 json!({ "type": "long", "logicalType": logical_type }),
1518 extras,
1519 ));
1520 }
1521 let logical_type = match (unit, tz.is_some()) {
1522 (TimeUnit::Millisecond, true) => "timestamp-millis",
1523 (TimeUnit::Millisecond, false) => "local-timestamp-millis",
1524 (TimeUnit::Microsecond, true) => "timestamp-micros",
1525 (TimeUnit::Microsecond, false) => "local-timestamp-micros",
1526 (TimeUnit::Nanosecond, true) => "timestamp-nanos",
1527 (TimeUnit::Nanosecond, false) => "local-timestamp-nanos",
1528 (TimeUnit::Second, has_tz) => {
1529 if !strip {
1531 extras.insert("arrowTimeUnit".into(), Value::String("second".into()));
1532 }
1533 let ts_logical_type = if has_tz {
1534 "timestamp-millis"
1535 } else {
1536 "local-timestamp-millis"
1537 };
1538 return Ok((
1539 json!({ "type": "long", "logicalType": ts_logical_type }),
1540 extras,
1541 ));
1542 }
1543 };
1544 if !strip && matches!(unit, TimeUnit::Nanosecond) {
1545 extras.insert("arrowTimeUnit".into(), Value::String("nanosecond".into()));
1546 }
1547 json!({ "type": "long", "logicalType": logical_type })
1548 }
1549 #[cfg(not(feature = "avro_custom_types"))]
1550 DataType::Duration(_unit) => Value::String("long".into()),
1551 #[cfg(feature = "avro_custom_types")]
1552 DataType::Duration(unit) => {
1553 let logical_type = match unit {
1556 TimeUnit::Second => "arrow.duration-seconds",
1557 TimeUnit::Millisecond => "arrow.duration-millis",
1558 TimeUnit::Microsecond => "arrow.duration-micros",
1559 TimeUnit::Nanosecond => "arrow.duration-nanos",
1560 };
1561 json!({ "type": "long", "logicalType": logical_type })
1562 }
1563 #[cfg(not(feature = "avro_custom_types"))]
1564 DataType::Interval(IntervalUnit::MonthDayNano) => {
1565 let chosen_name = metadata
1567 .get(AVRO_NAME_METADATA_KEY)
1568 .map(|s| sanitise_avro_name(s))
1569 .unwrap_or_else(|| name_gen.make_unique(field_name));
1570 let mut obj = JsonMap::from_iter([
1571 ("type".into(), json!("fixed")),
1572 ("name".into(), json!(chosen_name)),
1573 ("size".into(), json!(12)),
1574 ("logicalType".into(), json!("duration")),
1575 ]);
1576 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1577 obj.insert("namespace".into(), json!(ns));
1578 }
1579 json!(obj)
1580 }
1581 #[cfg(feature = "avro_custom_types")]
1582 DataType::Interval(IntervalUnit::MonthDayNano) => {
1583 let chosen_name = metadata
1586 .get(AVRO_NAME_METADATA_KEY)
1587 .map(|s| sanitise_avro_name(s))
1588 .unwrap_or_else(|| name_gen.make_unique(field_name));
1589 let mut obj = JsonMap::from_iter([
1590 ("type".into(), json!("fixed")),
1591 ("name".into(), json!(chosen_name)),
1592 ("size".into(), json!(16)),
1593 ("logicalType".into(), json!("arrow.interval-month-day-nano")),
1594 ]);
1595 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1596 obj.insert("namespace".into(), json!(ns));
1597 }
1598 json!(obj)
1599 }
1600 #[cfg(not(feature = "avro_custom_types"))]
1601 DataType::Interval(IntervalUnit::YearMonth) => {
1602 let chosen_name = metadata
1604 .get(AVRO_NAME_METADATA_KEY)
1605 .map(|s| sanitise_avro_name(s))
1606 .unwrap_or_else(|| name_gen.make_unique(field_name));
1607 let mut extras = JsonMap::from_iter([
1608 ("type".into(), json!("fixed")),
1609 ("name".into(), json!(chosen_name)),
1610 ("size".into(), json!(12)),
1611 ("logicalType".into(), json!("duration")),
1612 ]);
1613 if !strip {
1614 extras.insert(
1615 "arrowIntervalUnit".into(),
1616 Value::String("yearmonth".into()),
1617 );
1618 }
1619 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1620 extras.insert("namespace".into(), json!(ns));
1621 }
1622 json!(extras)
1623 }
1624 #[cfg(feature = "avro_custom_types")]
1625 DataType::Interval(IntervalUnit::YearMonth) => {
1626 let chosen_name = metadata
1627 .get(AVRO_NAME_METADATA_KEY)
1628 .map(|s| sanitise_avro_name(s))
1629 .unwrap_or_else(|| name_gen.make_unique(field_name));
1630 let mut obj = JsonMap::from_iter([
1631 ("type".into(), json!("fixed")),
1632 ("name".into(), json!(chosen_name)),
1633 ("size".into(), json!(4)),
1634 ("logicalType".into(), json!("arrow.interval-year-month")),
1635 ]);
1636 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1637 obj.insert("namespace".into(), json!(ns));
1638 }
1639 json!(obj)
1640 }
1641 #[cfg(not(feature = "avro_custom_types"))]
1642 DataType::Interval(IntervalUnit::DayTime) => {
1643 let chosen_name = metadata
1645 .get(AVRO_NAME_METADATA_KEY)
1646 .map(|s| sanitise_avro_name(s))
1647 .unwrap_or_else(|| name_gen.make_unique(field_name));
1648 let mut obj = JsonMap::from_iter([
1649 ("type".into(), json!("fixed")),
1650 ("name".into(), json!(chosen_name)),
1651 ("size".into(), json!(12)),
1652 ("logicalType".into(), json!("duration")),
1653 ]);
1654 if !strip {
1655 obj.insert("arrowIntervalUnit".into(), Value::String("daytime".into()));
1656 }
1657 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1658 obj.insert("namespace".into(), json!(ns));
1659 }
1660 json!(obj)
1661 }
1662 #[cfg(feature = "avro_custom_types")]
1663 DataType::Interval(IntervalUnit::DayTime) => {
1664 let chosen_name = metadata
1665 .get(AVRO_NAME_METADATA_KEY)
1666 .map(|s| sanitise_avro_name(s))
1667 .unwrap_or_else(|| name_gen.make_unique(field_name));
1668 let mut obj = JsonMap::from_iter([
1669 ("type".into(), json!("fixed")),
1670 ("name".into(), json!(chosen_name)),
1671 ("size".into(), json!(8)),
1672 ("logicalType".into(), json!("arrow.interval-day-time")),
1673 ]);
1674 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1675 obj.insert("namespace".into(), json!(ns));
1676 }
1677 json!(obj)
1678 }
1679 DataType::List(child) | DataType::LargeList(child) => {
1680 if matches!(dt, DataType::LargeList(_)) && !strip {
1681 extras.insert("arrowLargeList".into(), Value::Bool(true));
1682 }
1683 let items_schema = process_datatype(
1684 child.data_type(),
1685 child.name(),
1686 child.metadata(),
1687 name_gen,
1688 null_order,
1689 child.is_nullable(),
1690 strip,
1691 )?;
1692 json!({
1693 "type": "array",
1694 "items": items_schema
1695 })
1696 }
1697 DataType::ListView(child) | DataType::LargeListView(child) => {
1698 if matches!(dt, DataType::LargeListView(_)) && !strip {
1699 extras.insert("arrowLargeList".into(), Value::Bool(true));
1700 }
1701 if !strip {
1702 extras.insert("arrowListView".into(), Value::Bool(true));
1703 }
1704 let items_schema = process_datatype(
1705 child.data_type(),
1706 child.name(),
1707 child.metadata(),
1708 name_gen,
1709 null_order,
1710 child.is_nullable(),
1711 strip,
1712 )?;
1713 json!({
1714 "type": "array",
1715 "items": items_schema
1716 })
1717 }
1718 DataType::FixedSizeList(child, len) => {
1719 if !strip {
1720 extras.insert("arrowFixedSize".into(), json!(len));
1721 }
1722 let items_schema = process_datatype(
1723 child.data_type(),
1724 child.name(),
1725 child.metadata(),
1726 name_gen,
1727 null_order,
1728 child.is_nullable(),
1729 strip,
1730 )?;
1731 json!({
1732 "type": "array",
1733 "items": items_schema
1734 })
1735 }
1736 DataType::Map(entries, _) => {
1737 let value_field = match entries.data_type() {
1738 DataType::Struct(fs) => &fs[1],
1739 _ => {
1740 return Err(ArrowError::SchemaError(
1741 "Map 'entries' field must be Struct(key,value)".into(),
1742 ));
1743 }
1744 };
1745 let values_schema = process_datatype(
1746 value_field.data_type(),
1747 value_field.name(),
1748 value_field.metadata(),
1749 name_gen,
1750 null_order,
1751 value_field.is_nullable(),
1752 strip,
1753 )?;
1754 json!({
1755 "type": "map",
1756 "values": values_schema
1757 })
1758 }
1759 DataType::Struct(fields) => {
1760 let avro_fields = fields
1761 .iter()
1762 .map(|field| arrow_field_to_avro(field, name_gen, null_order, strip))
1763 .collect::<Result<Vec<_>, _>>()?;
1764 let chosen_name = metadata
1766 .get(AVRO_NAME_METADATA_KEY)
1767 .map(|s| sanitise_avro_name(s))
1768 .unwrap_or_else(|| name_gen.make_unique(field_name));
1769 let mut obj = JsonMap::from_iter([
1770 ("type".into(), json!("record")),
1771 ("name".into(), json!(chosen_name)),
1772 ("fields".into(), Value::Array(avro_fields)),
1773 ]);
1774 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1775 obj.insert("namespace".into(), json!(ns));
1776 }
1777 Value::Object(obj)
1778 }
1779 DataType::Dictionary(_, value) => {
1780 if let Some(j) = metadata.get(AVRO_ENUM_SYMBOLS_METADATA_KEY) {
1781 let symbols: Vec<&str> =
1782 serde_json::from_str(j).map_err(|e| ArrowError::ParseError(e.to_string()))?;
1783 let chosen_name = metadata
1785 .get(AVRO_NAME_METADATA_KEY)
1786 .map(|s| sanitise_avro_name(s))
1787 .unwrap_or_else(|| name_gen.make_unique(field_name));
1788 let mut obj = JsonMap::from_iter([
1789 ("type".into(), json!("enum")),
1790 ("name".into(), json!(chosen_name)),
1791 ("symbols".into(), json!(symbols)),
1792 ]);
1793 if let Some(ns) = metadata.get(AVRO_NAMESPACE_METADATA_KEY) {
1794 obj.insert("namespace".into(), json!(ns));
1795 }
1796 Value::Object(obj)
1797 } else {
1798 process_datatype(
1799 value.as_ref(),
1800 field_name,
1801 metadata,
1802 name_gen,
1803 null_order,
1804 false,
1805 strip,
1806 )?
1807 }
1808 }
1809 #[cfg(feature = "avro_custom_types")]
1810 DataType::RunEndEncoded(run_ends, values) => {
1811 let bits = match run_ends.data_type() {
1812 DataType::Int16 => 16,
1813 DataType::Int32 => 32,
1814 DataType::Int64 => 64,
1815 other => {
1816 return Err(ArrowError::SchemaError(format!(
1817 "RunEndEncoded requires Int16/Int32/Int64 for run_ends, found: {other:?}"
1818 )));
1819 }
1820 };
1821 let (value_schema, value_extras) = datatype_to_avro(
1823 values.data_type(),
1824 values.name(),
1825 values.metadata(),
1826 name_gen,
1827 null_order,
1828 strip,
1829 )?;
1830 let mut merged = merge_extras(value_schema, value_extras);
1831 if values.is_nullable() {
1832 merged = wrap_nullable(merged, null_order);
1833 }
1834 let mut extras = JsonMap::new();
1835 extras.insert("logicalType".into(), json!("arrow.run-end-encoded"));
1836 extras.insert("arrow.runEndIndexBits".into(), json!(bits));
1837 return Ok((merged, extras));
1838 }
1839 #[cfg(not(feature = "avro_custom_types"))]
1840 DataType::RunEndEncoded(_run_ends, values) => {
1841 let (value_schema, _extras) = datatype_to_avro(
1842 values.data_type(),
1843 values.name(),
1844 values.metadata(),
1845 name_gen,
1846 null_order,
1847 strip,
1848 )?;
1849 return Ok((value_schema, JsonMap::new()));
1850 }
1851 DataType::Union(fields, mode) => {
1852 let mut branches: Vec<Value> = Vec::with_capacity(fields.len());
1853 let mut type_ids: Vec<i32> = Vec::with_capacity(fields.len());
1854 for (type_id, field_ref) in fields.iter() {
1855 let (branch_schema, _branch_extras) = datatype_to_avro(
1857 field_ref.data_type(),
1858 field_ref.name(),
1859 field_ref.metadata(),
1860 name_gen,
1861 null_order,
1862 strip,
1863 )?;
1864 if matches!(branch_schema, Value::Array(_)) {
1866 return Err(ArrowError::SchemaError(
1867 "Avro union may not immediately contain another union".into(),
1868 ));
1869 }
1870 branches.push(branch_schema);
1871 type_ids.push(type_id as i32);
1872 }
1873 let mut seen: HashSet<String> = HashSet::with_capacity(branches.len());
1874 for b in &branches {
1875 let sig = union_branch_signature(b)?;
1876 if !seen.insert(sig) {
1877 return Err(ArrowError::SchemaError(
1878 "Avro union contains duplicate branch types (disallowed by spec)".into(),
1879 ));
1880 }
1881 }
1882 if !strip {
1883 extras.insert(
1884 "arrowUnionMode".into(),
1885 Value::String(
1886 match mode {
1887 UnionMode::Sparse => "sparse",
1888 UnionMode::Dense => "dense",
1889 }
1890 .to_string(),
1891 ),
1892 );
1893 extras.insert(
1894 "arrowUnionTypeIds".into(),
1895 Value::Array(type_ids.into_iter().map(|id| json!(id)).collect()),
1896 );
1897 }
1898 Value::Array(branches)
1899 }
1900 #[cfg(not(feature = "small_decimals"))]
1901 other => {
1902 return Err(ArrowError::NotYetImplemented(format!(
1903 "Arrow type {other:?} has no Avro representation"
1904 )));
1905 }
1906 };
1907 Ok((val, extras))
1908}
1909
1910fn process_datatype(
1911 dt: &DataType,
1912 field_name: &str,
1913 metadata: &Metadata,
1914 name_gen: &mut NameGenerator,
1915 null_order: Nullability,
1916 is_nullable: bool,
1917 strip: bool,
1918) -> Result<Value, ArrowError> {
1919 let (schema, extras) = datatype_to_avro(dt, field_name, metadata, name_gen, null_order, strip)?;
1920 let mut merged = merge_extras(schema, extras);
1921 if is_nullable {
1922 merged = wrap_nullable(merged, null_order)
1923 }
1924 Ok(merged)
1925}
1926
1927fn arrow_field_to_avro(
1928 field: &ArrowField,
1929 name_gen: &mut NameGenerator,
1930 null_order: Nullability,
1931 strip: bool,
1932) -> Result<Value, ArrowError> {
1933 let avro_name = sanitise_avro_name(field.name());
1934 let schema_value = process_datatype(
1935 field.data_type(),
1936 &avro_name,
1937 field.metadata(),
1938 name_gen,
1939 null_order,
1940 field.is_nullable(),
1941 strip,
1942 )?;
1943 let mut map = JsonMap::with_capacity(field.metadata().len() + 3);
1945 map.insert("name".into(), Value::String(avro_name));
1946 map.insert("type".into(), schema_value);
1947 for (meta_key, meta_val) in field.metadata() {
1949 if is_internal_arrow_key(meta_key) {
1950 continue;
1951 }
1952 match meta_key.as_str() {
1953 AVRO_DOC_METADATA_KEY => {
1954 map.insert("doc".into(), Value::String(meta_val.clone()));
1955 }
1956 AVRO_FIELD_DEFAULT_METADATA_KEY => {
1957 let default_value = serde_json::from_str(meta_val)
1958 .unwrap_or_else(|_| Value::String(meta_val.clone()));
1959 map.insert("default".into(), default_value);
1960 }
1961 _ => {
1962 let json_val = serde_json::from_str(meta_val)
1963 .unwrap_or_else(|_| Value::String(meta_val.clone()));
1964 map.insert(meta_key.clone(), json_val);
1965 }
1966 }
1967 }
1968 Ok(Value::Object(map))
1969}
1970
1971#[cfg(test)]
1972mod tests {
1973 use super::*;
1974 use crate::codec::{AvroField, AvroFieldBuilder};
1975 use arrow_schema::{DataType, Fields, SchemaBuilder, TimeUnit, UnionFields};
1976 use serde_json::json;
1977 use std::sync::Arc;
1978
1979 fn int_schema() -> Schema<'static> {
1980 Schema::TypeName(TypeName::Primitive(PrimitiveType::Int))
1981 }
1982
1983 fn record_schema() -> Schema<'static> {
1984 Schema::Complex(ComplexType::Record(Record {
1985 name: "record1",
1986 namespace: Some("test.namespace"),
1987 doc: Some(Cow::from("A test record")),
1988 aliases: vec![],
1989 fields: vec![
1990 Field {
1991 name: "field1",
1992 doc: Some(Cow::from("An integer field")),
1993 r#type: int_schema(),
1994 default: None,
1995 aliases: vec![],
1996 },
1997 Field {
1998 name: "field2",
1999 doc: None,
2000 r#type: Schema::TypeName(TypeName::Primitive(PrimitiveType::String)),
2001 default: None,
2002 aliases: vec![],
2003 },
2004 ],
2005 attributes: Attributes::default(),
2006 }))
2007 }
2008
2009 fn single_field_schema(field: ArrowField) -> arrow_schema::Schema {
2010 let mut sb = SchemaBuilder::new();
2011 sb.push(field);
2012 sb.finish()
2013 }
2014
2015 fn assert_json_contains(avro_json: &str, needle: &str) {
2016 assert!(
2017 avro_json.contains(needle),
2018 "JSON did not contain `{needle}` : {avro_json}"
2019 )
2020 }
2021
2022 #[test]
2023 fn test_deserialize() {
2024 let t: Schema = serde_json::from_str("\"string\"").unwrap();
2025 assert_eq!(
2026 t,
2027 Schema::TypeName(TypeName::Primitive(PrimitiveType::String))
2028 );
2029
2030 let t: Schema = serde_json::from_str("[\"int\", \"null\"]").unwrap();
2031 assert_eq!(
2032 t,
2033 Schema::Union(vec![
2034 Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)),
2035 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2036 ])
2037 );
2038
2039 let t: Type = serde_json::from_str(
2040 r#"{
2041 "type":"long",
2042 "logicalType":"timestamp-micros"
2043 }"#,
2044 )
2045 .unwrap();
2046
2047 let timestamp = Type {
2048 r#type: TypeName::Primitive(PrimitiveType::Long),
2049 attributes: Attributes {
2050 logical_type: Some("timestamp-micros"),
2051 additional: Default::default(),
2052 },
2053 };
2054
2055 assert_eq!(t, timestamp);
2056
2057 let t: ComplexType = serde_json::from_str(
2058 r#"{
2059 "type":"fixed",
2060 "name":"fixed",
2061 "namespace":"topLevelRecord.value",
2062 "size":11,
2063 "logicalType":"decimal",
2064 "precision":25,
2065 "scale":2
2066 }"#,
2067 )
2068 .unwrap();
2069
2070 let decimal = ComplexType::Fixed(Fixed {
2071 name: "fixed",
2072 namespace: Some("topLevelRecord.value"),
2073 aliases: vec![],
2074 size: 11,
2075 attributes: Attributes {
2076 logical_type: Some("decimal"),
2077 additional: vec![("precision", json!(25)), ("scale", json!(2))]
2078 .into_iter()
2079 .collect(),
2080 },
2081 });
2082
2083 assert_eq!(t, decimal);
2084
2085 let schema: Schema = serde_json::from_str(
2086 r#"{
2087 "type":"record",
2088 "name":"topLevelRecord",
2089 "fields":[
2090 {
2091 "name":"value",
2092 "type":[
2093 {
2094 "type":"fixed",
2095 "name":"fixed",
2096 "namespace":"topLevelRecord.value",
2097 "size":11,
2098 "logicalType":"decimal",
2099 "precision":25,
2100 "scale":2
2101 },
2102 "null"
2103 ]
2104 }
2105 ]
2106 }"#,
2107 )
2108 .unwrap();
2109
2110 assert_eq!(
2111 schema,
2112 Schema::Complex(ComplexType::Record(Record {
2113 name: "topLevelRecord",
2114 namespace: None,
2115 doc: None,
2116 aliases: vec![],
2117 fields: vec![Field {
2118 name: "value",
2119 doc: None,
2120 r#type: Schema::Union(vec![
2121 Schema::Complex(decimal),
2122 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2123 ]),
2124 default: None,
2125 aliases: vec![],
2126 },],
2127 attributes: Default::default(),
2128 }))
2129 );
2130
2131 let schema: Schema = serde_json::from_str(
2132 r#"{
2133 "type": "record",
2134 "name": "LongList",
2135 "aliases": ["LinkedLongs"],
2136 "fields" : [
2137 {"name": "value", "type": "long"},
2138 {"name": "next", "type": ["null", "LongList"]}
2139 ]
2140 }"#,
2141 )
2142 .unwrap();
2143
2144 assert_eq!(
2145 schema,
2146 Schema::Complex(ComplexType::Record(Record {
2147 name: "LongList",
2148 namespace: None,
2149 doc: None,
2150 aliases: vec!["LinkedLongs"],
2151 fields: vec![
2152 Field {
2153 name: "value",
2154 doc: None,
2155 r#type: Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)),
2156 default: None,
2157 aliases: vec![],
2158 },
2159 Field {
2160 name: "next",
2161 doc: None,
2162 r#type: Schema::Union(vec![
2163 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2164 Schema::TypeName(TypeName::Ref("LongList")),
2165 ]),
2166 default: None,
2167 aliases: vec![],
2168 }
2169 ],
2170 attributes: Attributes::default(),
2171 }))
2172 );
2173
2174 let err = AvroField::try_from(&schema).unwrap_err().to_string();
2176 assert_eq!(err, "Parser error: Failed to resolve .LongList");
2177
2178 let schema: Schema = serde_json::from_str(
2179 r#"{
2180 "type":"record",
2181 "name":"topLevelRecord",
2182 "fields":[
2183 {
2184 "name":"id",
2185 "type":[
2186 "int",
2187 "null"
2188 ]
2189 },
2190 {
2191 "name":"timestamp_col",
2192 "type":[
2193 {
2194 "type":"long",
2195 "logicalType":"timestamp-micros"
2196 },
2197 "null"
2198 ]
2199 }
2200 ]
2201 }"#,
2202 )
2203 .unwrap();
2204
2205 assert_eq!(
2206 schema,
2207 Schema::Complex(ComplexType::Record(Record {
2208 name: "topLevelRecord",
2209 namespace: None,
2210 doc: None,
2211 aliases: vec![],
2212 fields: vec![
2213 Field {
2214 name: "id",
2215 doc: None,
2216 r#type: Schema::Union(vec![
2217 Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)),
2218 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2219 ]),
2220 default: None,
2221 aliases: vec![],
2222 },
2223 Field {
2224 name: "timestamp_col",
2225 doc: None,
2226 r#type: Schema::Union(vec![
2227 Schema::Type(timestamp),
2228 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2229 ]),
2230 default: None,
2231 aliases: vec![],
2232 }
2233 ],
2234 attributes: Default::default(),
2235 }))
2236 );
2237 let codec = AvroField::try_from(&schema).unwrap();
2238 let expected_arrow_field = arrow_schema::Field::new(
2239 "topLevelRecord",
2240 DataType::Struct(Fields::from(vec![
2241 arrow_schema::Field::new("id", DataType::Int32, true),
2242 arrow_schema::Field::new(
2243 "timestamp_col",
2244 DataType::Timestamp(TimeUnit::Microsecond, Some("+00:00".into())),
2245 true,
2246 ),
2247 ])),
2248 false,
2249 )
2250 .with_metadata(std::collections::HashMap::from([(
2251 AVRO_NAME_METADATA_KEY.to_string(),
2252 "topLevelRecord".to_string(),
2253 )]));
2254
2255 assert_eq!(codec.field(), expected_arrow_field);
2256
2257 let schema: Schema = serde_json::from_str(
2258 r#"{
2259 "type": "record",
2260 "name": "HandshakeRequest", "namespace":"org.apache.avro.ipc",
2261 "fields": [
2262 {"name": "clientHash", "type": {"type": "fixed", "name": "MD5", "size": 16}},
2263 {"name": "clientProtocol", "type": ["null", "string"]},
2264 {"name": "serverHash", "type": "MD5"},
2265 {"name": "meta", "type": ["null", {"type": "map", "values": "bytes"}]}
2266 ]
2267 }"#,
2268 )
2269 .unwrap();
2270
2271 assert_eq!(
2272 schema,
2273 Schema::Complex(ComplexType::Record(Record {
2274 name: "HandshakeRequest",
2275 namespace: Some("org.apache.avro.ipc"),
2276 doc: None,
2277 aliases: vec![],
2278 fields: vec![
2279 Field {
2280 name: "clientHash",
2281 doc: None,
2282 r#type: Schema::Complex(ComplexType::Fixed(Fixed {
2283 name: "MD5",
2284 namespace: None,
2285 aliases: vec![],
2286 size: 16,
2287 attributes: Default::default(),
2288 })),
2289 default: None,
2290 aliases: vec![],
2291 },
2292 Field {
2293 name: "clientProtocol",
2294 doc: None,
2295 r#type: Schema::Union(vec![
2296 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2297 Schema::TypeName(TypeName::Primitive(PrimitiveType::String)),
2298 ]),
2299 default: None,
2300 aliases: vec![],
2301 },
2302 Field {
2303 name: "serverHash",
2304 doc: None,
2305 r#type: Schema::TypeName(TypeName::Ref("MD5")),
2306 default: None,
2307 aliases: vec![],
2308 },
2309 Field {
2310 name: "meta",
2311 doc: None,
2312 r#type: Schema::Union(vec![
2313 Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)),
2314 Schema::Complex(ComplexType::Map(Map {
2315 values: Box::new(Schema::TypeName(TypeName::Primitive(
2316 PrimitiveType::Bytes
2317 ))),
2318 attributes: Default::default(),
2319 })),
2320 ]),
2321 default: None,
2322 aliases: vec![],
2323 }
2324 ],
2325 attributes: Default::default(),
2326 }))
2327 );
2328 }
2329
2330 #[test]
2331 fn test_canonical_form_generation_comprehensive_record() {
2332 let json_str = r#"{
2334 "type": "record",
2335 "name": "E2eComprehensive",
2336 "namespace": "org.apache.arrow.avrotests.v1",
2337 "doc": "Comprehensive Avro writer schema to exercise arrow-avro Reader/Decoder paths.",
2338 "fields": [
2339 {"name": "id", "type": "long", "doc": "Primary row id", "aliases": ["identifier"]},
2340 {"name": "flag", "type": "boolean", "default": true, "doc": "A sample boolean with default true"},
2341 {"name": "ratio_f32", "type": "float", "default": 0.0, "doc": "Float32 example"},
2342 {"name": "ratio_f64", "type": "double", "default": 0.0, "doc": "Float64 example"},
2343 {"name": "count_i32", "type": "int", "default": 0, "doc": "Int32 example"},
2344 {"name": "count_i64", "type": "long", "default": 0, "doc": "Int64 example"},
2345 {"name": "opt_i32_nullfirst", "type": ["null", "int"], "default": null, "doc": "Nullable int (null-first)"},
2346 {"name": "opt_str_nullsecond", "type": ["string", "null"], "default": "", "aliases": ["old_opt_str"], "doc": "Nullable string (null-second). Default is empty string."},
2347 {"name": "tri_union_prim", "type": ["int", "string", "boolean"], "default": 0, "doc": "Union[int, string, boolean] with default on first branch (int=0)."},
2348 {"name": "str_utf8", "type": "string", "default": "default", "doc": "Plain Utf8 string (Reader may use Utf8View)."},
2349 {"name": "raw_bytes", "type": "bytes", "default": "", "doc": "Raw bytes field"},
2350 {"name": "fx16_plain", "type": {"type": "fixed", "name": "Fx16", "namespace": "org.apache.arrow.avrotests.v1.types", "aliases": ["Fixed16Old"], "size": 16}, "doc": "Plain fixed(16)"},
2351 {"name": "dec_bytes_s10_2", "type": {"type": "bytes", "logicalType": "decimal", "precision": 10, "scale": 2}, "doc": "Decimal encoded on bytes, precision 10, scale 2"},
2352 {"name": "dec_fix_s20_4", "type": {"type": "fixed", "name": "DecFix20", "namespace": "org.apache.arrow.avrotests.v1.types", "size": 20, "logicalType": "decimal", "precision": 20, "scale": 4}, "doc": "Decimal encoded on fixed(20), precision 20, scale 4"},
2353 {"name": "uuid_str", "type": {"type": "string", "logicalType": "uuid"}, "doc": "UUID logical type on string"},
2354 {"name": "d_date", "type": {"type": "int", "logicalType": "date"}, "doc": "Date32: days since 1970-01-01"},
2355 {"name": "t_millis", "type": {"type": "int", "logicalType": "time-millis"}, "doc": "Time32-millis"},
2356 {"name": "t_micros", "type": {"type": "long", "logicalType": "time-micros"}, "doc": "Time64-micros"},
2357 {"name": "ts_millis_utc", "type": {"type": "long", "logicalType": "timestamp-millis"}, "doc": "Timestamp ms (UTC)"},
2358 {"name": "ts_micros_utc", "type": {"type": "long", "logicalType": "timestamp-micros"}, "doc": "Timestamp µs (UTC)"},
2359 {"name": "ts_millis_local", "type": {"type": "long", "logicalType": "local-timestamp-millis"}, "doc": "Local timestamp ms"},
2360 {"name": "ts_micros_local", "type": {"type": "long", "logicalType": "local-timestamp-micros"}, "doc": "Local timestamp µs"},
2361 {"name": "interval_mdn", "type": {"type": "fixed", "name": "Dur12", "namespace": "org.apache.arrow.avrotests.v1.types", "size": 12, "logicalType": "duration"}, "doc": "Duration: fixed(12) little-endian (months, days, millis)"},
2362 {"name": "status", "type": {"type": "enum", "name": "Status", "namespace": "org.apache.arrow.avrotests.v1.types", "symbols": ["UNKNOWN", "NEW", "PROCESSING", "DONE"], "aliases": ["State"], "doc": "Processing status enum with default"}, "default": "UNKNOWN", "doc": "Enum field using default when resolving"},
2363 {"name": "arr_union", "type": {"type": "array", "items": ["long", "string", "null"]}, "default": [], "doc": "Array whose items are a union[long,string,null]"},
2364 {"name": "map_union", "type": {"type": "map", "values": ["null", "double", "string"]}, "default": {}, "doc": "Map whose values are a union[null,double,string]"},
2365 {"name": "address", "type": {"type": "record", "name": "Address", "namespace": "org.apache.arrow.avrotests.v1.types", "doc": "Postal address with defaults and field alias", "fields": [
2366 {"name": "street", "type": "string", "default": "", "aliases": ["street_name"], "doc": "Street (field alias = street_name)"},
2367 {"name": "zip", "type": "int", "default": 0, "doc": "ZIP/postal code"},
2368 {"name": "country", "type": "string", "default": "US", "doc": "Country code"}
2369 ]}, "doc": "Embedded Address record"},
2370 {"name": "maybe_auth", "type": {"type": "record", "name": "MaybeAuth", "namespace": "org.apache.arrow.avrotests.v1.types", "doc": "Optional auth token model", "fields": [
2371 {"name": "user", "type": "string", "doc": "Username"},
2372 {"name": "token", "type": ["null", "bytes"], "default": null, "doc": "Nullable auth token"}
2373 ]}},
2374 {"name": "union_enum_record_array_map", "type": [
2375 {"type": "enum", "name": "Color", "namespace": "org.apache.arrow.avrotests.v1.types", "symbols": ["RED", "GREEN", "BLUE"], "doc": "Color enum"},
2376 {"type": "record", "name": "RecA", "namespace": "org.apache.arrow.avrotests.v1.types", "fields": [{"name": "a", "type": "int"}, {"name": "b", "type": "string"}]},
2377 {"type": "record", "name": "RecB", "namespace": "org.apache.arrow.avrotests.v1.types", "fields": [{"name": "x", "type": "long"}, {"name": "y", "type": "bytes"}]},
2378 {"type": "array", "items": "long"},
2379 {"type": "map", "values": "string"}
2380 ], "doc": "Union of enum, two records, array, and map"},
2381 {"name": "union_date_or_fixed4", "type": [
2382 {"type": "int", "logicalType": "date"},
2383 {"type": "fixed", "name": "Fx4", "size": 4}
2384 ], "doc": "Union of date(int) or fixed(4)"},
2385 {"name": "union_interval_or_string", "type": [
2386 {"type": "fixed", "name": "Dur12U", "size": 12, "logicalType": "duration"},
2387 "string"
2388 ], "doc": "Union of duration(fixed12) or string"},
2389 {"name": "union_uuid_or_fixed10", "type": [
2390 {"type": "string", "logicalType": "uuid"},
2391 {"type": "fixed", "name": "Fx10", "size": 10}
2392 ], "doc": "Union of UUID string or fixed(10)"},
2393 {"name": "array_records_with_union", "type": {"type": "array", "items": {
2394 "type": "record", "name": "KV", "namespace": "org.apache.arrow.avrotests.v1.types",
2395 "fields": [
2396 {"name": "key", "type": "string"},
2397 {"name": "val", "type": ["null", "int", "long"], "default": null}
2398 ]
2399 }}, "doc": "Array<record{key, val: union[null,int,long]}>", "default": []},
2400 {"name": "union_map_or_array_int", "type": [
2401 {"type": "map", "values": "int"},
2402 {"type": "array", "items": "int"}
2403 ], "doc": "Union[map<string,int>, array<int>]"},
2404 {"name": "renamed_with_default", "type": "int", "default": 42, "aliases": ["old_count"], "doc": "Field with alias and default"},
2405 {"name": "person", "type": {"type": "record", "name": "PersonV2", "namespace": "com.example.v2", "aliases": ["com.example.Person"], "doc": "Person record with alias pointing to previous namespace/name", "fields": [
2406 {"name": "name", "type": "string"},
2407 {"name": "age", "type": "int", "default": 0}
2408 ]}, "doc": "Record using type alias for schema evolution tests"}
2409 ]
2410 }"#;
2411 let avro = AvroSchema::new(json_str.to_string());
2412 let parsed = avro.schema().expect("schema should deserialize");
2413 let expected_canonical_form = r#"{"name":"org.apache.arrow.avrotests.v1.E2eComprehensive","type":"record","fields":[{"name":"id","type":"long"},{"name":"flag","type":"boolean"},{"name":"ratio_f32","type":"float"},{"name":"ratio_f64","type":"double"},{"name":"count_i32","type":"int"},{"name":"count_i64","type":"long"},{"name":"opt_i32_nullfirst","type":["null","int"]},{"name":"opt_str_nullsecond","type":["string","null"]},{"name":"tri_union_prim","type":["int","string","boolean"]},{"name":"str_utf8","type":"string"},{"name":"raw_bytes","type":"bytes"},{"name":"fx16_plain","type":{"name":"org.apache.arrow.avrotests.v1.types.Fx16","type":"fixed","size":16}},{"name":"dec_bytes_s10_2","type":"bytes"},{"name":"dec_fix_s20_4","type":{"name":"org.apache.arrow.avrotests.v1.types.DecFix20","type":"fixed","size":20}},{"name":"uuid_str","type":"string"},{"name":"d_date","type":"int"},{"name":"t_millis","type":"int"},{"name":"t_micros","type":"long"},{"name":"ts_millis_utc","type":"long"},{"name":"ts_micros_utc","type":"long"},{"name":"ts_millis_local","type":"long"},{"name":"ts_micros_local","type":"long"},{"name":"interval_mdn","type":{"name":"org.apache.arrow.avrotests.v1.types.Dur12","type":"fixed","size":12}},{"name":"status","type":{"name":"org.apache.arrow.avrotests.v1.types.Status","type":"enum","symbols":["UNKNOWN","NEW","PROCESSING","DONE"]}},{"name":"arr_union","type":{"type":"array","items":["long","string","null"]}},{"name":"map_union","type":{"type":"map","values":["null","double","string"]}},{"name":"address","type":{"name":"org.apache.arrow.avrotests.v1.types.Address","type":"record","fields":[{"name":"street","type":"string"},{"name":"zip","type":"int"},{"name":"country","type":"string"}]}},{"name":"maybe_auth","type":{"name":"org.apache.arrow.avrotests.v1.types.MaybeAuth","type":"record","fields":[{"name":"user","type":"string"},{"name":"token","type":["null","bytes"]}]}},{"name":"union_enum_record_array_map","type":[{"name":"org.apache.arrow.avrotests.v1.types.Color","type":"enum","symbols":["RED","GREEN","BLUE"]},{"name":"org.apache.arrow.avrotests.v1.types.RecA","type":"record","fields":[{"name":"a","type":"int"},{"name":"b","type":"string"}]},{"name":"org.apache.arrow.avrotests.v1.types.RecB","type":"record","fields":[{"name":"x","type":"long"},{"name":"y","type":"bytes"}]},{"type":"array","items":"long"},{"type":"map","values":"string"}]},{"name":"union_date_or_fixed4","type":["int",{"name":"org.apache.arrow.avrotests.v1.Fx4","type":"fixed","size":4}]},{"name":"union_interval_or_string","type":[{"name":"org.apache.arrow.avrotests.v1.Dur12U","type":"fixed","size":12},"string"]},{"name":"union_uuid_or_fixed10","type":["string",{"name":"org.apache.arrow.avrotests.v1.Fx10","type":"fixed","size":10}]},{"name":"array_records_with_union","type":{"type":"array","items":{"name":"org.apache.arrow.avrotests.v1.types.KV","type":"record","fields":[{"name":"key","type":"string"},{"name":"val","type":["null","int","long"]}]}}},{"name":"union_map_or_array_int","type":[{"type":"map","values":"int"},{"type":"array","items":"int"}]},{"name":"renamed_with_default","type":"int"},{"name":"person","type":{"name":"com.example.v2.PersonV2","type":"record","fields":[{"name":"name","type":"string"},{"name":"age","type":"int"}]}}]}"#;
2414 let canonical_form =
2415 AvroSchema::generate_canonical_form(&parsed).expect("canonical form should be built");
2416 assert_eq!(
2417 canonical_form, expected_canonical_form,
2418 "Canonical form must match Avro spec PCF exactly"
2419 );
2420 }
2421
2422 #[test]
2423 fn test_new_schema_store() {
2424 let store = SchemaStore::new();
2425 assert!(store.schemas.is_empty());
2426 }
2427
2428 #[test]
2429 fn test_try_from_schemas_rabin() {
2430 let int_avro_schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2431 let record_avro_schema = AvroSchema::new(serde_json::to_string(&record_schema()).unwrap());
2432 let mut schemas: HashMap<Fingerprint, AvroSchema> = HashMap::new();
2433 schemas.insert(
2434 int_avro_schema
2435 .fingerprint(FingerprintAlgorithm::Rabin)
2436 .unwrap(),
2437 int_avro_schema.clone(),
2438 );
2439 schemas.insert(
2440 record_avro_schema
2441 .fingerprint(FingerprintAlgorithm::Rabin)
2442 .unwrap(),
2443 record_avro_schema.clone(),
2444 );
2445 let store = SchemaStore::try_from(schemas).unwrap();
2446 let int_fp = int_avro_schema
2447 .fingerprint(FingerprintAlgorithm::Rabin)
2448 .unwrap();
2449 assert_eq!(store.lookup(&int_fp).cloned(), Some(int_avro_schema));
2450 let rec_fp = record_avro_schema
2451 .fingerprint(FingerprintAlgorithm::Rabin)
2452 .unwrap();
2453 assert_eq!(store.lookup(&rec_fp).cloned(), Some(record_avro_schema));
2454 }
2455
2456 #[test]
2457 fn test_try_from_with_duplicates() {
2458 let int_avro_schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2459 let record_avro_schema = AvroSchema::new(serde_json::to_string(&record_schema()).unwrap());
2460 let mut schemas: HashMap<Fingerprint, AvroSchema> = HashMap::new();
2461 schemas.insert(
2462 int_avro_schema
2463 .fingerprint(FingerprintAlgorithm::Rabin)
2464 .unwrap(),
2465 int_avro_schema.clone(),
2466 );
2467 schemas.insert(
2468 record_avro_schema
2469 .fingerprint(FingerprintAlgorithm::Rabin)
2470 .unwrap(),
2471 record_avro_schema.clone(),
2472 );
2473 schemas.insert(
2475 int_avro_schema
2476 .fingerprint(FingerprintAlgorithm::Rabin)
2477 .unwrap(),
2478 int_avro_schema.clone(),
2479 );
2480 let store = SchemaStore::try_from(schemas).unwrap();
2481 assert_eq!(store.schemas.len(), 2);
2482 let int_fp = int_avro_schema
2483 .fingerprint(FingerprintAlgorithm::Rabin)
2484 .unwrap();
2485 assert_eq!(store.lookup(&int_fp).cloned(), Some(int_avro_schema));
2486 }
2487
2488 #[test]
2489 fn test_register_and_lookup_rabin() {
2490 let mut store = SchemaStore::new();
2491 let schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2492 let fp_enum = store.register(schema.clone()).unwrap();
2493 match fp_enum {
2494 Fingerprint::Rabin(fp_val) => {
2495 assert_eq!(
2496 store.lookup(&Fingerprint::Rabin(fp_val)).cloned(),
2497 Some(schema.clone())
2498 );
2499 assert!(
2500 store
2501 .lookup(&Fingerprint::Rabin(fp_val.wrapping_add(1)))
2502 .is_none()
2503 );
2504 }
2505 Fingerprint::Id(_id) => {
2506 unreachable!("This test should only generate Rabin fingerprints")
2507 }
2508 Fingerprint::Id64(_id) => {
2509 unreachable!("This test should only generate Rabin fingerprints")
2510 }
2511 #[cfg(feature = "md5")]
2512 Fingerprint::MD5(_id) => {
2513 unreachable!("This test should only generate Rabin fingerprints")
2514 }
2515 #[cfg(feature = "sha256")]
2516 Fingerprint::SHA256(_id) => {
2517 unreachable!("This test should only generate Rabin fingerprints")
2518 }
2519 }
2520 }
2521
2522 #[test]
2523 fn test_set_and_lookup_id() {
2524 let mut store = SchemaStore::new();
2525 let schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2526 let id = 42u32;
2527 let fp = Fingerprint::Id(id);
2528 let out_fp = store.set(fp, schema.clone()).unwrap();
2529 assert_eq!(out_fp, fp);
2530 assert_eq!(store.lookup(&fp).cloned(), Some(schema.clone()));
2531 assert!(store.lookup(&Fingerprint::Id(id.wrapping_add(1))).is_none());
2532 }
2533
2534 #[test]
2535 fn test_set_and_lookup_id64() {
2536 let mut store = SchemaStore::new();
2537 let schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2538 let id64: u64 = 0xDEAD_BEEF_DEAD_BEEF;
2539 let fp = Fingerprint::Id64(id64);
2540 let out_fp = store.set(fp, schema.clone()).unwrap();
2541 assert_eq!(out_fp, fp, "set should return the same Id64 fingerprint");
2542 assert_eq!(
2543 store.lookup(&fp).cloned(),
2544 Some(schema.clone()),
2545 "lookup should find the schema by Id64"
2546 );
2547 assert!(
2548 store
2549 .lookup(&Fingerprint::Id64(id64.wrapping_add(1)))
2550 .is_none(),
2551 "lookup with a different Id64 must return None"
2552 );
2553 }
2554
2555 #[test]
2556 fn test_fingerprint_id64_conversions() {
2557 let algo_from_fp = FingerprintAlgorithm::from(&Fingerprint::Id64(123));
2558 assert_eq!(algo_from_fp, FingerprintAlgorithm::Id64);
2559 let fp_from_algo = Fingerprint::from(FingerprintAlgorithm::Id64);
2560 assert!(matches!(fp_from_algo, Fingerprint::Id64(0)));
2561 let strategy_from_fp = FingerprintStrategy::from(Fingerprint::Id64(5));
2562 assert!(matches!(strategy_from_fp, FingerprintStrategy::Id64(0)));
2563 let algo_from_strategy = FingerprintAlgorithm::from(strategy_from_fp);
2564 assert_eq!(algo_from_strategy, FingerprintAlgorithm::Id64);
2565 }
2566
2567 #[test]
2568 fn test_register_duplicate_schema() {
2569 let mut store = SchemaStore::new();
2570 let schema1 = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2571 let schema2 = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2572 let fingerprint1 = store.register(schema1).unwrap();
2573 let fingerprint2 = store.register(schema2).unwrap();
2574 assert_eq!(fingerprint1, fingerprint2);
2575 assert_eq!(store.schemas.len(), 1);
2576 }
2577
2578 #[test]
2579 fn test_set_and_lookup_with_provided_fingerprint() {
2580 let mut store = SchemaStore::new();
2581 let schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2582 let fp = schema.fingerprint(FingerprintAlgorithm::Rabin).unwrap();
2583 let out_fp = store.set(fp, schema.clone()).unwrap();
2584 assert_eq!(out_fp, fp);
2585 assert_eq!(store.lookup(&fp).cloned(), Some(schema));
2586 }
2587
2588 #[test]
2589 fn test_set_duplicate_same_schema_ok() {
2590 let mut store = SchemaStore::new();
2591 let schema = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2592 let fp = schema.fingerprint(FingerprintAlgorithm::Rabin).unwrap();
2593 let _ = store.set(fp, schema.clone()).unwrap();
2594 let _ = store.set(fp, schema.clone()).unwrap();
2595 assert_eq!(store.schemas.len(), 1);
2596 }
2597
2598 #[test]
2599 fn test_set_duplicate_different_schema_collision_error() {
2600 let mut store = SchemaStore::new();
2601 let schema1 = AvroSchema::new(serde_json::to_string(&int_schema()).unwrap());
2602 let schema2 = AvroSchema::new(serde_json::to_string(&record_schema()).unwrap());
2603 let fp = Fingerprint::Id(123);
2605 let _ = store.set(fp, schema1).unwrap();
2606 let err = store.set(fp, schema2).unwrap_err();
2607 let msg = format!("{err}");
2608 assert!(msg.contains("Schema fingerprint collision"));
2609 }
2610
2611 #[test]
2612 fn test_canonical_form_generation_primitive() {
2613 let schema = int_schema();
2614 let canonical_form = AvroSchema::generate_canonical_form(&schema).unwrap();
2615 assert_eq!(canonical_form, r#""int""#);
2616 }
2617
2618 #[test]
2619 fn test_canonical_form_generation_record() {
2620 let schema = record_schema();
2621 let expected_canonical_form = r#"{"name":"test.namespace.record1","type":"record","fields":[{"name":"field1","type":"int"},{"name":"field2","type":"string"}]}"#;
2622 let canonical_form = AvroSchema::generate_canonical_form(&schema).unwrap();
2623 assert_eq!(canonical_form, expected_canonical_form);
2624 }
2625
2626 #[test]
2627 fn test_fingerprint_calculation() {
2628 let canonical_form = r#"{"fields":[{"name":"a","type":"long"},{"name":"b","type":"string"}],"name":"test","type":"record"}"#;
2629 let expected_fingerprint = 10505236152925314060;
2630 let fingerprint = compute_fingerprint_rabin(canonical_form);
2631 assert_eq!(fingerprint, expected_fingerprint);
2632 }
2633
2634 #[test]
2635 fn test_register_and_lookup_complex_schema() {
2636 let mut store = SchemaStore::new();
2637 let schema = AvroSchema::new(serde_json::to_string(&record_schema()).unwrap());
2638 let canonical_form = r#"{"name":"test.namespace.record1","type":"record","fields":[{"name":"field1","type":"int"},{"name":"field2","type":"string"}]}"#;
2639 let expected_fingerprint = Fingerprint::Rabin(compute_fingerprint_rabin(canonical_form));
2640 let fingerprint = store.register(schema.clone()).unwrap();
2641 assert_eq!(fingerprint, expected_fingerprint);
2642 let looked_up = store.lookup(&fingerprint).cloned();
2643 assert_eq!(looked_up, Some(schema));
2644 }
2645
2646 #[test]
2647 fn test_fingerprints_returns_all_keys() {
2648 let mut store = SchemaStore::new();
2649 let fp_int = store
2650 .register(AvroSchema::new(
2651 serde_json::to_string(&int_schema()).unwrap(),
2652 ))
2653 .unwrap();
2654 let fp_record = store
2655 .register(AvroSchema::new(
2656 serde_json::to_string(&record_schema()).unwrap(),
2657 ))
2658 .unwrap();
2659 let fps = store.fingerprints();
2660 assert_eq!(fps.len(), 2);
2661 assert!(fps.contains(&fp_int));
2662 assert!(fps.contains(&fp_record));
2663 }
2664
2665 #[test]
2666 fn test_canonical_form_strips_attributes() {
2667 let schema_with_attrs = Schema::Complex(ComplexType::Record(Record {
2668 name: "record_with_attrs",
2669 namespace: None,
2670 doc: Some(Cow::from("This doc should be stripped")),
2671 aliases: vec!["alias1", "alias2"],
2672 fields: vec![Field {
2673 name: "f1",
2674 doc: Some(Cow::from("field doc")),
2675 r#type: Schema::Type(Type {
2676 r#type: TypeName::Primitive(PrimitiveType::Bytes),
2677 attributes: Attributes {
2678 logical_type: None,
2679 additional: HashMap::from([("precision", json!(4))]),
2680 },
2681 }),
2682 default: None,
2683 aliases: vec![],
2684 }],
2685 attributes: Attributes {
2686 logical_type: None,
2687 additional: HashMap::from([("custom_attr", json!("value"))]),
2688 },
2689 }));
2690 let expected_canonical_form = r#"{"name":"record_with_attrs","type":"record","fields":[{"name":"f1","type":"bytes"}]}"#;
2691 let canonical_form = AvroSchema::generate_canonical_form(&schema_with_attrs).unwrap();
2692 assert_eq!(canonical_form, expected_canonical_form);
2693 }
2694
2695 #[cfg(not(feature = "avro_custom_types"))]
2696 #[test]
2697 fn test_primitive_mappings() {
2698 let cases = vec![
2699 (DataType::Boolean, "\"boolean\""),
2700 (DataType::Int8, "\"int\""),
2701 (DataType::Int16, "\"int\""),
2702 (DataType::Int32, "\"int\""),
2703 (DataType::Int64, "\"long\""),
2704 (DataType::UInt8, "\"int\""),
2705 (DataType::UInt16, "\"int\""),
2706 (DataType::UInt32, "\"long\""),
2707 (DataType::UInt64, "\"long\""),
2708 (DataType::Float16, "\"float\""),
2709 (DataType::Float32, "\"float\""),
2710 (DataType::Float64, "\"double\""),
2711 (DataType::Utf8, "\"string\""),
2712 (DataType::Binary, "\"bytes\""),
2713 ];
2714 for (dt, avro_token) in cases {
2715 let field = ArrowField::new("col", dt.clone(), false);
2716 let arrow_schema = single_field_schema(field);
2717 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
2718 assert_json_contains(&avro.json_string, avro_token);
2719 }
2720 }
2721
2722 #[cfg(feature = "avro_custom_types")]
2723 #[test]
2724 fn test_primitive_mappings() {
2725 let cases = vec![
2726 (DataType::Boolean, "\"boolean\""),
2727 (DataType::Int8, "\"logicalType\":\"arrow.int8\""),
2728 (DataType::Int16, "\"logicalType\":\"arrow.int16\""),
2729 (DataType::Int32, "\"int\""),
2730 (DataType::Int64, "\"long\""),
2731 (DataType::UInt8, "\"logicalType\":\"arrow.uint8\""),
2732 (DataType::UInt16, "\"logicalType\":\"arrow.uint16\""),
2733 (DataType::UInt32, "\"logicalType\":\"arrow.uint32\""),
2734 (DataType::UInt64, "\"logicalType\":\"arrow.uint64\""),
2735 (DataType::Float16, "\"logicalType\":\"arrow.float16\""),
2736 (DataType::Float32, "\"float\""),
2737 (DataType::Float64, "\"double\""),
2738 (DataType::Utf8, "\"string\""),
2739 (DataType::Binary, "\"bytes\""),
2740 ];
2741 for (dt, avro_token) in cases {
2742 let field = ArrowField::new("col", dt.clone(), false);
2743 let arrow_schema = single_field_schema(field);
2744 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
2745 assert_json_contains(&avro.json_string, avro_token);
2746 }
2747 }
2748
2749 #[cfg(feature = "avro_custom_types")]
2750 #[test]
2751 fn test_custom_fixed_logical_types_preserve_namespace_metadata() {
2752 let namespace = "com.example.types";
2753
2754 let mut md_u64 = HashMap::new();
2755 md_u64.insert(AVRO_NAME_METADATA_KEY.to_string(), "U64Type".to_string());
2756 md_u64.insert(
2757 AVRO_NAMESPACE_METADATA_KEY.to_string(),
2758 namespace.to_string(),
2759 );
2760
2761 let mut md_f16 = HashMap::new();
2762 md_f16.insert(AVRO_NAME_METADATA_KEY.to_string(), "F16Type".to_string());
2763 md_f16.insert(
2764 AVRO_NAMESPACE_METADATA_KEY.to_string(),
2765 namespace.to_string(),
2766 );
2767
2768 let mut md_iv_ym = HashMap::new();
2769 md_iv_ym.insert(AVRO_NAME_METADATA_KEY.to_string(), "IvYmType".to_string());
2770 md_iv_ym.insert(
2771 AVRO_NAMESPACE_METADATA_KEY.to_string(),
2772 namespace.to_string(),
2773 );
2774
2775 let mut md_iv_dt = HashMap::new();
2776 md_iv_dt.insert(AVRO_NAME_METADATA_KEY.to_string(), "IvDtType".to_string());
2777 md_iv_dt.insert(
2778 AVRO_NAMESPACE_METADATA_KEY.to_string(),
2779 namespace.to_string(),
2780 );
2781
2782 let arrow_schema = ArrowSchema::new(vec![
2783 ArrowField::new("u64_col", DataType::UInt64, false).with_metadata(md_u64),
2784 ArrowField::new("f16_col", DataType::Float16, false).with_metadata(md_f16),
2785 ArrowField::new(
2786 "iv_ym_col",
2787 DataType::Interval(IntervalUnit::YearMonth),
2788 false,
2789 )
2790 .with_metadata(md_iv_ym),
2791 ArrowField::new(
2792 "iv_dt_col",
2793 DataType::Interval(IntervalUnit::DayTime),
2794 false,
2795 )
2796 .with_metadata(md_iv_dt),
2797 ]);
2798
2799 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
2800 let root: Value = serde_json::from_str(&avro.json_string).unwrap();
2801 let fields = root
2802 .get("fields")
2803 .and_then(|f| f.as_array())
2804 .expect("record fields array");
2805
2806 let expected = [
2807 ("u64_col", "arrow.uint64"),
2808 ("f16_col", "arrow.float16"),
2809 ("iv_ym_col", "arrow.interval-year-month"),
2810 ("iv_dt_col", "arrow.interval-day-time"),
2811 ];
2812
2813 for (field_name, logical_type) in expected {
2814 let field = fields
2815 .iter()
2816 .find(|f| f.get("name").and_then(Value::as_str) == Some(field_name))
2817 .unwrap_or_else(|| panic!("missing field {field_name}"));
2818 let ty = field
2819 .get("type")
2820 .and_then(Value::as_object)
2821 .unwrap_or_else(|| panic!("field {field_name} type must be object"));
2822
2823 assert_eq!(ty.get("type").and_then(Value::as_str), Some("fixed"));
2824 assert_eq!(
2825 ty.get("logicalType").and_then(Value::as_str),
2826 Some(logical_type)
2827 );
2828 assert_eq!(
2829 ty.get("namespace").and_then(Value::as_str),
2830 Some(namespace),
2831 "field {field_name} must preserve avro.namespace metadata"
2832 );
2833 }
2834 }
2835
2836 #[cfg(feature = "avro_custom_types")]
2837 #[test]
2838 fn test_custom_fixed_logical_types_omit_namespace_without_metadata() {
2839 let mut md_u64 = HashMap::new();
2840 md_u64.insert(AVRO_NAME_METADATA_KEY.to_string(), "U64Type".to_string());
2841
2842 let mut md_f16 = HashMap::new();
2843 md_f16.insert(AVRO_NAME_METADATA_KEY.to_string(), "F16Type".to_string());
2844
2845 let mut md_iv_ym = HashMap::new();
2846 md_iv_ym.insert(AVRO_NAME_METADATA_KEY.to_string(), "IvYmType".to_string());
2847
2848 let mut md_iv_dt = HashMap::new();
2849 md_iv_dt.insert(AVRO_NAME_METADATA_KEY.to_string(), "IvDtType".to_string());
2850
2851 let arrow_schema = ArrowSchema::new(vec![
2852 ArrowField::new("u64_col", DataType::UInt64, false).with_metadata(md_u64),
2853 ArrowField::new("f16_col", DataType::Float16, false).with_metadata(md_f16),
2854 ArrowField::new(
2855 "iv_ym_col",
2856 DataType::Interval(IntervalUnit::YearMonth),
2857 false,
2858 )
2859 .with_metadata(md_iv_ym),
2860 ArrowField::new(
2861 "iv_dt_col",
2862 DataType::Interval(IntervalUnit::DayTime),
2863 false,
2864 )
2865 .with_metadata(md_iv_dt),
2866 ]);
2867
2868 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
2869 let root: Value = serde_json::from_str(&avro.json_string).unwrap();
2870 let fields = root
2871 .get("fields")
2872 .and_then(|f| f.as_array())
2873 .expect("record fields array");
2874
2875 for field_name in ["u64_col", "f16_col", "iv_ym_col", "iv_dt_col"] {
2876 let field = fields
2877 .iter()
2878 .find(|f| f.get("name").and_then(Value::as_str) == Some(field_name))
2879 .unwrap_or_else(|| panic!("missing field {field_name}"));
2880 let ty = field
2881 .get("type")
2882 .and_then(Value::as_object)
2883 .unwrap_or_else(|| panic!("field {field_name} type must be object"));
2884
2885 assert_eq!(ty.get("type").and_then(Value::as_str), Some("fixed"));
2886 assert!(
2887 !ty.contains_key("namespace"),
2888 "field {field_name} should not include namespace when metadata lacks avro.namespace"
2889 );
2890 }
2891 }
2892
2893 #[test]
2894 fn test_temporal_mappings() {
2895 let cases = vec![
2896 (DataType::Date32, "\"logicalType\":\"date\""),
2897 (
2898 DataType::Time32(TimeUnit::Millisecond),
2899 "\"logicalType\":\"time-millis\"",
2900 ),
2901 (
2902 DataType::Time64(TimeUnit::Microsecond),
2903 "\"logicalType\":\"time-micros\"",
2904 ),
2905 (
2906 DataType::Timestamp(TimeUnit::Millisecond, None),
2907 "\"logicalType\":\"local-timestamp-millis\"",
2908 ),
2909 (
2910 DataType::Timestamp(TimeUnit::Microsecond, Some("+00:00".into())),
2911 "\"logicalType\":\"timestamp-micros\"",
2912 ),
2913 ];
2914 for (dt, needle) in cases {
2915 let field = ArrowField::new("ts", dt.clone(), true);
2916 let arrow_schema = single_field_schema(field);
2917 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
2918 assert_json_contains(&avro.json_string, needle);
2919 }
2920 }
2921
2922 #[test]
2923 fn test_decimal_and_uuid() {
2924 let decimal_field = ArrowField::new("amount", DataType::Decimal128(25, 2), false);
2925 let dec_schema = single_field_schema(decimal_field);
2926 let avro_dec = AvroSchema::try_from(&dec_schema).unwrap();
2927 assert_json_contains(&avro_dec.json_string, "\"logicalType\":\"decimal\"");
2928 assert_json_contains(&avro_dec.json_string, "\"precision\":25");
2929 assert_json_contains(&avro_dec.json_string, "\"scale\":2");
2930 let mut md = HashMap::new();
2931 md.insert("logicalType".into(), "uuid".into());
2932 let uuid_field =
2933 ArrowField::new("id", DataType::FixedSizeBinary(16), false).with_metadata(md);
2934 let uuid_schema = single_field_schema(uuid_field);
2935 let avro_uuid = AvroSchema::try_from(&uuid_schema).unwrap();
2936 assert_json_contains(&avro_uuid.json_string, "\"logicalType\":\"uuid\"");
2937 }
2938
2939 #[cfg(not(feature = "avro_custom_types"))]
2940 #[test]
2941 fn test_interval_month_day_nano_duration_schema() {
2942 let interval_field = ArrowField::new(
2943 "span",
2944 DataType::Interval(IntervalUnit::MonthDayNano),
2945 false,
2946 );
2947 let s = single_field_schema(interval_field);
2948 let avro = AvroSchema::try_from(&s).unwrap();
2949 assert_json_contains(&avro.json_string, "\"logicalType\":\"duration\"");
2950 assert_json_contains(&avro.json_string, "\"size\":12");
2951 }
2952
2953 #[cfg(feature = "avro_custom_types")]
2954 #[test]
2955 fn test_interval_month_day_nano_custom_schema() {
2956 let interval_field = ArrowField::new(
2957 "span",
2958 DataType::Interval(IntervalUnit::MonthDayNano),
2959 false,
2960 );
2961 let s = single_field_schema(interval_field);
2962 let avro = AvroSchema::try_from(&s).unwrap();
2963 assert_json_contains(
2964 &avro.json_string,
2965 "\"logicalType\":\"arrow.interval-month-day-nano\"",
2966 );
2967 assert_json_contains(&avro.json_string, "\"size\":16");
2968 }
2969
2970 #[cfg(feature = "avro_custom_types")]
2971 #[test]
2972 fn test_duration_custom_logical_type() {
2973 let dur_field = ArrowField::new("latency", DataType::Duration(TimeUnit::Nanosecond), false);
2974 let s2 = single_field_schema(dur_field);
2975 let avro2 = AvroSchema::try_from(&s2).unwrap();
2976 assert_json_contains(
2977 &avro2.json_string,
2978 "\"logicalType\":\"arrow.duration-nanos\"",
2979 );
2980 }
2981
2982 #[test]
2983 fn test_complex_types() {
2984 let list_dt = DataType::List(Arc::new(ArrowField::new("item", DataType::Int32, true)));
2985 let list_schema = single_field_schema(ArrowField::new("numbers", list_dt, false));
2986 let avro_list = AvroSchema::try_from(&list_schema).unwrap();
2987 assert_json_contains(&avro_list.json_string, "\"type\":\"array\"");
2988 assert_json_contains(&avro_list.json_string, "\"items\"");
2989 let value_field = ArrowField::new(
2990 arrow_schema::Field::MAP_VALUE_FIELD_DEFAULT_NAME,
2991 DataType::Boolean,
2992 true,
2993 );
2994 let entries_struct = ArrowField::new(
2995 arrow_schema::Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
2996 DataType::Struct(Fields::from(vec![
2997 ArrowField::new(
2998 arrow_schema::Field::MAP_KEY_FIELD_DEFAULT_NAME,
2999 DataType::Utf8,
3000 false,
3001 ),
3002 value_field.clone(),
3003 ])),
3004 false,
3005 );
3006 let map_dt = DataType::Map(Arc::new(entries_struct), false);
3007 let map_schema = single_field_schema(ArrowField::new("props", map_dt, false));
3008 let avro_map = AvroSchema::try_from(&map_schema).unwrap();
3009 assert_json_contains(&avro_map.json_string, "\"type\":\"map\"");
3010 assert_json_contains(&avro_map.json_string, "\"values\"");
3011 let struct_dt = DataType::Struct(Fields::from(vec![
3012 ArrowField::new("f1", DataType::Int64, false),
3013 ArrowField::new("f2", DataType::Utf8, true),
3014 ]));
3015 let struct_schema = single_field_schema(ArrowField::new("person", struct_dt, true));
3016 let avro_struct = AvroSchema::try_from(&struct_schema).unwrap();
3017 assert_json_contains(&avro_struct.json_string, "\"type\":\"record\"");
3018 assert_json_contains(&avro_struct.json_string, "\"null\"");
3019 }
3020
3021 #[test]
3022 fn test_enum_dictionary() {
3023 let mut md = HashMap::new();
3024 md.insert(
3025 AVRO_ENUM_SYMBOLS_METADATA_KEY.into(),
3026 "[\"OPEN\",\"CLOSED\"]".into(),
3027 );
3028 let enum_dt = DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8));
3029 let field = ArrowField::new("status", enum_dt, false).with_metadata(md);
3030 let schema = single_field_schema(field);
3031 let avro = AvroSchema::try_from(&schema).unwrap();
3032 assert_json_contains(&avro.json_string, "\"type\":\"enum\"");
3033 assert_json_contains(&avro.json_string, "\"symbols\":[\"OPEN\",\"CLOSED\"]");
3034 }
3035
3036 #[test]
3037 fn test_run_end_encoded() {
3038 let ree_dt = DataType::RunEndEncoded(
3039 Arc::new(ArrowField::new("run_ends", DataType::Int32, false)),
3040 Arc::new(ArrowField::new("values", DataType::Utf8, false)),
3041 );
3042 let s = single_field_schema(ArrowField::new("text", ree_dt, false));
3043 let avro = AvroSchema::try_from(&s).unwrap();
3044 assert_json_contains(&avro.json_string, "\"string\"");
3045 }
3046
3047 #[test]
3048 fn test_dense_union() {
3049 let uf: UnionFields = vec![
3050 (2i8, Arc::new(ArrowField::new("a", DataType::Int32, false))),
3051 (7i8, Arc::new(ArrowField::new("b", DataType::Utf8, true))),
3052 ]
3053 .into_iter()
3054 .collect();
3055 let union_dt = DataType::Union(uf, UnionMode::Dense);
3056 let s = single_field_schema(ArrowField::new("u", union_dt, false));
3057 let avro =
3058 AvroSchema::try_from(&s).expect("Arrow Union -> Avro union conversion should succeed");
3059 let v: serde_json::Value = serde_json::from_str(&avro.json_string).unwrap();
3060 let fields = v
3061 .get("fields")
3062 .and_then(|x| x.as_array())
3063 .expect("fields array");
3064 let u_field = fields
3065 .iter()
3066 .find(|f| f.get("name").and_then(|n| n.as_str()) == Some("u"))
3067 .expect("field 'u'");
3068 let union = u_field.get("type").expect("u.type");
3069 let arr = union.as_array().expect("u.type must be Avro union array");
3070 assert_eq!(arr.len(), 2, "expected two union branches");
3071 let first = &arr[0];
3072 let obj = first
3073 .as_object()
3074 .expect("first branch should be an object with metadata");
3075 assert_eq!(obj.get("type").and_then(|t| t.as_str()), Some("int"));
3076 assert_eq!(
3077 obj.get("arrowUnionMode").and_then(|m| m.as_str()),
3078 Some("dense")
3079 );
3080 let type_ids: Vec<i64> = obj
3081 .get("arrowUnionTypeIds")
3082 .and_then(|a| a.as_array())
3083 .expect("arrowUnionTypeIds array")
3084 .iter()
3085 .map(|n| n.as_i64().expect("i64"))
3086 .collect();
3087 assert_eq!(type_ids, vec![2, 7], "type id ordering should be preserved");
3088 assert_eq!(arr[1], Value::String("string".into()));
3089 }
3090
3091 #[test]
3092 fn round_trip_primitive() {
3093 let arrow_schema = ArrowSchema::new(vec![ArrowField::new("f1", DataType::Int32, false)]);
3094 let avro_schema = AvroSchema::try_from(&arrow_schema).unwrap();
3095 let decoded = avro_schema.schema().unwrap();
3096 assert!(matches!(decoded, Schema::Complex(_)));
3097 }
3098
3099 #[test]
3100 fn test_name_generator_sanitization_and_uniqueness() {
3101 let f1 = ArrowField::new("weird-name", DataType::FixedSizeBinary(8), false);
3102 let f2 = ArrowField::new("weird name", DataType::FixedSizeBinary(8), false);
3103 let f3 = ArrowField::new("123bad", DataType::FixedSizeBinary(8), false);
3104 let arrow_schema = ArrowSchema::new(vec![f1, f2, f3]);
3105 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
3106 assert_json_contains(&avro.json_string, "\"name\":\"weird_name\"");
3107 assert_json_contains(&avro.json_string, "\"name\":\"weird_name_1\"");
3108 assert_json_contains(&avro.json_string, "\"name\":\"_123bad\"");
3109 }
3110
3111 #[cfg(not(feature = "avro_custom_types"))]
3112 #[test]
3113 fn test_date64_logical_type_mapping() {
3114 let field = ArrowField::new("d", DataType::Date64, true);
3115 let schema = single_field_schema(field);
3116 let avro = AvroSchema::try_from(&schema).unwrap();
3117 assert_json_contains(
3118 &avro.json_string,
3119 "\"logicalType\":\"local-timestamp-millis\"",
3120 );
3121 }
3122
3123 #[cfg(feature = "avro_custom_types")]
3124 #[test]
3125 fn test_date64_logical_type_mapping_custom() {
3126 let field = ArrowField::new("d", DataType::Date64, true);
3127 let schema = single_field_schema(field);
3128 let avro = AvroSchema::try_from(&schema).unwrap();
3129 assert_json_contains(&avro.json_string, "\"logicalType\":\"arrow.date64\"");
3130 }
3131
3132 #[cfg(feature = "avro_custom_types")]
3133 #[test]
3134 fn test_duration_list_extras_propagated() {
3135 let child = ArrowField::new("lat", DataType::Duration(TimeUnit::Microsecond), false);
3136 let list_dt = DataType::List(Arc::new(child));
3137 let arrow_schema = single_field_schema(ArrowField::new("durations", list_dt, false));
3138 let avro = AvroSchema::try_from(&arrow_schema).unwrap();
3139 assert_json_contains(
3140 &avro.json_string,
3141 "\"logicalType\":\"arrow.duration-micros\"",
3142 );
3143 }
3144
3145 #[cfg(not(feature = "avro_custom_types"))]
3146 #[test]
3147 fn test_interval_yearmonth_extra() {
3148 let field = ArrowField::new("iv", DataType::Interval(IntervalUnit::YearMonth), false);
3149 let schema = single_field_schema(field);
3150 let avro = AvroSchema::try_from(&schema).unwrap();
3151 assert_json_contains(&avro.json_string, "\"arrowIntervalUnit\":\"yearmonth\"");
3152 }
3153
3154 #[cfg(not(feature = "avro_custom_types"))]
3155 #[test]
3156 fn test_interval_daytime_extra() {
3157 let field = ArrowField::new("iv_dt", DataType::Interval(IntervalUnit::DayTime), false);
3158 let schema = single_field_schema(field);
3159 let avro = AvroSchema::try_from(&schema).unwrap();
3160 assert_json_contains(&avro.json_string, "\"arrowIntervalUnit\":\"daytime\"");
3161 }
3162
3163 #[cfg(feature = "avro_custom_types")]
3164 #[test]
3165 fn test_interval_yearmonth_custom() {
3166 let field = ArrowField::new("iv", DataType::Interval(IntervalUnit::YearMonth), false);
3167 let schema = single_field_schema(field);
3168 let avro = AvroSchema::try_from(&schema).unwrap();
3169 assert_json_contains(
3170 &avro.json_string,
3171 "\"logicalType\":\"arrow.interval-year-month\"",
3172 );
3173 }
3174
3175 #[cfg(feature = "avro_custom_types")]
3176 #[test]
3177 fn test_interval_daytime_custom() {
3178 let field = ArrowField::new("iv_dt", DataType::Interval(IntervalUnit::DayTime), false);
3179 let schema = single_field_schema(field);
3180 let avro = AvroSchema::try_from(&schema).unwrap();
3181 assert_json_contains(
3182 &avro.json_string,
3183 "\"logicalType\":\"arrow.interval-day-time\"",
3184 );
3185 }
3186
3187 #[test]
3188 fn test_fixed_size_list_extra() {
3189 let child = ArrowField::new("item", DataType::Int32, false);
3190 let dt = DataType::FixedSizeList(Arc::new(child), 3);
3191 let schema = single_field_schema(ArrowField::new("triples", dt, false));
3192 let avro = AvroSchema::try_from(&schema).unwrap();
3193 assert_json_contains(&avro.json_string, "\"arrowFixedSize\":3");
3194 }
3195
3196 #[cfg(feature = "avro_custom_types")]
3197 #[test]
3198 fn test_map_duration_value_extra() {
3199 let val_field = ArrowField::new(
3200 ArrowField::MAP_VALUE_FIELD_DEFAULT_NAME,
3201 DataType::Duration(TimeUnit::Second),
3202 true,
3203 );
3204 let entries_struct = ArrowField::new(
3205 ArrowField::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3206 DataType::Struct(Fields::from(vec![
3207 ArrowField::new(
3208 ArrowField::MAP_KEY_FIELD_DEFAULT_NAME,
3209 DataType::Utf8,
3210 false,
3211 ),
3212 val_field,
3213 ])),
3214 false,
3215 );
3216 let map_dt = DataType::Map(Arc::new(entries_struct), false);
3217 let schema = single_field_schema(ArrowField::new("metrics", map_dt, false));
3218 let avro = AvroSchema::try_from(&schema).unwrap();
3219 assert_json_contains(
3220 &avro.json_string,
3221 "\"logicalType\":\"arrow.duration-seconds\"",
3222 );
3223 }
3224
3225 #[test]
3226 fn test_schema_with_non_string_defaults_decodes_successfully() {
3227 let schema_json = r#"{
3228 "type": "record",
3229 "name": "R",
3230 "fields": [
3231 {"name": "a", "type": "int", "default": 0},
3232 {"name": "b", "type": {"type": "array", "items": "long"}, "default": [1, 2, 3]},
3233 {"name": "c", "type": {"type": "map", "values": "double"}, "default": {"x": 1.5, "y": 2.5}},
3234 {"name": "inner", "type": {"type": "record", "name": "Inner", "fields": [
3235 {"name": "flag", "type": "boolean", "default": true},
3236 {"name": "name", "type": "string", "default": "hi"}
3237 ]}, "default": {"flag": false, "name": "d"}},
3238 {"name": "u", "type": ["int", "null"], "default": 42}
3239 ]
3240 }"#;
3241 let schema: Schema = serde_json::from_str(schema_json).expect("schema should parse");
3242 match &schema {
3243 Schema::Complex(ComplexType::Record(_)) => {}
3244 other => panic!("expected record schema, got: {other:?}"),
3245 }
3246 let field = crate::codec::AvroField::try_from(&schema)
3248 .expect("Avro->Arrow conversion should succeed");
3249 let arrow_field = field.field();
3250 let expected_list_item = ArrowField::new(
3252 arrow_schema::Field::LIST_FIELD_DEFAULT_NAME,
3253 DataType::Int64,
3254 false,
3255 );
3256 let expected_b = ArrowField::new("b", DataType::List(Arc::new(expected_list_item)), false);
3257
3258 let expected_map_value = ArrowField::new(
3259 arrow_schema::Field::MAP_VALUE_FIELD_DEFAULT_NAME,
3260 DataType::Float64,
3261 false,
3262 );
3263 let expected_entries = ArrowField::new(
3264 arrow_schema::Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3265 DataType::Struct(Fields::from(vec![
3266 ArrowField::new(
3267 arrow_schema::Field::MAP_KEY_FIELD_DEFAULT_NAME,
3268 DataType::Utf8,
3269 false,
3270 ),
3271 expected_map_value,
3272 ])),
3273 false,
3274 );
3275 let expected_c =
3276 ArrowField::new("c", DataType::Map(Arc::new(expected_entries), false), false);
3277 let mut inner_md = std::collections::HashMap::new();
3278 inner_md.insert(AVRO_NAME_METADATA_KEY.to_string(), "Inner".to_string());
3279 let expected_inner = ArrowField::new(
3280 "inner",
3281 DataType::Struct(Fields::from(vec![
3282 ArrowField::new("flag", DataType::Boolean, false),
3283 ArrowField::new("name", DataType::Utf8, false),
3284 ])),
3285 false,
3286 )
3287 .with_metadata(inner_md);
3288 let mut root_md = std::collections::HashMap::new();
3289 root_md.insert(AVRO_NAME_METADATA_KEY.to_string(), "R".to_string());
3290 let expected = ArrowField::new(
3291 "R",
3292 DataType::Struct(Fields::from(vec![
3293 ArrowField::new("a", DataType::Int32, false),
3294 expected_b,
3295 expected_c,
3296 expected_inner,
3297 ArrowField::new("u", DataType::Int32, true),
3298 ])),
3299 false,
3300 )
3301 .with_metadata(root_md);
3302 assert_eq!(arrow_field, expected);
3303 }
3304
3305 #[test]
3306 fn default_order_is_consistent() {
3307 let arrow_schema = ArrowSchema::new(vec![ArrowField::new("s", DataType::Utf8, true)]);
3308 let a = AvroSchema::try_from(&arrow_schema).unwrap().json_string;
3309 let b = AvroSchema::from_arrow_with_options(&arrow_schema, None);
3310 assert_eq!(a, b.unwrap().json_string);
3311 }
3312
3313 #[test]
3314 fn test_union_branch_missing_name_errors() {
3315 for t in ["record", "enum", "fixed"] {
3316 let branch = json!({ "type": t });
3317 let err = union_branch_signature(&branch).unwrap_err().to_string();
3318 assert!(
3319 err.contains(&format!("Union branch '{t}' missing required 'name'")),
3320 "expected missing-name error for {t}, got: {err}"
3321 );
3322 }
3323 }
3324
3325 #[test]
3326 fn test_union_branch_named_type_signature_includes_name() {
3327 let rec = json!({ "type": "record", "name": "Foo" });
3328 assert_eq!(union_branch_signature(&rec).unwrap(), "N:record:Foo");
3329 let en = json!({ "type": "enum", "name": "Color", "symbols": ["R", "G", "B"] });
3330 assert_eq!(union_branch_signature(&en).unwrap(), "N:enum:Color");
3331 let fx = json!({ "type": "fixed", "name": "Bytes16", "size": 16 });
3332 assert_eq!(union_branch_signature(&fx).unwrap(), "N:fixed:Bytes16");
3333 }
3334
3335 #[test]
3336 fn test_record_field_alias_resolution_without_default() {
3337 let writer_json = r#"{
3338 "type":"record",
3339 "name":"R",
3340 "fields":[{"name":"old","type":"int"}]
3341 }"#;
3342 let reader_json = r#"{
3343 "type":"record",
3344 "name":"R",
3345 "fields":[{"name":"new","aliases":["old"],"type":"int"}]
3346 }"#;
3347 let writer: Schema = serde_json::from_str(writer_json).unwrap();
3348 let reader: Schema = serde_json::from_str(reader_json).unwrap();
3349 let resolved = AvroFieldBuilder::new(&writer)
3350 .with_reader_schema(&reader)
3351 .with_utf8view(false)
3352 .with_strict_mode(false)
3353 .build()
3354 .unwrap();
3355 let expected = ArrowField::new(
3356 "R",
3357 DataType::Struct(Fields::from(vec![ArrowField::new(
3358 "new",
3359 DataType::Int32,
3360 false,
3361 )])),
3362 false,
3363 )
3364 .with_metadata(HashMap::from_iter([(
3365 "avro.name".to_owned(),
3366 "R".to_owned(),
3367 )]));
3368 assert_eq!(resolved.field(), expected);
3369 }
3370
3371 #[test]
3372 fn test_record_field_alias_ambiguous_in_strict_mode_errors() {
3373 let writer_json = r#"{
3374 "type":"record",
3375 "name":"R",
3376 "fields":[
3377 {"name":"a","type":"int","aliases":["old"]},
3378 {"name":"b","type":"int","aliases":["old"]}
3379 ]
3380 }"#;
3381 let reader_json = r#"{
3382 "type":"record",
3383 "name":"R",
3384 "fields":[{"name":"target","type":"int","aliases":["old"]}]
3385 }"#;
3386 let writer: Schema = serde_json::from_str(writer_json).unwrap();
3387 let reader: Schema = serde_json::from_str(reader_json).unwrap();
3388 let err = AvroFieldBuilder::new(&writer)
3389 .with_reader_schema(&reader)
3390 .with_utf8view(false)
3391 .with_strict_mode(true)
3392 .build()
3393 .unwrap_err()
3394 .to_string();
3395 assert!(
3396 err.contains("Ambiguous alias 'old'"),
3397 "expected ambiguous-alias error, got: {err}"
3398 );
3399 }
3400
3401 #[test]
3402 fn test_pragmatic_writer_field_alias_mapping_non_strict() {
3403 let writer_json = r#"{
3404 "type":"record",
3405 "name":"R",
3406 "fields":[{"name":"before","type":"int","aliases":["now"]}]
3407 }"#;
3408 let reader_json = r#"{
3409 "type":"record",
3410 "name":"R",
3411 "fields":[{"name":"now","type":"int"}]
3412 }"#;
3413 let writer: Schema = serde_json::from_str(writer_json).unwrap();
3414 let reader: Schema = serde_json::from_str(reader_json).unwrap();
3415 let resolved = AvroFieldBuilder::new(&writer)
3416 .with_reader_schema(&reader)
3417 .with_utf8view(false)
3418 .with_strict_mode(false)
3419 .build()
3420 .unwrap();
3421 let expected = ArrowField::new(
3422 "R",
3423 DataType::Struct(Fields::from(vec![ArrowField::new(
3424 "now",
3425 DataType::Int32,
3426 false,
3427 )])),
3428 false,
3429 )
3430 .with_metadata(HashMap::from_iter([(
3431 "avro.name".to_owned(),
3432 "R".to_owned(),
3433 )]));
3434 assert_eq!(resolved.field(), expected);
3435 }
3436
3437 #[test]
3438 fn test_missing_reader_field_null_first_no_default_is_ok() {
3439 let writer_json = r#"{
3440 "type":"record",
3441 "name":"R",
3442 "fields":[{"name":"a","type":"int"}]
3443 }"#;
3444 let reader_json = r#"{
3445 "type":"record",
3446 "name":"R",
3447 "fields":[
3448 {"name":"a","type":"int"},
3449 {"name":"b","type":["null","int"]}
3450 ]
3451 }"#;
3452 let writer: Schema = serde_json::from_str(writer_json).unwrap();
3453 let reader: Schema = serde_json::from_str(reader_json).unwrap();
3454 let resolved = AvroFieldBuilder::new(&writer)
3455 .with_reader_schema(&reader)
3456 .with_utf8view(false)
3457 .with_strict_mode(false)
3458 .build()
3459 .unwrap();
3460 let expected = ArrowField::new(
3461 "R",
3462 DataType::Struct(Fields::from(vec![
3463 ArrowField::new("a", DataType::Int32, false),
3464 ArrowField::new("b", DataType::Int32, true).with_metadata(HashMap::from([(
3465 AVRO_FIELD_DEFAULT_METADATA_KEY.to_string(),
3466 "null".to_string(),
3467 )])),
3468 ])),
3469 false,
3470 )
3471 .with_metadata(HashMap::from_iter([(
3472 "avro.name".to_owned(),
3473 "R".to_owned(),
3474 )]));
3475 assert_eq!(resolved.field(), expected);
3476 }
3477
3478 #[test]
3479 fn test_missing_reader_field_null_second_without_default_errors() {
3480 let writer_json = r#"{
3481 "type":"record",
3482 "name":"R",
3483 "fields":[{"name":"a","type":"int"}]
3484 }"#;
3485 let reader_json = r#"{
3486 "type":"record",
3487 "name":"R",
3488 "fields":[
3489 {"name":"a","type":"int"},
3490 {"name":"b","type":["int","null"]}
3491 ]
3492 }"#;
3493 let writer: Schema = serde_json::from_str(writer_json).unwrap();
3494 let reader: Schema = serde_json::from_str(reader_json).unwrap();
3495 let err = AvroFieldBuilder::new(&writer)
3496 .with_reader_schema(&reader)
3497 .with_utf8view(false)
3498 .with_strict_mode(false)
3499 .build()
3500 .unwrap_err()
3501 .to_string();
3502 assert!(
3503 err.contains("must have a default value"),
3504 "expected missing-default error, got: {err}"
3505 );
3506 }
3507
3508 #[test]
3509 fn test_from_arrow_with_options_respects_schema_metadata_when_not_stripping() {
3510 let field = ArrowField::new("x", DataType::Int32, true);
3511 let injected_json =
3512 r#"{"type":"record","name":"Injected","fields":[{"name":"ignored","type":"int"}]}"#
3513 .to_string();
3514 let mut md = HashMap::new();
3515 md.insert(SCHEMA_METADATA_KEY.to_string(), injected_json.clone());
3516 md.insert("custom".to_string(), "123".to_string());
3517 let arrow_schema = ArrowSchema::new_with_metadata(vec![field], md);
3518 let opts = AvroSchemaOptions {
3519 null_order: Some(Nullability::NullSecond),
3520 strip_metadata: false,
3521 };
3522 let out = AvroSchema::from_arrow_with_options(&arrow_schema, Some(opts)).unwrap();
3523 assert_eq!(
3524 out.json_string, injected_json,
3525 "When strip_metadata=false and avro.schema is present, return the embedded JSON verbatim"
3526 );
3527 let v: Value = serde_json::from_str(&out.json_string).unwrap();
3528 assert_eq!(v.get("type").and_then(|t| t.as_str()), Some("record"));
3529 assert_eq!(v.get("name").and_then(|n| n.as_str()), Some("Injected"));
3530 }
3531
3532 #[test]
3533 fn test_from_arrow_with_options_ignores_schema_metadata_when_stripping_and_keeps_passthrough() {
3534 let field = ArrowField::new("x", DataType::Int32, true);
3535 let injected_json =
3536 r#"{"type":"record","name":"Injected","fields":[{"name":"ignored","type":"int"}]}"#
3537 .to_string();
3538 let mut md = HashMap::new();
3539 md.insert(SCHEMA_METADATA_KEY.to_string(), injected_json);
3540 md.insert("custom_meta".to_string(), "7".to_string());
3541 let arrow_schema = ArrowSchema::new_with_metadata(vec![field], md);
3542 let opts = AvroSchemaOptions {
3543 null_order: Some(Nullability::NullFirst),
3544 strip_metadata: true,
3545 };
3546 let out = AvroSchema::from_arrow_with_options(&arrow_schema, Some(opts)).unwrap();
3547 assert_json_contains(&out.json_string, "\"type\":\"record\"");
3548 assert_json_contains(&out.json_string, "\"name\":\"topLevelRecord\"");
3549 assert_json_contains(&out.json_string, "\"custom_meta\":7");
3550 }
3551
3552 #[test]
3553 fn test_from_arrow_with_options_null_first_for_nullable_primitive() {
3554 let field = ArrowField::new("s", DataType::Utf8, true);
3555 let arrow_schema = single_field_schema(field);
3556 let opts = AvroSchemaOptions {
3557 null_order: Some(Nullability::NullFirst),
3558 strip_metadata: true,
3559 };
3560 let out = AvroSchema::from_arrow_with_options(&arrow_schema, Some(opts)).unwrap();
3561 let v: Value = serde_json::from_str(&out.json_string).unwrap();
3562 let arr = v["fields"][0]["type"]
3563 .as_array()
3564 .expect("nullable primitive should be Avro union array");
3565 assert_eq!(arr[0], Value::String("null".into()));
3566 assert_eq!(arr[1], Value::String("string".into()));
3567 }
3568
3569 #[test]
3570 fn test_from_arrow_with_options_null_second_for_nullable_primitive() {
3571 let field = ArrowField::new("s", DataType::Utf8, true);
3572 let arrow_schema = single_field_schema(field);
3573 let opts = AvroSchemaOptions {
3574 null_order: Some(Nullability::NullSecond),
3575 strip_metadata: true,
3576 };
3577 let out = AvroSchema::from_arrow_with_options(&arrow_schema, Some(opts)).unwrap();
3578 let v: Value = serde_json::from_str(&out.json_string).unwrap();
3579 let arr = v["fields"][0]["type"]
3580 .as_array()
3581 .expect("nullable primitive should be Avro union array");
3582 assert_eq!(arr[0], Value::String("string".into()));
3583 assert_eq!(arr[1], Value::String("null".into()));
3584 }
3585
3586 #[test]
3587 fn test_from_arrow_with_options_union_extras_respected_by_strip_metadata() {
3588 let uf: UnionFields = vec![
3589 (2i8, Arc::new(ArrowField::new("a", DataType::Int32, false))),
3590 (7i8, Arc::new(ArrowField::new("b", DataType::Utf8, false))),
3591 ]
3592 .into_iter()
3593 .collect();
3594 let union_dt = DataType::Union(uf, UnionMode::Dense);
3595 let arrow_schema = single_field_schema(ArrowField::new("u", union_dt, true));
3596 let with_extras = AvroSchema::from_arrow_with_options(
3597 &arrow_schema,
3598 Some(AvroSchemaOptions {
3599 null_order: Some(Nullability::NullFirst),
3600 strip_metadata: false,
3601 }),
3602 )
3603 .unwrap();
3604 let v_with: Value = serde_json::from_str(&with_extras.json_string).unwrap();
3605 let union_arr = v_with["fields"][0]["type"].as_array().expect("union array");
3606 let first_obj = union_arr
3607 .iter()
3608 .find(|b| b.is_object())
3609 .expect("expected an object branch with extras");
3610 let obj = first_obj.as_object().unwrap();
3611 assert_eq!(obj.get("type").and_then(|t| t.as_str()), Some("int"));
3612 assert_eq!(
3613 obj.get("arrowUnionMode").and_then(|m| m.as_str()),
3614 Some("dense")
3615 );
3616 let type_ids: Vec<i64> = obj["arrowUnionTypeIds"]
3617 .as_array()
3618 .expect("arrowUnionTypeIds array")
3619 .iter()
3620 .map(|n| n.as_i64().expect("i64"))
3621 .collect();
3622 assert_eq!(type_ids, vec![2, 7]);
3623 let stripped = AvroSchema::from_arrow_with_options(
3624 &arrow_schema,
3625 Some(AvroSchemaOptions {
3626 null_order: Some(Nullability::NullFirst),
3627 strip_metadata: true,
3628 }),
3629 )
3630 .unwrap();
3631 let v_stripped: Value = serde_json::from_str(&stripped.json_string).unwrap();
3632 let union_arr2 = v_stripped["fields"][0]["type"]
3633 .as_array()
3634 .expect("union array");
3635 assert!(
3636 !union_arr2.iter().any(|b| b
3637 .as_object()
3638 .is_some_and(|m| m.contains_key("arrowUnionMode"))),
3639 "extras must be removed when strip_metadata=true"
3640 );
3641 assert_eq!(union_arr2[0], Value::String("null".into()));
3642 assert_eq!(union_arr2[1], Value::String("int".into()));
3643 assert_eq!(union_arr2[2], Value::String("string".into()));
3644 }
3645
3646 #[test]
3647 fn test_project_empty_projection() {
3648 let schema_json = r#"{
3649 "type": "record",
3650 "name": "Test",
3651 "fields": [
3652 {"name": "a", "type": "int"},
3653 {"name": "b", "type": "string"}
3654 ]
3655 }"#;
3656 let schema = AvroSchema::new(schema_json.to_string());
3657 let projected = schema.project(&[]).unwrap();
3658 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3659 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3660 assert!(
3661 fields.is_empty(),
3662 "Empty projection should yield empty fields"
3663 );
3664 }
3665
3666 #[test]
3667 fn test_project_single_field() {
3668 let schema_json = r#"{
3669 "type": "record",
3670 "name": "Test",
3671 "fields": [
3672 {"name": "a", "type": "int"},
3673 {"name": "b", "type": "string"},
3674 {"name": "c", "type": "long"}
3675 ]
3676 }"#;
3677 let schema = AvroSchema::new(schema_json.to_string());
3678 let projected = schema.project(&[1]).unwrap();
3679 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3680 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3681 assert_eq!(fields.len(), 1);
3682 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("b"));
3683 }
3684
3685 #[test]
3686 fn test_project_multiple_fields() {
3687 let schema_json = r#"{
3688 "type": "record",
3689 "name": "Test",
3690 "fields": [
3691 {"name": "a", "type": "int"},
3692 {"name": "b", "type": "string"},
3693 {"name": "c", "type": "long"},
3694 {"name": "d", "type": "boolean"}
3695 ]
3696 }"#;
3697 let schema = AvroSchema::new(schema_json.to_string());
3698 let projected = schema.project(&[0, 2, 3]).unwrap();
3699 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3700 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3701 assert_eq!(fields.len(), 3);
3702 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("a"));
3703 assert_eq!(fields[1].get("name").and_then(|n| n.as_str()), Some("c"));
3704 assert_eq!(fields[2].get("name").and_then(|n| n.as_str()), Some("d"));
3705 }
3706
3707 #[test]
3708 fn test_project_all_fields() {
3709 let schema_json = r#"{
3710 "type": "record",
3711 "name": "Test",
3712 "fields": [
3713 {"name": "a", "type": "int"},
3714 {"name": "b", "type": "string"}
3715 ]
3716 }"#;
3717 let schema = AvroSchema::new(schema_json.to_string());
3718 let projected = schema.project(&[0, 1]).unwrap();
3719 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3720 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3721 assert_eq!(fields.len(), 2);
3722 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("a"));
3723 assert_eq!(fields[1].get("name").and_then(|n| n.as_str()), Some("b"));
3724 }
3725
3726 #[test]
3727 fn test_project_reorder_fields() {
3728 let schema_json = r#"{
3729 "type": "record",
3730 "name": "Test",
3731 "fields": [
3732 {"name": "a", "type": "int"},
3733 {"name": "b", "type": "string"},
3734 {"name": "c", "type": "long"}
3735 ]
3736 }"#;
3737 let schema = AvroSchema::new(schema_json.to_string());
3738 let projected = schema.project(&[2, 0, 1]).unwrap();
3740 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3741 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3742 assert_eq!(fields.len(), 3);
3743 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("c"));
3744 assert_eq!(fields[1].get("name").and_then(|n| n.as_str()), Some("a"));
3745 assert_eq!(fields[2].get("name").and_then(|n| n.as_str()), Some("b"));
3746 }
3747
3748 #[test]
3749 fn test_project_preserves_record_metadata() {
3750 let schema_json = r#"{
3751 "type": "record",
3752 "name": "MyRecord",
3753 "namespace": "com.example",
3754 "doc": "A test record",
3755 "aliases": ["OldRecord"],
3756 "fields": [
3757 {"name": "a", "type": "int"},
3758 {"name": "b", "type": "string"}
3759 ]
3760 }"#;
3761 let schema = AvroSchema::new(schema_json.to_string());
3762 let projected = schema.project(&[0]).unwrap();
3763 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3764 assert_eq!(v.get("name").and_then(|n| n.as_str()), Some("MyRecord"));
3765 assert_eq!(
3766 v.get("namespace").and_then(|n| n.as_str()),
3767 Some("com.example")
3768 );
3769 assert_eq!(v.get("doc").and_then(|n| n.as_str()), Some("A test record"));
3770 assert!(v.get("aliases").is_some());
3771 }
3772
3773 #[test]
3774 fn test_project_preserves_field_metadata() {
3775 let schema_json = r#"{
3776 "type": "record",
3777 "name": "Test",
3778 "fields": [
3779 {"name": "a", "type": "int", "doc": "Field A", "default": 0},
3780 {"name": "b", "type": "string"}
3781 ]
3782 }"#;
3783 let schema = AvroSchema::new(schema_json.to_string());
3784 let projected = schema.project(&[0]).unwrap();
3785 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3786 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3787 assert_eq!(
3788 fields[0].get("doc").and_then(|d| d.as_str()),
3789 Some("Field A")
3790 );
3791 assert_eq!(fields[0].get("default").and_then(|d| d.as_i64()), Some(0));
3792 }
3793
3794 #[test]
3795 fn test_project_with_nested_record() {
3796 let schema_json = r#"{
3797 "type": "record",
3798 "name": "Outer",
3799 "fields": [
3800 {"name": "id", "type": "int"},
3801 {"name": "inner", "type": {
3802 "type": "record",
3803 "name": "Inner",
3804 "fields": [
3805 {"name": "x", "type": "int"},
3806 {"name": "y", "type": "string"}
3807 ]
3808 }},
3809 {"name": "value", "type": "double"}
3810 ]
3811 }"#;
3812 let schema = AvroSchema::new(schema_json.to_string());
3813 let projected = schema.project(&[1]).unwrap();
3814 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3815 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3816 assert_eq!(fields.len(), 1);
3817 assert_eq!(
3818 fields[0].get("name").and_then(|n| n.as_str()),
3819 Some("inner")
3820 );
3821 let inner_type = fields[0].get("type").unwrap();
3823 assert_eq!(
3824 inner_type.get("type").and_then(|t| t.as_str()),
3825 Some("record")
3826 );
3827 assert_eq!(
3828 inner_type.get("name").and_then(|n| n.as_str()),
3829 Some("Inner")
3830 );
3831 }
3832
3833 #[test]
3834 fn test_project_with_complex_field_types() {
3835 let schema_json = r#"{
3836 "type": "record",
3837 "name": "Test",
3838 "fields": [
3839 {"name": "arr", "type": {"type": "array", "items": "int"}},
3840 {"name": "map", "type": {"type": "map", "values": "string"}},
3841 {"name": "union", "type": ["null", "int"]}
3842 ]
3843 }"#;
3844 let schema = AvroSchema::new(schema_json.to_string());
3845 let projected = schema.project(&[0, 2]).unwrap();
3846 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
3847 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
3848 assert_eq!(fields.len(), 2);
3849 let arr_type = fields[0].get("type").unwrap();
3851 assert_eq!(arr_type.get("type").and_then(|t| t.as_str()), Some("array"));
3852 let union_type = fields[1].get("type").unwrap();
3854 assert!(union_type.is_array());
3855 }
3856
3857 #[test]
3858 fn test_project_error_invalid_json() {
3859 let schema = AvroSchema::new("not valid json".to_string());
3860 let err = schema.project(&[0]).unwrap_err();
3861 let msg = err.to_string();
3862 assert!(
3863 msg.contains("Invalid Avro schema JSON"),
3864 "Expected parse error, got: {msg}"
3865 );
3866 }
3867
3868 #[test]
3869 fn test_project_error_not_object() {
3870 let schema = AvroSchema::new(r#""string""#.to_string());
3872 let err = schema.project(&[0]).unwrap_err();
3873 let msg = err.to_string();
3874 assert!(
3875 msg.contains("must be a JSON object"),
3876 "Expected object error, got: {msg}"
3877 );
3878 }
3879
3880 #[test]
3881 fn test_project_error_array_schema() {
3882 let schema = AvroSchema::new(r#"["null", "int"]"#.to_string());
3884 let err = schema.project(&[0]).unwrap_err();
3885 let msg = err.to_string();
3886 assert!(
3887 msg.contains("must be a JSON object"),
3888 "Expected object error for array schema, got: {msg}"
3889 );
3890 }
3891
3892 #[test]
3893 fn test_project_error_type_not_record() {
3894 let schema_json = r#"{
3895 "type": "enum",
3896 "name": "Color",
3897 "symbols": ["RED", "GREEN", "BLUE"]
3898 }"#;
3899 let schema = AvroSchema::new(schema_json.to_string());
3900 let err = schema.project(&[0]).unwrap_err();
3901 let msg = err.to_string();
3902 assert!(
3903 msg.contains("must be an Avro record") && msg.contains("'enum'"),
3904 "Expected type mismatch error, got: {msg}"
3905 );
3906 }
3907
3908 #[test]
3909 fn test_project_error_type_array() {
3910 let schema_json = r#"{
3911 "type": "array",
3912 "items": "int"
3913 }"#;
3914 let schema = AvroSchema::new(schema_json.to_string());
3915 let err = schema.project(&[0]).unwrap_err();
3916 let msg = err.to_string();
3917 assert!(
3918 msg.contains("must be an Avro record") && msg.contains("'array'"),
3919 "Expected type mismatch error for array type, got: {msg}"
3920 );
3921 }
3922
3923 #[test]
3924 fn test_project_error_type_fixed() {
3925 let schema_json = r#"{
3926 "type": "fixed",
3927 "name": "MD5",
3928 "size": 16
3929 }"#;
3930 let schema = AvroSchema::new(schema_json.to_string());
3931 let err = schema.project(&[0]).unwrap_err();
3932 let msg = err.to_string();
3933 assert!(
3934 msg.contains("must be an Avro record") && msg.contains("'fixed'"),
3935 "Expected type mismatch error for fixed type, got: {msg}"
3936 );
3937 }
3938
3939 #[test]
3940 fn test_project_error_type_map() {
3941 let schema_json = r#"{
3942 "type": "map",
3943 "values": "string"
3944 }"#;
3945 let schema = AvroSchema::new(schema_json.to_string());
3946 let err = schema.project(&[0]).unwrap_err();
3947 let msg = err.to_string();
3948 assert!(
3949 msg.contains("must be an Avro record") && msg.contains("'map'"),
3950 "Expected type mismatch error for map type, got: {msg}"
3951 );
3952 }
3953
3954 #[test]
3955 fn test_project_error_missing_type_field() {
3956 let schema_json = r#"{
3957 "name": "Test",
3958 "fields": [{"name": "a", "type": "int"}]
3959 }"#;
3960 let schema = AvroSchema::new(schema_json.to_string());
3961 let err = schema.project(&[0]).unwrap_err();
3962 let msg = err.to_string();
3963 assert!(
3964 msg.contains("missing required 'type' field"),
3965 "Expected missing type error, got: {msg}"
3966 );
3967 }
3968
3969 #[test]
3970 fn test_project_error_missing_fields() {
3971 let schema_json = r#"{
3972 "type": "record",
3973 "name": "Test"
3974 }"#;
3975 let schema = AvroSchema::new(schema_json.to_string());
3976 let err = schema.project(&[0]).unwrap_err();
3977 let msg = err.to_string();
3978 assert!(
3979 msg.contains("missing required 'fields'"),
3980 "Expected missing fields error, got: {msg}"
3981 );
3982 }
3983
3984 #[test]
3985 fn test_project_error_fields_not_array() {
3986 let schema_json = r#"{
3987 "type": "record",
3988 "name": "Test",
3989 "fields": "not an array"
3990 }"#;
3991 let schema = AvroSchema::new(schema_json.to_string());
3992 let err = schema.project(&[0]).unwrap_err();
3993 let msg = err.to_string();
3994 assert!(
3995 msg.contains("'fields' must be an array"),
3996 "Expected fields array error, got: {msg}"
3997 );
3998 }
3999
4000 #[test]
4001 fn test_project_error_index_out_of_bounds() {
4002 let schema_json = r#"{
4003 "type": "record",
4004 "name": "Test",
4005 "fields": [
4006 {"name": "a", "type": "int"},
4007 {"name": "b", "type": "string"}
4008 ]
4009 }"#;
4010 let schema = AvroSchema::new(schema_json.to_string());
4011 let err = schema.project(&[5]).unwrap_err();
4012 let msg = err.to_string();
4013 assert!(
4014 msg.contains("out of bounds") && msg.contains('5') && msg.contains('2'),
4015 "Expected out of bounds error, got: {msg}"
4016 );
4017 }
4018
4019 #[test]
4020 fn test_project_error_index_out_of_bounds_edge() {
4021 let schema_json = r#"{
4022 "type": "record",
4023 "name": "Test",
4024 "fields": [
4025 {"name": "a", "type": "int"}
4026 ]
4027 }"#;
4028 let schema = AvroSchema::new(schema_json.to_string());
4029 let err = schema.project(&[1]).unwrap_err();
4031 let msg = err.to_string();
4032 assert!(
4033 msg.contains("out of bounds") && msg.contains('1'),
4034 "Expected out of bounds error for edge case, got: {msg}"
4035 );
4036 }
4037
4038 #[test]
4039 fn test_project_error_duplicate_index() {
4040 let schema_json = r#"{
4041 "type": "record",
4042 "name": "Test",
4043 "fields": [
4044 {"name": "a", "type": "int"},
4045 {"name": "b", "type": "string"},
4046 {"name": "c", "type": "long"}
4047 ]
4048 }"#;
4049 let schema = AvroSchema::new(schema_json.to_string());
4050 let err = schema.project(&[0, 1, 0]).unwrap_err();
4051 let msg = err.to_string();
4052 assert!(
4053 msg.contains("Duplicate projection index") && msg.contains('0'),
4054 "Expected duplicate index error, got: {msg}"
4055 );
4056 }
4057
4058 #[test]
4059 fn test_project_error_duplicate_index_consecutive() {
4060 let schema_json = r#"{
4061 "type": "record",
4062 "name": "Test",
4063 "fields": [
4064 {"name": "a", "type": "int"},
4065 {"name": "b", "type": "string"}
4066 ]
4067 }"#;
4068 let schema = AvroSchema::new(schema_json.to_string());
4069 let err = schema.project(&[1, 1]).unwrap_err();
4070 let msg = err.to_string();
4071 assert!(
4072 msg.contains("Duplicate projection index") && msg.contains('1'),
4073 "Expected duplicate index error for consecutive duplicates, got: {msg}"
4074 );
4075 }
4076
4077 #[test]
4078 fn test_project_with_empty_fields() {
4079 let schema_json = r#"{
4080 "type": "record",
4081 "name": "EmptyRecord",
4082 "fields": []
4083 }"#;
4084 let schema = AvroSchema::new(schema_json.to_string());
4085 let projected = schema.project(&[]).unwrap();
4087 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
4088 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
4089 assert!(fields.is_empty());
4090 }
4091
4092 #[test]
4093 fn test_project_empty_fields_index_out_of_bounds() {
4094 let schema_json = r#"{
4095 "type": "record",
4096 "name": "EmptyRecord",
4097 "fields": []
4098 }"#;
4099 let schema = AvroSchema::new(schema_json.to_string());
4100 let err = schema.project(&[0]).unwrap_err();
4101 let msg = err.to_string();
4102 assert!(
4103 msg.contains("out of bounds") && msg.contains("0 fields"),
4104 "Expected out of bounds error for empty record, got: {msg}"
4105 );
4106 }
4107
4108 #[test]
4109 fn test_project_result_is_valid_avro_schema() {
4110 let schema_json = r#"{
4111 "type": "record",
4112 "name": "Test",
4113 "namespace": "com.example",
4114 "fields": [
4115 {"name": "id", "type": "long"},
4116 {"name": "name", "type": "string"},
4117 {"name": "active", "type": "boolean"}
4118 ]
4119 }"#;
4120 let schema = AvroSchema::new(schema_json.to_string());
4121 let projected = schema.project(&[0, 2]).unwrap();
4122 let parsed = projected.schema();
4124 assert!(parsed.is_ok(), "Projected schema should be valid Avro");
4125 match parsed.unwrap() {
4126 Schema::Complex(ComplexType::Record(r)) => {
4127 assert_eq!(r.name, "Test");
4128 assert_eq!(r.namespace, Some("com.example"));
4129 assert_eq!(r.fields.len(), 2);
4130 assert_eq!(r.fields[0].name, "id");
4131 assert_eq!(r.fields[1].name, "active");
4132 }
4133 _ => panic!("Expected Record schema"),
4134 }
4135 }
4136
4137 #[test]
4138 fn test_project_non_contiguous_indices() {
4139 let schema_json = r#"{
4140 "type": "record",
4141 "name": "Test",
4142 "fields": [
4143 {"name": "f0", "type": "int"},
4144 {"name": "f1", "type": "int"},
4145 {"name": "f2", "type": "int"},
4146 {"name": "f3", "type": "int"},
4147 {"name": "f4", "type": "int"}
4148 ]
4149 }"#;
4150 let schema = AvroSchema::new(schema_json.to_string());
4151 let projected = schema.project(&[0, 2, 4]).unwrap();
4153 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
4154 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
4155 assert_eq!(fields.len(), 3);
4156 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("f0"));
4157 assert_eq!(fields[1].get("name").and_then(|n| n.as_str()), Some("f2"));
4158 assert_eq!(fields[2].get("name").and_then(|n| n.as_str()), Some("f4"));
4159 }
4160
4161 #[test]
4162 fn test_project_single_field_from_many() {
4163 let schema_json = r#"{
4164 "type": "record",
4165 "name": "BigRecord",
4166 "fields": [
4167 {"name": "f0", "type": "int"},
4168 {"name": "f1", "type": "int"},
4169 {"name": "f2", "type": "int"},
4170 {"name": "f3", "type": "int"},
4171 {"name": "f4", "type": "int"},
4172 {"name": "f5", "type": "int"},
4173 {"name": "f6", "type": "int"},
4174 {"name": "f7", "type": "int"},
4175 {"name": "f8", "type": "int"},
4176 {"name": "f9", "type": "int"}
4177 ]
4178 }"#;
4179 let schema = AvroSchema::new(schema_json.to_string());
4180 let projected = schema.project(&[9]).unwrap();
4182 let v: Value = serde_json::from_str(&projected.json_string).unwrap();
4183 let fields = v.get("fields").and_then(|f| f.as_array()).unwrap();
4184 assert_eq!(fields.len(), 1);
4185 assert_eq!(fields[0].get("name").and_then(|n| n.as_str()), Some("f9"));
4186 }
4187}