From 6b3174467e20a63e8400bfe7f22a884dabf71312 Mon Sep 17 00:00:00 2001 From: Ariel Miculas Date: Mon, 23 Mar 2026 16:03:59 +0200 Subject: [PATCH 1/4] fix: use writer types in Skipper for resolved named record types When a writer-only field references a named Avro type that was previously resolved against a reader schema, `parse_type` returns the registered reader-resolved type from the shared resolver. This caused two problems: 1. The Skipper built its struct sub-skippers from the reader's field list, which omits writer-only fields. Their bytes were never consumed, leaving the cursor at the wrong position for all subsequent records. 2. Reader fields carry resolution-induced nullability (e.g. a writer plain `long` matched against a reader `["null", long]` gains `nullability = Some(NullFirst)`). The Skipper read a union-tag byte that was never written, causing "Unexpected EOF" errors. Fix: store the writer's data type in `ResolvedField::ToReader` alongside the reader index. The Skipper's `Codec::Struct` arm now iterates `rec.writer_fields` and uses the writer type from every entry - both `ToReader(_, wdt)` and `Skip(wdt)` - so it always follows the writer's wire format. --- arrow-avro/src/codec.rs | 16 +- arrow-avro/src/reader/record.rs | 576 +++++++++++++++++++++++++++++++- 2 files changed, 577 insertions(+), 15 deletions(-) diff --git a/arrow-avro/src/codec.rs b/arrow-avro/src/codec.rs index 92a0ed051951..455d4c2f71ae 100644 --- a/arrow-avro/src/codec.rs +++ b/arrow-avro/src/codec.rs @@ -94,7 +94,9 @@ pub(crate) struct ResolvedRecord { #[derive(Debug, Clone, PartialEq)] pub(crate) enum ResolvedField { /// Resolves to a field indexed in the reader schema. - ToReader(usize), + /// The `AvroDataType` is the writer's type for this field, used by the Skipper + /// to correctly consume writer bytes when the whole record is being skipped. + ToReader(usize, AvroDataType), /// For fields present in the writer's schema but not the reader's, this stores their data type. /// This is needed to correctly skip over these fields during deserialization. Skip(AvroDataType), @@ -2341,10 +2343,10 @@ impl<'a> Maker<'a> { .iter() .enumerate() .map(|(writer_index, writer_field)| { + let dt = self.parse_type(&writer_field.r#type, writer_ns)?; if let Some(reader_index) = writer_to_reader[writer_index] { - Ok(ResolvedField::ToReader(reader_index)) + Ok(ResolvedField::ToReader(reader_index, dt)) } else { - let dt = self.parse_type(&writer_field.r#type, writer_ns)?; Ok(ResolvedField::Skip(dt)) } }) @@ -2888,7 +2890,7 @@ mod tests { default_fields, }) => { assert_eq!(writer_fields.len(), 1); - assert_eq!(writer_fields[0], ResolvedField::ToReader(0)); + assert!(matches!(writer_fields[0], ResolvedField::ToReader(0, _))); assert_eq!(default_fields.len(), 1); assert_eq!(default_fields[0], 1); } @@ -2981,7 +2983,7 @@ mod tests { default_fields, }) => { assert_eq!(writer_fields.len(), 1); - assert_eq!(writer_fields[0], ResolvedField::ToReader(0)); + assert!(matches!(writer_fields[0], ResolvedField::ToReader(0, _))); assert_eq!(default_fields.len(), 1); assert_eq!(default_fields[0], 1); } @@ -3802,9 +3804,9 @@ mod tests { assert!(matches!( &rec.writer_fields[..], &[ - ResolvedField::ToReader(1), + ResolvedField::ToReader(1, _), ResolvedField::Skip(_), - ResolvedField::ToReader(0), + ResolvedField::ToReader(0, _), ] )); assert_eq!(rec.default_fields.as_ref(), &[2usize, 3usize]); diff --git a/arrow-avro/src/reader/record.rs b/arrow-avro/src/reader/record.rs index 97cdeed20fc6..52ad777be721 100644 --- a/arrow-avro/src/reader/record.rs +++ b/arrow-avro/src/reader/record.rs @@ -2458,7 +2458,7 @@ impl<'a> ProjectorBuilder<'a> { .writer_fields .iter() .map(|field| match field { - ResolvedField::ToReader(index) => Ok(FieldProjection::ToReader(*index)), + ResolvedField::ToReader(index, _) => Ok(FieldProjection::ToReader(*index)), ResolvedField::Skip(datatype) => { let skipper = Skipper::from_avro(datatype)?; Ok(FieldProjection::Skip(skipper)) @@ -2568,12 +2568,27 @@ impl Skipper { Codec::Uuid => Self::UuidString, // encoded as string Codec::Enum(_) => Self::Enum, Codec::List(item) => Self::List(Box::new(Skipper::from_avro(item)?)), - Codec::Struct(fields) => Self::Struct( - fields - .iter() - .map(|f| Skipper::from_avro(f.data_type())) - .collect::>()?, - ), + Codec::Struct(fields) => { + if let Some(ResolutionInfo::Record(rec)) = dt.resolution.as_ref() { + Self::Struct( + rec.writer_fields + .iter() + .map(|wf| match wf { + ResolvedField::ToReader(_, wdt) | ResolvedField::Skip(wdt) => { + Skipper::from_avro(wdt) + } + }) + .collect::>()?, + ) + } else { + Self::Struct( + fields + .iter() + .map(|f| Skipper::from_avro(f.data_type())) + .collect::>()?, + ) + } + } Codec::Map(values) => Self::Map(Box::new(Skipper::from_avro(values)?)), Codec::Interval => Self::DurationFixed12, Codec::Union(encodings, _, _) => { @@ -2715,7 +2730,8 @@ impl Skipper { mod tests { use super::*; use crate::codec::AvroFieldBuilder; - use crate::schema::{Attributes, ComplexType, Field, PrimitiveType, Record, Schema, TypeName}; + use crate::schema::{Attributes, ComplexType, Enum as AvroEnum, Field, PrimitiveType, Record, Schema, TypeName}; + use crate::schema::Array as AvroArray; use arrow_array::cast::AsArray; use indexmap::IndexMap; use std::collections::HashMap; @@ -5611,4 +5627,548 @@ mod tests { other => panic!("expected Timestamp(Nanosecond, None), got {other:?}"), } } + + /// When a Skip field references a named type by name, the Skipper must use the + /// writer's wire format for that type. Here the reader wraps the Timestamp's scalar + /// fields in nullable unions, but the writer wrote them as plain values; the Skipper + /// must not add a union-tag read for each field. + #[test] + fn test_skip_named_type_ref_uses_writer_schema_not_resolved() { + let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); + let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + + // Writer: Timestamp{seconds:long, nanos:int} (plain, no union wrappers) + let timestamp_writer = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, + doc: None, + aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer: Event{time:"Timestamp", code:int} — `time` is a named-type reference + let event_writer = Schema::Complex(ComplexType::Record(Record { + name: "Event", + namespace: None, + doc: None, + aliases: vec![], + fields: vec![ + Field { name: "time", r#type: Schema::TypeName(TypeName::Ref("Timestamp")), default: None, doc: None, aliases: vec![] }, + Field { name: "code", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer root: { ts: Timestamp, events: array } + let writer_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, + doc: None, + aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, + Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Timestamp: inner scalars wrapped in ["null", T] + let timestamp_reader = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, + doc: None, + aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader root: only `ts` (nullable wrapper), `events` absent → Skip + let reader_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, + doc: None, + aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: Schema::Union(vec![null, timestamp_reader]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let field = AvroFieldBuilder::new(&writer_schema) + .with_reader_schema(&reader_schema) + .build() + .expect("schema resolution must succeed"); + + let mut decoder = Decoder::try_new(field.data_type()) + .expect("decoder must build"); + + // Encode one writer record: + // ts = Timestamp{seconds=100, nanos=5} — plain longs/ints, no union tags + // events = [Event{time=Timestamp{seconds=200, nanos=1}, code=42}] + let mut data = Vec::new(); + data.extend_from_slice(&encode_avro_long(100)); // ts.seconds + data.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds (plain) + data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos (plain) + data.extend_from_slice(&encode_avro_int(42)); // event.code + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + + let mut cur = AvroCursor::new(&data); + decoder + .decode(&mut cur) + .expect("decode must not corrupt cursor due to spurious Nullable wrappers"); + assert_eq!( + cur.position(), + data.len(), + "cursor must consume exactly the writer-encoded bytes" + ); + } + + + /// When a Skip field references a named type that has more fields in the writer + /// schema than in the reader schema, the Skipper must consume all writer fields + /// including writer-only ones. Here the writer's Timestamp has a third field + /// `tz_offset` that the reader omits; the Skipper must still consume those bytes. + #[test] + fn test_skip_named_type_ref_with_writer_extra_fields_consumes_all_bytes() { + let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); + let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + + // Writer Timestamp has an extra field `tz_offset` that the reader omits + let timestamp_writer = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + Field { name: "tz_offset", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let event_writer = Schema::Complex(ComplexType::Record(Record { + name: "Event", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "time", r#type: Schema::TypeName(TypeName::Ref("Timestamp")), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let writer_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, + Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Timestamp: only seconds and nanos (no tz_offset) + let timestamp_reader = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let reader_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: Schema::Union(vec![null, timestamp_reader]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let field = AvroFieldBuilder::new(&writer_schema) + .with_reader_schema(&reader_schema) + .build() + .expect("schema resolution must succeed"); + + let mut decoder = Decoder::try_new(field.data_type()) + .expect("decoder must build"); + + // Encode one writer record: + // ts = Timestamp{seconds=100, nanos=5, tz_offset=3600} + // events = [Event{time=Timestamp{seconds=200, nanos=1, tz_offset=7200}}] + // The Skipper for `events` must consume ALL three fields per Timestamp, + // including tz_offset, even though the resolved Timestamp only has two reader fields. + let mut data = Vec::new(); + data.extend_from_slice(&encode_avro_long(100)); // ts.seconds + data.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data.extend_from_slice(&encode_avro_int(3600)); // ts.tz_offset (writer-only, decoded via writer_fields) + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds + data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos + data.extend_from_slice(&encode_avro_int(7200)); // event.time.tz_offset (must be skipped!) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + + let mut cur = AvroCursor::new(&data); + decoder + .decode(&mut cur) + .expect("decode must not fail"); + assert_eq!( + cur.position(), + data.len(), + "cursor must consume all writer-encoded bytes including writer-only tz_offset in skipped events" + ); + } + + /// When a Skip field references a named type whose inner struct field was written + /// as a plain record but the reader wraps it in a nullable union, the Skipper must + /// use the writer's plain encoding and not read a union tag for that inner field. + #[test] + fn test_skip_nested_struct_record_resolution_no_spurious_nullable_wrapper() { + let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + + // Writer: InnerRecord { v: int } + let inner_writer = Schema::Complex(ComplexType::Record(Record { + name: "InnerRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "v", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer: OuterRecord { inner: InnerRecord } + let outer_writer = Schema::Complex(ComplexType::Record(Record { + name: "OuterRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "inner", r#type: inner_writer, default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer: Event { data: "OuterRecord" } (named type ref — will look up registered type) + let event_writer = Schema::Complex(ComplexType::Record(Record { + name: "Event", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "data", r#type: Schema::TypeName(TypeName::Ref("OuterRecord")), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer Root: { outer: OuterRecord, events: array } + // `events` is omitted from the reader schema → will be skipped. + let writer_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "outer", r#type: outer_writer, default: None, doc: None, aliases: vec![] }, + Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader: InnerRecord { v: int } + let inner_reader = Schema::Complex(ComplexType::Record(Record { + name: "InnerRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "v", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader: OuterRecord { inner: ["null", InnerRecord] } — inner wrapped in nullable + let outer_reader = Schema::Complex(ComplexType::Record(Record { + name: "OuterRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "inner", r#type: Schema::Union(vec![null.clone(), inner_reader]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Root: { outer: ["null", OuterRecord] } — events omitted (Skip) + let reader_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "outer", r#type: Schema::Union(vec![null, outer_reader]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let field = AvroFieldBuilder::new(&writer_schema) + .with_reader_schema(&reader_schema) + .build() + .expect("schema resolution must succeed"); + + let mut decoder = Decoder::try_new(field.data_type()) + .expect("decoder must build"); + + // Encode one writer record (all plain, no union tags anywhere): + // outer.inner.v = 42 + // events = [Event { data = OuterRecord { inner = InnerRecord { v = 7 } } }] + // The Skipper for `events` must skip InnerRecord as a plain struct (no union tag), + // even though the registered OuterRecord has inner with nullability+Record resolution. + let mut data = Vec::new(); + data.extend_from_slice(&encode_avro_int(42)); // outer.inner.v + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_int(7)); // event.data.inner.v (skipped, no union tag) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + + let mut cur = AvroCursor::new(&data); + decoder + .decode(&mut cur) + .expect("decode must not fail"); + assert_eq!( + cur.position(), + data.len(), + "cursor must consume all bytes; spurious Nullable wrapper would consume an extra byte" + ); + } + + /// When a Skip field references a named type whose inner enum field was written as a + /// plain index but the reader wraps it in a nullable union (with symbol reordering), + /// the Skipper must skip the enum as a plain VLQ int without reading a union tag. + #[test] + fn test_skip_enum_mapping_nullable_no_spurious_nullable_wrapper() { + let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); + + // Writer: SomeRecord { status: StatusEnum{OK, ERROR} } + let status_writer = Schema::Complex(ComplexType::Enum(AvroEnum { + name: "StatusEnum", + namespace: None, doc: None, aliases: vec![], + symbols: vec!["OK", "ERROR"], + default: None, + attributes: Attributes::default(), + })); + let some_record_writer = Schema::Complex(ComplexType::Record(Record { + name: "SomeRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "status", r#type: status_writer, default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer: Event { data: "SomeRecord" } (named type ref — looks up registered type) + let event_writer = Schema::Complex(ComplexType::Record(Record { + name: "Event", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "data", r#type: Schema::TypeName(TypeName::Ref("SomeRecord")), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer Root: { outer: SomeRecord, events: array } + // `events` is omitted from the reader schema → will be skipped. + let writer_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "outer", r#type: some_record_writer, default: None, doc: None, aliases: vec![] }, + Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader: SomeRecord { status: ["null", StatusEnum{ERROR, OK}] } (reordered + nullable) + let status_reader = Schema::Complex(ComplexType::Enum(AvroEnum { + name: "StatusEnum", + namespace: None, doc: None, aliases: vec![], + symbols: vec!["ERROR", "OK"], + default: None, + attributes: Attributes::default(), + })); + let some_record_reader = Schema::Complex(ComplexType::Record(Record { + name: "SomeRecord", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "status", r#type: Schema::Union(vec![null, status_reader]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Root: { outer: SomeRecord } — events omitted (Skip) + let reader_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "outer", r#type: some_record_reader, default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let field = AvroFieldBuilder::new(&writer_schema) + .with_reader_schema(&reader_schema) + .build() + .expect("schema resolution must succeed"); + + let mut decoder = Decoder::try_new(field.data_type()) + .expect("decoder must build"); + + // Encode one writer record (all plain, no union tags anywhere): + // outer.status = OK (writer index 0) + // events = [Event { data = SomeRecord { status = ERROR (writer index 1) } }] + // The Skipper for `events` must skip the enum as a plain VLQ int (no union tag). + // Using ERROR (index 1, encoded as 0x02) is deliberate: a spurious Nullable wrapper + // would read 0x02 as union tag 1 ("non-null") and then consume an additional byte + // for the inner value, reading into the end-of-array marker and leaving the cursor + // at the wrong position (or causing an EOF error). + let mut data = Vec::new(); + data.extend_from_slice(&encode_avro_int(0)); // outer.status: OK = writer index 0 + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_int(1)); // event.data.status: ERROR = writer index 1 (skipped) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + + let mut cur = AvroCursor::new(&data); + decoder + .decode(&mut cur) + .expect("decode must not fail"); + assert_eq!( + cur.position(), + data.len(), + "cursor must consume all bytes; spurious Nullable wrapper would consume an extra byte" + ); + } + + /// When a writer-only Skip field is a genuine nullable union `["null", "TypeRef"]`, + /// the Skipper must preserve the `Nullable` wrapper and read the union tag byte. + /// This is distinct from resolution-induced nullability: here the writer itself + /// wrote a union tag, so the tag byte must be consumed. + #[test] + fn test_skip_nullable_named_type_ref_union_not_stripped_of_nullable_wrapper() { + let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); + let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + + // Writer Timestamp: plain {seconds: long, nanos: int} + let timestamp_writer = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Writer Root: + // ts: Timestamp — matched by reader, triggers resolver registration + // extra: ["null", "Timestamp"] — writer-only genuine nullable union (will be Skipped) + let writer_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, + Field { name: "extra", r#type: Schema::Union(vec![ + null.clone(), + Schema::TypeName(TypeName::Ref("Timestamp")), + ]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Timestamp: fields are nullable (writer plain → reader nullable via resolution) + let timestamp_reader = Schema::Complex(ComplexType::Record(Record { + name: "Timestamp", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, + Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + // Reader Root: only `ts` (extra is writer-only → will be Skipped) + let reader_schema = Schema::Complex(ComplexType::Record(Record { + name: "Root", + namespace: None, doc: None, aliases: vec![], + fields: vec![ + Field { name: "ts", r#type: timestamp_reader, default: None, doc: None, aliases: vec![] }, + ], + attributes: Attributes::default(), + })); + + let field = AvroFieldBuilder::new(&writer_schema) + .with_reader_schema(&reader_schema) + .build() + .expect("schema resolution must succeed"); + + let mut decoder = Decoder::try_new(field.data_type()) + .expect("decoder must build"); + + // --- Test case 1: extra = null --- + // Writer encodes: + // ts.seconds = 100 (plain long, no union tag) + // ts.nanos = 5 (plain int, no union tag) + // extra = null (union tag 0 = 0x00) + // The Skipper for `extra` must read the union tag byte (0x00 → null branch, no payload). + // Before the fix, the Nullable wrapper was stripped and the Skipper treated `extra` as a + // plain Timestamp struct, consuming ts.seconds bytes of the next record as if they were + // struct fields — corrupting the cursor. + let mut data_null = Vec::new(); + data_null.extend_from_slice(&encode_avro_long(100)); // ts.seconds + data_null.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data_null.push(0x00); // extra: union tag 0 = null + + let mut cur = AvroCursor::new(&data_null); + decoder + .decode(&mut cur) + .expect("decode must not fail (null branch)"); + assert_eq!( + cur.position(), + data_null.len(), + "cursor must consume exactly the writer bytes including the null union tag" + ); + + // --- Test case 2: extra = non-null Timestamp --- + // Writer encodes: + // ts.seconds = 100 + // ts.nanos = 5 + // extra = Timestamp{seconds=200, nanos=3} + // → union tag 2 (0x04, zigzag for index 1) + seconds=200 + nanos=3 + // The Skipper for `extra` must read the union tag byte, then skip the Timestamp payload. + let mut data_nonnull = Vec::new(); + data_nonnull.extend_from_slice(&encode_avro_long(100)); // ts.seconds + data_nonnull.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data_nonnull.extend_from_slice(&encode_avro_int(1)); // extra: union tag zigzag(1)=0x02 + data_nonnull.extend_from_slice(&encode_avro_long(200)); // extra.seconds + data_nonnull.extend_from_slice(&encode_avro_int(3)); // extra.nanos + + let mut cur2 = AvroCursor::new(&data_nonnull); + decoder + .decode(&mut cur2) + .expect("decode must not fail (non-null branch)"); + assert_eq!( + cur2.position(), + data_nonnull.len(), + "cursor must consume all writer bytes including union tag and Timestamp payload" + ); + } } From 082b2f508fe81c9bcd8e1c57afe90737d7873870 Mon Sep 17 00:00:00 2001 From: Ariel Miculas Date: Tue, 31 Mar 2026 13:50:50 +0300 Subject: [PATCH 2/4] fix: formatting --- arrow-avro/src/reader/record.rs | 552 +++++++++++++++++++++++--------- 1 file changed, 397 insertions(+), 155 deletions(-) diff --git a/arrow-avro/src/reader/record.rs b/arrow-avro/src/reader/record.rs index 52ad777be721..9db0be385a7a 100644 --- a/arrow-avro/src/reader/record.rs +++ b/arrow-avro/src/reader/record.rs @@ -2730,8 +2730,10 @@ impl Skipper { mod tests { use super::*; use crate::codec::AvroFieldBuilder; - use crate::schema::{Attributes, ComplexType, Enum as AvroEnum, Field, PrimitiveType, Record, Schema, TypeName}; use crate::schema::Array as AvroArray; + use crate::schema::{ + Attributes, ComplexType, Enum as AvroEnum, Field, PrimitiveType, Record, Schema, TypeName, + }; use arrow_array::cast::AsArray; use indexmap::IndexMap; use std::collections::HashMap; @@ -5645,8 +5647,20 @@ mod tests { doc: None, aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: long.clone(), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5658,8 +5672,20 @@ mod tests { doc: None, aliases: vec![], fields: vec![ - Field { name: "time", r#type: Schema::TypeName(TypeName::Ref("Timestamp")), default: None, doc: None, aliases: vec![] }, - Field { name: "code", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + Field { + name: "time", + r#type: Schema::TypeName(TypeName::Ref("Timestamp")), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "code", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5671,11 +5697,23 @@ mod tests { doc: None, aliases: vec![], fields: vec![ - Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, - Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), default: None, doc: None, aliases: vec![] }, + Field { + name: "ts", + r#type: timestamp_writer, + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "events", + r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5687,8 +5725,20 @@ mod tests { doc: None, aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: Schema::Union(vec![null.clone(), long.clone()]), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: Schema::Union(vec![null.clone(), int.clone()]), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5699,9 +5749,13 @@ mod tests { namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "ts", r#type: Schema::Union(vec![null, timestamp_reader]), default: None, doc: None, aliases: vec![] }, - ], + fields: vec![Field { + name: "ts", + r#type: Schema::Union(vec![null, timestamp_reader]), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -5710,20 +5764,19 @@ mod tests { .build() .expect("schema resolution must succeed"); - let mut decoder = Decoder::try_new(field.data_type()) - .expect("decoder must build"); + let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); // Encode one writer record: // ts = Timestamp{seconds=100, nanos=5} — plain longs/ints, no union tags // events = [Event{time=Timestamp{seconds=200, nanos=1}, code=42}] let mut data = Vec::new(); data.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds (plain) - data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos (plain) - data.extend_from_slice(&encode_avro_int(42)); // event.code - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos (plain) + data.extend_from_slice(&encode_avro_int(42)); // event.code + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array let mut cur = AvroCursor::new(&data); decoder @@ -5736,7 +5789,6 @@ mod tests { ); } - /// When a Skip field references a named type that has more fields in the writer /// schema than in the reader schema, the Skipper must consume all writer fields /// including writer-only ones. Here the writer's Timestamp has a third field @@ -5745,38 +5797,78 @@ mod tests { fn test_skip_named_type_ref_with_writer_extra_fields_consumes_all_bytes() { let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); // Writer Timestamp has an extra field `tz_offset` that the reader omits let timestamp_writer = Schema::Complex(ComplexType::Record(Record { name: "Timestamp", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, - Field { name: "tz_offset", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: long.clone(), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "tz_offset", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); let event_writer = Schema::Complex(ComplexType::Record(Record { name: "Event", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "time", r#type: Schema::TypeName(TypeName::Ref("Timestamp")), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "time", + r#type: Schema::TypeName(TypeName::Ref("Timestamp")), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); let writer_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, - Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), default: None, doc: None, aliases: vec![] }, + Field { + name: "ts", + r#type: timestamp_writer, + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "events", + r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5784,20 +5876,40 @@ mod tests { // Reader Timestamp: only seconds and nanos (no tz_offset) let timestamp_reader = Schema::Complex(ComplexType::Record(Record { name: "Timestamp", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: Schema::Union(vec![null.clone(), long.clone()]), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: Schema::Union(vec![null.clone(), int.clone()]), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); let reader_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "ts", r#type: Schema::Union(vec![null, timestamp_reader]), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "ts", + r#type: Schema::Union(vec![null, timestamp_reader]), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -5806,8 +5918,7 @@ mod tests { .build() .expect("schema resolution must succeed"); - let mut decoder = Decoder::try_new(field.data_type()) - .expect("decoder must build"); + let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); // Encode one writer record: // ts = Timestamp{seconds=100, nanos=5, tz_offset=3600} @@ -5815,19 +5926,17 @@ mod tests { // The Skipper for `events` must consume ALL three fields per Timestamp, // including tz_offset, even though the resolved Timestamp only has two reader fields. let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data.extend_from_slice(&encode_avro_int(3600)); // ts.tz_offset (writer-only, decoded via writer_fields) - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds - data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos - data.extend_from_slice(&encode_avro_int(7200)); // event.time.tz_offset (must be skipped!) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + data.extend_from_slice(&encode_avro_long(100)); // ts.seconds + data.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data.extend_from_slice(&encode_avro_int(3600)); // ts.tz_offset (writer-only, decoded via writer_fields) + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds + data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos + data.extend_from_slice(&encode_avro_int(7200)); // event.time.tz_offset (must be skipped!) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array let mut cur = AvroCursor::new(&data); - decoder - .decode(&mut cur) - .expect("decode must not fail"); + decoder.decode(&mut cur).expect("decode must not fail"); assert_eq!( cur.position(), data.len(), @@ -5841,35 +5950,53 @@ mod tests { #[test] fn test_skip_nested_struct_record_resolution_no_spurious_nullable_wrapper() { let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); // Writer: InnerRecord { v: int } let inner_writer = Schema::Complex(ComplexType::Record(Record { name: "InnerRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "v", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "v", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Writer: OuterRecord { inner: InnerRecord } let outer_writer = Schema::Complex(ComplexType::Record(Record { name: "OuterRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "inner", r#type: inner_writer, default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "inner", + r#type: inner_writer, + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Writer: Event { data: "OuterRecord" } (named type ref — will look up registered type) let event_writer = Schema::Complex(ComplexType::Record(Record { name: "Event", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "data", r#type: Schema::TypeName(TypeName::Ref("OuterRecord")), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "data", + r#type: Schema::TypeName(TypeName::Ref("OuterRecord")), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -5877,13 +6004,27 @@ mod tests { // `events` is omitted from the reader schema → will be skipped. let writer_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "outer", r#type: outer_writer, default: None, doc: None, aliases: vec![] }, - Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), default: None, doc: None, aliases: vec![] }, + Field { + name: "outer", + r#type: outer_writer, + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "events", + r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -5891,30 +6032,48 @@ mod tests { // Reader: InnerRecord { v: int } let inner_reader = Schema::Complex(ComplexType::Record(Record { name: "InnerRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "v", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "v", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Reader: OuterRecord { inner: ["null", InnerRecord] } — inner wrapped in nullable let outer_reader = Schema::Complex(ComplexType::Record(Record { name: "OuterRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "inner", r#type: Schema::Union(vec![null.clone(), inner_reader]), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "inner", + r#type: Schema::Union(vec![null.clone(), inner_reader]), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Reader Root: { outer: ["null", OuterRecord] } — events omitted (Skip) let reader_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "outer", r#type: Schema::Union(vec![null, outer_reader]), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "outer", + r#type: Schema::Union(vec![null, outer_reader]), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -5923,8 +6082,7 @@ mod tests { .build() .expect("schema resolution must succeed"); - let mut decoder = Decoder::try_new(field.data_type()) - .expect("decoder must build"); + let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); // Encode one writer record (all plain, no union tags anywhere): // outer.inner.v = 42 @@ -5932,15 +6090,13 @@ mod tests { // The Skipper for `events` must skip InnerRecord as a plain struct (no union tag), // even though the registered OuterRecord has inner with nullability+Record resolution. let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_int(42)); // outer.inner.v - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_int(7)); // event.data.inner.v (skipped, no union tag) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + data.extend_from_slice(&encode_avro_int(42)); // outer.inner.v + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_int(7)); // event.data.inner.v (skipped, no union tag) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array let mut cur = AvroCursor::new(&data); - decoder - .decode(&mut cur) - .expect("decode must not fail"); + decoder.decode(&mut cur).expect("decode must not fail"); assert_eq!( cur.position(), data.len(), @@ -5958,27 +6114,41 @@ mod tests { // Writer: SomeRecord { status: StatusEnum{OK, ERROR} } let status_writer = Schema::Complex(ComplexType::Enum(AvroEnum { name: "StatusEnum", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], symbols: vec!["OK", "ERROR"], default: None, attributes: Attributes::default(), })); let some_record_writer = Schema::Complex(ComplexType::Record(Record { name: "SomeRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "status", r#type: status_writer, default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "status", + r#type: status_writer, + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Writer: Event { data: "SomeRecord" } (named type ref — looks up registered type) let event_writer = Schema::Complex(ComplexType::Record(Record { name: "Event", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "data", r#type: Schema::TypeName(TypeName::Ref("SomeRecord")), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "data", + r#type: Schema::TypeName(TypeName::Ref("SomeRecord")), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -5986,13 +6156,27 @@ mod tests { // `events` is omitted from the reader schema → will be skipped. let writer_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "outer", r#type: some_record_writer, default: None, doc: None, aliases: vec![] }, - Field { name: "events", r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), default: None, doc: None, aliases: vec![] }, + Field { + name: "outer", + r#type: some_record_writer, + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "events", + r#type: Schema::Complex(ComplexType::Array(AvroArray { + items: Box::new(event_writer), + attributes: Attributes::default(), + })), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -6000,27 +6184,41 @@ mod tests { // Reader: SomeRecord { status: ["null", StatusEnum{ERROR, OK}] } (reordered + nullable) let status_reader = Schema::Complex(ComplexType::Enum(AvroEnum { name: "StatusEnum", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], symbols: vec!["ERROR", "OK"], default: None, attributes: Attributes::default(), })); let some_record_reader = Schema::Complex(ComplexType::Record(Record { name: "SomeRecord", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "status", r#type: Schema::Union(vec![null, status_reader]), default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "status", + r#type: Schema::Union(vec![null, status_reader]), + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); // Reader Root: { outer: SomeRecord } — events omitted (Skip) let reader_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "outer", r#type: some_record_reader, default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "outer", + r#type: some_record_reader, + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -6029,8 +6227,7 @@ mod tests { .build() .expect("schema resolution must succeed"); - let mut decoder = Decoder::try_new(field.data_type()) - .expect("decoder must build"); + let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); // Encode one writer record (all plain, no union tags anywhere): // outer.status = OK (writer index 0) @@ -6041,15 +6238,13 @@ mod tests { // for the inner value, reading into the end-of-array marker and leaving the cursor // at the wrong position (or causing an EOF error). let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_int(0)); // outer.status: OK = writer index 0 - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_int(1)); // event.data.status: ERROR = writer index 1 (skipped) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array + data.extend_from_slice(&encode_avro_int(0)); // outer.status: OK = writer index 0 + data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 + data.extend_from_slice(&encode_avro_int(1)); // event.data.status: ERROR = writer index 1 (skipped) + data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array let mut cur = AvroCursor::new(&data); - decoder - .decode(&mut cur) - .expect("decode must not fail"); + decoder.decode(&mut cur).expect("decode must not fail"); assert_eq!( cur.position(), data.len(), @@ -6065,15 +6260,29 @@ mod tests { fn test_skip_nullable_named_type_ref_union_not_stripped_of_nullable_wrapper() { let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); + let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); // Writer Timestamp: plain {seconds: long, nanos: int} let timestamp_writer = Schema::Complex(ComplexType::Record(Record { name: "Timestamp", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: long.clone(), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: int.clone(), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: long.clone(), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: int.clone(), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -6083,13 +6292,27 @@ mod tests { // extra: ["null", "Timestamp"] — writer-only genuine nullable union (will be Skipped) let writer_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "ts", r#type: timestamp_writer, default: None, doc: None, aliases: vec![] }, - Field { name: "extra", r#type: Schema::Union(vec![ - null.clone(), - Schema::TypeName(TypeName::Ref("Timestamp")), - ]), default: None, doc: None, aliases: vec![] }, + Field { + name: "ts", + r#type: timestamp_writer, + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "extra", + r#type: Schema::Union(vec![ + null.clone(), + Schema::TypeName(TypeName::Ref("Timestamp")), + ]), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -6097,10 +6320,24 @@ mod tests { // Reader Timestamp: fields are nullable (writer plain → reader nullable via resolution) let timestamp_reader = Schema::Complex(ComplexType::Record(Record { name: "Timestamp", - namespace: None, doc: None, aliases: vec![], + namespace: None, + doc: None, + aliases: vec![], fields: vec![ - Field { name: "seconds", r#type: Schema::Union(vec![null.clone(), long.clone()]), default: None, doc: None, aliases: vec![] }, - Field { name: "nanos", r#type: Schema::Union(vec![null.clone(), int.clone()]), default: None, doc: None, aliases: vec![] }, + Field { + name: "seconds", + r#type: Schema::Union(vec![null.clone(), long.clone()]), + default: None, + doc: None, + aliases: vec![], + }, + Field { + name: "nanos", + r#type: Schema::Union(vec![null.clone(), int.clone()]), + default: None, + doc: None, + aliases: vec![], + }, ], attributes: Attributes::default(), })); @@ -6108,10 +6345,16 @@ mod tests { // Reader Root: only `ts` (extra is writer-only → will be Skipped) let reader_schema = Schema::Complex(ComplexType::Record(Record { name: "Root", - namespace: None, doc: None, aliases: vec![], - fields: vec![ - Field { name: "ts", r#type: timestamp_reader, default: None, doc: None, aliases: vec![] }, - ], + namespace: None, + doc: None, + aliases: vec![], + fields: vec![Field { + name: "ts", + r#type: timestamp_reader, + default: None, + doc: None, + aliases: vec![], + }], attributes: Attributes::default(), })); @@ -6120,8 +6363,7 @@ mod tests { .build() .expect("schema resolution must succeed"); - let mut decoder = Decoder::try_new(field.data_type()) - .expect("decoder must build"); + let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); // --- Test case 1: extra = null --- // Writer encodes: @@ -6134,8 +6376,8 @@ mod tests { // struct fields — corrupting the cursor. let mut data_null = Vec::new(); data_null.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data_null.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data_null.push(0x00); // extra: union tag 0 = null + data_null.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data_null.push(0x00); // extra: union tag 0 = null let mut cur = AvroCursor::new(&data_null); decoder @@ -6156,10 +6398,10 @@ mod tests { // The Skipper for `extra` must read the union tag byte, then skip the Timestamp payload. let mut data_nonnull = Vec::new(); data_nonnull.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data_nonnull.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data_nonnull.extend_from_slice(&encode_avro_int(1)); // extra: union tag zigzag(1)=0x02 + data_nonnull.extend_from_slice(&encode_avro_int(5)); // ts.nanos + data_nonnull.extend_from_slice(&encode_avro_int(1)); // extra: union tag zigzag(1)=0x02 data_nonnull.extend_from_slice(&encode_avro_long(200)); // extra.seconds - data_nonnull.extend_from_slice(&encode_avro_int(3)); // extra.nanos + data_nonnull.extend_from_slice(&encode_avro_int(3)); // extra.nanos let mut cur2 = AvroCursor::new(&data_nonnull); decoder From 7ca866a3bffb18c6a85d7a6976721c250842e419 Mon Sep 17 00:00:00 2001 From: Ariel Miculas Date: Thu, 2 Apr 2026 13:28:04 +0300 Subject: [PATCH 3/4] chore: remove the previous skipper tests --- arrow-avro/src/reader/record.rs | 789 +------------------------------- 1 file changed, 1 insertion(+), 788 deletions(-) diff --git a/arrow-avro/src/reader/record.rs b/arrow-avro/src/reader/record.rs index 9db0be385a7a..b71a6bdc7cdb 100644 --- a/arrow-avro/src/reader/record.rs +++ b/arrow-avro/src/reader/record.rs @@ -2730,10 +2730,7 @@ impl Skipper { mod tests { use super::*; use crate::codec::AvroFieldBuilder; - use crate::schema::Array as AvroArray; - use crate::schema::{ - Attributes, ComplexType, Enum as AvroEnum, Field, PrimitiveType, Record, Schema, TypeName, - }; + use crate::schema::{Attributes, ComplexType, Field, PrimitiveType, Record, Schema, TypeName}; use arrow_array::cast::AsArray; use indexmap::IndexMap; use std::collections::HashMap; @@ -5629,788 +5626,4 @@ mod tests { other => panic!("expected Timestamp(Nanosecond, None), got {other:?}"), } } - - /// When a Skip field references a named type by name, the Skipper must use the - /// writer's wire format for that type. Here the reader wraps the Timestamp's scalar - /// fields in nullable unions, but the writer wrote them as plain values; the Skipper - /// must not add a union-tag read for each field. - #[test] - fn test_skip_named_type_ref_uses_writer_schema_not_resolved() { - let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); - - // Writer: Timestamp{seconds:long, nanos:int} (plain, no union wrappers) - let timestamp_writer = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: long.clone(), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Writer: Event{time:"Timestamp", code:int} — `time` is a named-type reference - let event_writer = Schema::Complex(ComplexType::Record(Record { - name: "Event", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "time", - r#type: Schema::TypeName(TypeName::Ref("Timestamp")), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "code", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Writer root: { ts: Timestamp, events: array } - let writer_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "ts", - r#type: timestamp_writer, - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "events", - r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader Timestamp: inner scalars wrapped in ["null", T] - let timestamp_reader = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: Schema::Union(vec![null.clone(), long.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: Schema::Union(vec![null.clone(), int.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader root: only `ts` (nullable wrapper), `events` absent → Skip - let reader_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "ts", - r#type: Schema::Union(vec![null, timestamp_reader]), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let field = AvroFieldBuilder::new(&writer_schema) - .with_reader_schema(&reader_schema) - .build() - .expect("schema resolution must succeed"); - - let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); - - // Encode one writer record: - // ts = Timestamp{seconds=100, nanos=5} — plain longs/ints, no union tags - // events = [Event{time=Timestamp{seconds=200, nanos=1}, code=42}] - let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds (plain) - data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos (plain) - data.extend_from_slice(&encode_avro_int(42)); // event.code - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array - - let mut cur = AvroCursor::new(&data); - decoder - .decode(&mut cur) - .expect("decode must not corrupt cursor due to spurious Nullable wrappers"); - assert_eq!( - cur.position(), - data.len(), - "cursor must consume exactly the writer-encoded bytes" - ); - } - - /// When a Skip field references a named type that has more fields in the writer - /// schema than in the reader schema, the Skipper must consume all writer fields - /// including writer-only ones. Here the writer's Timestamp has a third field - /// `tz_offset` that the reader omits; the Skipper must still consume those bytes. - #[test] - fn test_skip_named_type_ref_with_writer_extra_fields_consumes_all_bytes() { - let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); - - // Writer Timestamp has an extra field `tz_offset` that the reader omits - let timestamp_writer = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: long.clone(), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "tz_offset", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - let event_writer = Schema::Complex(ComplexType::Record(Record { - name: "Event", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "time", - r#type: Schema::TypeName(TypeName::Ref("Timestamp")), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let writer_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "ts", - r#type: timestamp_writer, - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "events", - r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader Timestamp: only seconds and nanos (no tz_offset) - let timestamp_reader = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: Schema::Union(vec![null.clone(), long.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: Schema::Union(vec![null.clone(), int.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - let reader_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "ts", - r#type: Schema::Union(vec![null, timestamp_reader]), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let field = AvroFieldBuilder::new(&writer_schema) - .with_reader_schema(&reader_schema) - .build() - .expect("schema resolution must succeed"); - - let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); - - // Encode one writer record: - // ts = Timestamp{seconds=100, nanos=5, tz_offset=3600} - // events = [Event{time=Timestamp{seconds=200, nanos=1, tz_offset=7200}}] - // The Skipper for `events` must consume ALL three fields per Timestamp, - // including tz_offset, even though the resolved Timestamp only has two reader fields. - let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data.extend_from_slice(&encode_avro_int(3600)); // ts.tz_offset (writer-only, decoded via writer_fields) - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_long(200)); // event.time.seconds - data.extend_from_slice(&encode_avro_int(1)); // event.time.nanos - data.extend_from_slice(&encode_avro_int(7200)); // event.time.tz_offset (must be skipped!) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array - - let mut cur = AvroCursor::new(&data); - decoder.decode(&mut cur).expect("decode must not fail"); - assert_eq!( - cur.position(), - data.len(), - "cursor must consume all writer-encoded bytes including writer-only tz_offset in skipped events" - ); - } - - /// When a Skip field references a named type whose inner struct field was written - /// as a plain record but the reader wraps it in a nullable union, the Skipper must - /// use the writer's plain encoding and not read a union tag for that inner field. - #[test] - fn test_skip_nested_struct_record_resolution_no_spurious_nullable_wrapper() { - let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); - - // Writer: InnerRecord { v: int } - let inner_writer = Schema::Complex(ComplexType::Record(Record { - name: "InnerRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "v", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Writer: OuterRecord { inner: InnerRecord } - let outer_writer = Schema::Complex(ComplexType::Record(Record { - name: "OuterRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "inner", - r#type: inner_writer, - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Writer: Event { data: "OuterRecord" } (named type ref — will look up registered type) - let event_writer = Schema::Complex(ComplexType::Record(Record { - name: "Event", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "data", - r#type: Schema::TypeName(TypeName::Ref("OuterRecord")), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Writer Root: { outer: OuterRecord, events: array } - // `events` is omitted from the reader schema → will be skipped. - let writer_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "outer", - r#type: outer_writer, - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "events", - r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader: InnerRecord { v: int } - let inner_reader = Schema::Complex(ComplexType::Record(Record { - name: "InnerRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "v", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Reader: OuterRecord { inner: ["null", InnerRecord] } — inner wrapped in nullable - let outer_reader = Schema::Complex(ComplexType::Record(Record { - name: "OuterRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "inner", - r#type: Schema::Union(vec![null.clone(), inner_reader]), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Reader Root: { outer: ["null", OuterRecord] } — events omitted (Skip) - let reader_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "outer", - r#type: Schema::Union(vec![null, outer_reader]), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let field = AvroFieldBuilder::new(&writer_schema) - .with_reader_schema(&reader_schema) - .build() - .expect("schema resolution must succeed"); - - let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); - - // Encode one writer record (all plain, no union tags anywhere): - // outer.inner.v = 42 - // events = [Event { data = OuterRecord { inner = InnerRecord { v = 7 } } }] - // The Skipper for `events` must skip InnerRecord as a plain struct (no union tag), - // even though the registered OuterRecord has inner with nullability+Record resolution. - let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_int(42)); // outer.inner.v - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_int(7)); // event.data.inner.v (skipped, no union tag) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array - - let mut cur = AvroCursor::new(&data); - decoder.decode(&mut cur).expect("decode must not fail"); - assert_eq!( - cur.position(), - data.len(), - "cursor must consume all bytes; spurious Nullable wrapper would consume an extra byte" - ); - } - - /// When a Skip field references a named type whose inner enum field was written as a - /// plain index but the reader wraps it in a nullable union (with symbol reordering), - /// the Skipper must skip the enum as a plain VLQ int without reading a union tag. - #[test] - fn test_skip_enum_mapping_nullable_no_spurious_nullable_wrapper() { - let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - - // Writer: SomeRecord { status: StatusEnum{OK, ERROR} } - let status_writer = Schema::Complex(ComplexType::Enum(AvroEnum { - name: "StatusEnum", - namespace: None, - doc: None, - aliases: vec![], - symbols: vec!["OK", "ERROR"], - default: None, - attributes: Attributes::default(), - })); - let some_record_writer = Schema::Complex(ComplexType::Record(Record { - name: "SomeRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "status", - r#type: status_writer, - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Writer: Event { data: "SomeRecord" } (named type ref — looks up registered type) - let event_writer = Schema::Complex(ComplexType::Record(Record { - name: "Event", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "data", - r#type: Schema::TypeName(TypeName::Ref("SomeRecord")), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Writer Root: { outer: SomeRecord, events: array } - // `events` is omitted from the reader schema → will be skipped. - let writer_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "outer", - r#type: some_record_writer, - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "events", - r#type: Schema::Complex(ComplexType::Array(AvroArray { - items: Box::new(event_writer), - attributes: Attributes::default(), - })), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader: SomeRecord { status: ["null", StatusEnum{ERROR, OK}] } (reordered + nullable) - let status_reader = Schema::Complex(ComplexType::Enum(AvroEnum { - name: "StatusEnum", - namespace: None, - doc: None, - aliases: vec![], - symbols: vec!["ERROR", "OK"], - default: None, - attributes: Attributes::default(), - })); - let some_record_reader = Schema::Complex(ComplexType::Record(Record { - name: "SomeRecord", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "status", - r#type: Schema::Union(vec![null, status_reader]), - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - // Reader Root: { outer: SomeRecord } — events omitted (Skip) - let reader_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "outer", - r#type: some_record_reader, - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let field = AvroFieldBuilder::new(&writer_schema) - .with_reader_schema(&reader_schema) - .build() - .expect("schema resolution must succeed"); - - let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); - - // Encode one writer record (all plain, no union tags anywhere): - // outer.status = OK (writer index 0) - // events = [Event { data = SomeRecord { status = ERROR (writer index 1) } }] - // The Skipper for `events` must skip the enum as a plain VLQ int (no union tag). - // Using ERROR (index 1, encoded as 0x02) is deliberate: a spurious Nullable wrapper - // would read 0x02 as union tag 1 ("non-null") and then consume an additional byte - // for the inner value, reading into the end-of-array marker and leaving the cursor - // at the wrong position (or causing an EOF error). - let mut data = Vec::new(); - data.extend_from_slice(&encode_avro_int(0)); // outer.status: OK = writer index 0 - data.extend_from_slice(&encode_avro_long(1)); // events: block of 1 - data.extend_from_slice(&encode_avro_int(1)); // event.data.status: ERROR = writer index 1 (skipped) - data.extend_from_slice(&encode_avro_long(0)); // events: end-of-array - - let mut cur = AvroCursor::new(&data); - decoder.decode(&mut cur).expect("decode must not fail"); - assert_eq!( - cur.position(), - data.len(), - "cursor must consume all bytes; spurious Nullable wrapper would consume an extra byte" - ); - } - - /// When a writer-only Skip field is a genuine nullable union `["null", "TypeRef"]`, - /// the Skipper must preserve the `Nullable` wrapper and read the union tag byte. - /// This is distinct from resolution-induced nullability: here the writer itself - /// wrote a union tag, so the tag byte must be consumed. - #[test] - fn test_skip_nullable_named_type_ref_union_not_stripped_of_nullable_wrapper() { - let null = Schema::TypeName(TypeName::Primitive(PrimitiveType::Null)); - let long = Schema::TypeName(TypeName::Primitive(PrimitiveType::Long)); - let int = Schema::TypeName(TypeName::Primitive(PrimitiveType::Int)); - - // Writer Timestamp: plain {seconds: long, nanos: int} - let timestamp_writer = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: long.clone(), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: int.clone(), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Writer Root: - // ts: Timestamp — matched by reader, triggers resolver registration - // extra: ["null", "Timestamp"] — writer-only genuine nullable union (will be Skipped) - let writer_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "ts", - r#type: timestamp_writer, - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "extra", - r#type: Schema::Union(vec![ - null.clone(), - Schema::TypeName(TypeName::Ref("Timestamp")), - ]), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader Timestamp: fields are nullable (writer plain → reader nullable via resolution) - let timestamp_reader = Schema::Complex(ComplexType::Record(Record { - name: "Timestamp", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![ - Field { - name: "seconds", - r#type: Schema::Union(vec![null.clone(), long.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - Field { - name: "nanos", - r#type: Schema::Union(vec![null.clone(), int.clone()]), - default: None, - doc: None, - aliases: vec![], - }, - ], - attributes: Attributes::default(), - })); - - // Reader Root: only `ts` (extra is writer-only → will be Skipped) - let reader_schema = Schema::Complex(ComplexType::Record(Record { - name: "Root", - namespace: None, - doc: None, - aliases: vec![], - fields: vec![Field { - name: "ts", - r#type: timestamp_reader, - default: None, - doc: None, - aliases: vec![], - }], - attributes: Attributes::default(), - })); - - let field = AvroFieldBuilder::new(&writer_schema) - .with_reader_schema(&reader_schema) - .build() - .expect("schema resolution must succeed"); - - let mut decoder = Decoder::try_new(field.data_type()).expect("decoder must build"); - - // --- Test case 1: extra = null --- - // Writer encodes: - // ts.seconds = 100 (plain long, no union tag) - // ts.nanos = 5 (plain int, no union tag) - // extra = null (union tag 0 = 0x00) - // The Skipper for `extra` must read the union tag byte (0x00 → null branch, no payload). - // Before the fix, the Nullable wrapper was stripped and the Skipper treated `extra` as a - // plain Timestamp struct, consuming ts.seconds bytes of the next record as if they were - // struct fields — corrupting the cursor. - let mut data_null = Vec::new(); - data_null.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data_null.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data_null.push(0x00); // extra: union tag 0 = null - - let mut cur = AvroCursor::new(&data_null); - decoder - .decode(&mut cur) - .expect("decode must not fail (null branch)"); - assert_eq!( - cur.position(), - data_null.len(), - "cursor must consume exactly the writer bytes including the null union tag" - ); - - // --- Test case 2: extra = non-null Timestamp --- - // Writer encodes: - // ts.seconds = 100 - // ts.nanos = 5 - // extra = Timestamp{seconds=200, nanos=3} - // → union tag 2 (0x04, zigzag for index 1) + seconds=200 + nanos=3 - // The Skipper for `extra` must read the union tag byte, then skip the Timestamp payload. - let mut data_nonnull = Vec::new(); - data_nonnull.extend_from_slice(&encode_avro_long(100)); // ts.seconds - data_nonnull.extend_from_slice(&encode_avro_int(5)); // ts.nanos - data_nonnull.extend_from_slice(&encode_avro_int(1)); // extra: union tag zigzag(1)=0x02 - data_nonnull.extend_from_slice(&encode_avro_long(200)); // extra.seconds - data_nonnull.extend_from_slice(&encode_avro_int(3)); // extra.nanos - - let mut cur2 = AvroCursor::new(&data_nonnull); - decoder - .decode(&mut cur2) - .expect("decode must not fail (non-null branch)"); - assert_eq!( - cur2.position(), - data_nonnull.len(), - "cursor must consume all writer bytes including union tag and Timestamp payload" - ); - } } From 2401361816f1d9a4d4510a43b6671fb1b274d8e9 Mon Sep 17 00:00:00 2001 From: Ariel Miculas Date: Thu, 2 Apr 2026 13:29:20 +0300 Subject: [PATCH 4/4] feat: add high-level skipper tests --- arrow-avro/src/reader/mod.rs | 215 +++++++++++++++++++++++++++++++++++ 1 file changed, 215 insertions(+) diff --git a/arrow-avro/src/reader/mod.rs b/arrow-avro/src/reader/mod.rs index 070204f2bcfb..6dbf5b1553c2 100644 --- a/arrow-avro/src/reader/mod.rs +++ b/arrow-avro/src/reader/mod.rs @@ -9590,4 +9590,219 @@ mod test { "entire RecordBatch mismatch (schema, all columns, all rows)" ); } + + // Build Avro OCF bytes whose schema contains a TypeName::Ref + // + // Schema written to the OCF header verbatim: + // ```text + // Root { + // ts: Timestamp { seconds: long, nanos: int }, + // extra: Event { time: "Timestamp" } <- TypeName::Ref + // } + // ``` + fn make_type_ref_ocf() -> Vec { + use apache_avro::{Schema as ApacheSchema, Writer as ApacheWriter, types::Value}; + let schema_json = r#"{ + "type": "record", "name": "Root", + "fields": [ + {"name": "ts", "type": {"type": "record", "name": "Timestamp", "fields": [ + {"name": "seconds", "type": "long"}, + {"name": "nanos", "type": "int"} + ]}}, + {"name": "extra", "type": {"type": "record", "name": "Event", "fields": [ + {"name": "time", "type": "Timestamp"} + ]}} + ] + }"#; + let schema = ApacheSchema::parse_str(schema_json).expect("valid schema"); + let mut out = Vec::new(); + { + let mut writer = ApacheWriter::new(&schema, &mut out); + let ts_val = |s: i64, n: i32| { + Value::Record(vec![ + ("seconds".into(), Value::Long(s)), + ("nanos".into(), Value::Int(n)), + ]) + }; + // Two rows: ts={1000,100}/extra.time={-1,-1} and ts={2000,200}/extra.time={-2,-2}. + for (ts_s, ts_n, ex_s, ex_n) in [(1000i64, 100i32, -1i64, -1i32), (2000, 200, -2, -2)] { + let row = Value::Record(vec![ + ("ts".into(), ts_val(ts_s, ts_n)), + ( + "extra".into(), + Value::Record(vec![("time".into(), ts_val(ex_s, ex_n))]), + ), + ]); + writer.append_value_ref(&row).expect("append row"); + } + writer.flush().expect("flush"); + } + out + } + + // writer-plain / reader-nullable mismatch. + // + // The writer schema uses a TypeName::Ref ("Timestamp" referenced in `extra.time`). + // The reader wraps `ts` in `["null", T]` unions and omits `extra`. + // The Skipper for `extra.time` resolves "Timestamp" via the resolver and must use + // the writer's plain field types (long, int) — not the nullable reader types - when + // consuming bytes. Without the fix, it skips union-encoded fields from plain data, + // reads the wrong number of bytes, and corrupts row 2's `ts.seconds`. + #[test] + fn test_nullable_reader_schema_vs_plain_writer_nested_struct() { + let bytes = make_type_ref_ocf(); + let reader_schema = AvroSchema::new( + r#"{"type":"record","name":"Root","fields":[ + {"name":"ts","type":["null",{"type":"record","name":"Timestamp","fields":[ + {"name":"seconds","type":["null","long"]}, + {"name":"nanos", "type":["null","int"]} + ]}]} + ]}"# + .to_string(), + ); + let mut reader = ReaderBuilder::new() + .with_reader_schema(reader_schema) + .build(Cursor::new(bytes)) + .expect("reader should build"); + let batch = reader + .next() + .expect("should have a batch") + .expect("reading should succeed"); + assert_eq!(batch.num_rows(), 2); + let ts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let seconds = ts + .column_by_name("seconds") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(seconds.value(0), 1000); + assert_eq!(seconds.value(1), 2000); + } + + // Skipper must consume all writer fields, including writer-only ones. + // + // The writer schema uses a TypeName::Ref ("Timestamp" referenced in `extra.time`). + // The reader requests only `ts.seconds` (no `nanos`, no `extra`). + // The Skipper for `extra.time` resolves "Timestamp" and must skip both `seconds` + // and `nanos` bytes. Without the fix it skips only `seconds`, leaving the `nanos` + // bytes in the buffer and corrupting row 2's `ts.seconds` read. + #[test] + fn test_skipper_consumes_writer_only_struct_fields() { + let bytes = make_type_ref_ocf(); + let reader_schema = AvroSchema::new( + r#"{"type":"record","name":"Root","fields":[ + {"name":"ts","type":{"type":"record","name":"Timestamp","fields":[ + {"name":"seconds","type":"long"} + ]}} + ]}"# + .to_string(), + ); + let mut reader = ReaderBuilder::new() + .with_reader_schema(reader_schema) + .build(Cursor::new(bytes)) + .expect("reader should build"); + let batch = reader + .next() + .expect("should have a batch") + .expect("Skipper must consume both seconds and nanos for extra.time"); + assert_eq!(batch.num_rows(), 2); + let ts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let seconds = ts + .column_by_name("seconds") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(seconds.value(0), 1000); + assert_eq!(seconds.value(1), 2000); + } + + // The Skipper for a skipped array field must consume all bytes of each element, + // including every field of a nested struct resolved via a TypeName::Ref. + // + // Writer: `Root { ts: Timestamp{seconds,nanos}, events: array }` + // Reader: only `ts` with nullable wrappers; `events` is absent (forces a Skip). + // The Skipper for `events` resolves each element's `time` field as "Timestamp" + // and must use the writer's plain {seconds,nanos} definition — not the + // nullable-wrapped reader type — when consuming bytes. + #[test] + fn test_skip_array_of_structs_uses_writer_schema_not_resolved() { + use apache_avro::{Schema as ApacheSchema, Writer as ApacheWriter, types::Value}; + let schema_json = r#"{ + "type": "record", "name": "Root", + "fields": [ + {"name": "ts", "type": {"type": "record", "name": "Timestamp", "fields": [ + {"name": "seconds", "type": "long"}, + {"name": "nanos", "type": "int"} + ]}}, + {"name": "events", "type": {"type": "array", "items": { + "type": "record", "name": "Event", "fields": [ + {"name": "time", "type": "Timestamp"} + ] + }}} + ] + }"#; + let schema = ApacheSchema::parse_str(schema_json).expect("valid schema"); + let mut bytes = Vec::new(); + { + let mut writer = ApacheWriter::new(&schema, &mut bytes); + // One row: ts={100, 5}, events=[{time={200, 1}}] + let ts_val = |s: i64, n: i32| { + Value::Record(vec![ + ("seconds".into(), Value::Long(s)), + ("nanos".into(), Value::Int(n)), + ]) + }; + let row = Value::Record(vec![ + ("ts".into(), ts_val(100, 5)), + ( + "events".into(), + Value::Array(vec![Value::Record(vec![("time".into(), ts_val(200, 1))])]), + ), + ]); + writer.append_value_ref(&row).expect("append row"); + writer.flush().expect("flush"); + } + + // Reader omits `events` (forces Skip) and wraps `ts` fields in nullable unions. + let reader_schema = AvroSchema::new( + r#"{"type":"record","name":"Root","fields":[ + {"name":"ts","type":["null",{"type":"record","name":"Timestamp","fields":[ + {"name":"seconds","type":["null","long"]}, + {"name":"nanos", "type":["null","int"]} + ]}]} + ]}"# + .to_string(), + ); + let mut reader = ReaderBuilder::new() + .with_reader_schema(reader_schema) + .build(Cursor::new(bytes)) + .expect("reader should build"); + let batch = reader + .next() + .expect("should have a batch") + .expect("Skipper must consume all events bytes using writer field types"); + assert_eq!(batch.num_rows(), 1); + let ts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let seconds = ts + .column_by_name("seconds") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(seconds.value(0), 100); + } }