Skip to main content

arrow_json/reader/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! JSON reader
19//!
20//! This JSON reader allows JSON records to be read into the Arrow memory
21//! model. Records are loaded in batches and are then converted from the record-oriented
22//! representation to the columnar arrow data model.
23//!
24//! The reader ignores whitespace between JSON values, including `\n` and `\r`, allowing
25//! parsing of sequences of one or more arbitrarily formatted JSON values, including
26//! but not limited to newline-delimited JSON.
27//!
28//! # Basic Usage
29//!
30//! [`Reader`] can be used directly with synchronous data sources, such as [`std::fs::File`]
31//!
32//! ```
33//! # use arrow_schema::*;
34//! # use std::fs::File;
35//! # use std::io::BufReader;
36//! # use std::sync::Arc;
37//!
38//! let schema = Arc::new(Schema::new(vec![
39//!     Field::new("a", DataType::Float64, false),
40//!     Field::new("b", DataType::Float64, false),
41//!     Field::new("c", DataType::Boolean, true),
42//! ]));
43//!
44//! let file = File::open("test/data/basic.json").unwrap();
45//!
46//! let mut json = arrow_json::ReaderBuilder::new(schema).build(BufReader::new(file)).unwrap();
47//! let batch = json.next().unwrap().unwrap();
48//! ```
49//!
50//! # Async Usage
51//!
52//! The lower-level [`Decoder`] can be integrated with various forms of async data streams,
53//! and is designed to be agnostic to the various different kinds of async IO primitives found
54//! within the Rust ecosystem.
55//!
56//! For example, see below for how it can be used with an arbitrary `Stream` of `Bytes`
57//!
58//! ```
59//! # use std::task::{Poll, ready};
60//! # use bytes::{Buf, Bytes};
61//! # use arrow_schema::ArrowError;
62//! # use futures::stream::{Stream, StreamExt};
63//! # use arrow_array::RecordBatch;
64//! # use arrow_json::reader::Decoder;
65//! #
66//! fn decode_stream<S: Stream<Item = Bytes> + Unpin>(
67//!     mut decoder: Decoder,
68//!     mut input: S,
69//! ) -> impl Stream<Item = Result<RecordBatch, ArrowError>> {
70//!     let mut buffered = Bytes::new();
71//!     futures::stream::poll_fn(move |cx| {
72//!         loop {
73//!             if buffered.is_empty() {
74//!                 buffered = match ready!(input.poll_next_unpin(cx)) {
75//!                     Some(b) => b,
76//!                     None => break,
77//!                 };
78//!             }
79//!             let decoded = match decoder.decode(buffered.as_ref()) {
80//!                 Ok(decoded) => decoded,
81//!                 Err(e) => return Poll::Ready(Some(Err(e))),
82//!             };
83//!             let read = buffered.len();
84//!             buffered.advance(decoded);
85//!             if decoded != read {
86//!                 break
87//!             }
88//!         }
89//!
90//!         Poll::Ready(decoder.flush().transpose())
91//!     })
92//! }
93//!
94//! ```
95//!
96//! In a similar vein, it can also be used with tokio-based IO primitives
97//!
98//! ```
99//! # use std::sync::Arc;
100//! # use arrow_schema::{DataType, Field, Schema};
101//! # use std::pin::Pin;
102//! # use std::task::{Poll, ready};
103//! # use futures::{Stream, TryStreamExt};
104//! # use tokio::io::AsyncBufRead;
105//! # use arrow_array::RecordBatch;
106//! # use arrow_json::reader::Decoder;
107//! # use arrow_schema::ArrowError;
108//! fn decode_stream<R: AsyncBufRead + Unpin>(
109//!     mut decoder: Decoder,
110//!     mut reader: R,
111//! ) -> impl Stream<Item = Result<RecordBatch, ArrowError>> {
112//!     futures::stream::poll_fn(move |cx| {
113//!         loop {
114//!             let b = match ready!(Pin::new(&mut reader).poll_fill_buf(cx)) {
115//!                 Ok(b) if b.is_empty() => break,
116//!                 Ok(b) => b,
117//!                 Err(e) => return Poll::Ready(Some(Err(e.into()))),
118//!             };
119//!             let read = b.len();
120//!             let decoded = match decoder.decode(b) {
121//!                 Ok(decoded) => decoded,
122//!                 Err(e) => return Poll::Ready(Some(Err(e))),
123//!             };
124//!             Pin::new(&mut reader).consume(decoded);
125//!             if decoded != read {
126//!                 break;
127//!             }
128//!         }
129//!
130//!         Poll::Ready(decoder.flush().transpose())
131//!     })
132//! }
133//! ```
134//!
135
136use std::borrow::Cow;
137use std::io::BufRead;
138use std::sync::Arc;
139
140use arrow_array::cast::AsArray;
141use arrow_array::timezone::Tz;
142use arrow_array::types::*;
143use arrow_array::{ArrayRef, RecordBatch, RecordBatchReader, downcast_integer};
144use arrow_schema::{ArrowError, DataType, FieldRef, Schema, SchemaRef, TimeUnit};
145use chrono::Utc;
146use serde_core::Serialize;
147
148use crate::StructMode;
149use crate::reader::binary_array::{
150    BinaryArrayDecoder, BinaryViewDecoder, FixedSizeBinaryArrayDecoder,
151};
152use crate::reader::boolean_array::BooleanArrayDecoder;
153use crate::reader::decimal_array::DecimalArrayDecoder;
154use crate::reader::list_array::{
155    FixedSizeListArrayDecoder, ListArrayDecoder, ListViewArrayDecoder,
156};
157use crate::reader::map_array::MapArrayDecoder;
158use crate::reader::null_array::NullArrayDecoder;
159use crate::reader::primitive_array::PrimitiveArrayDecoder;
160use crate::reader::run_end_array::RunEndEncodedArrayDecoder;
161use crate::reader::string_array::StringArrayDecoder;
162use crate::reader::string_view_array::StringViewArrayDecoder;
163use crate::reader::struct_array::StructArrayDecoder;
164use crate::reader::tape::{Tape, TapeDecoder};
165use crate::reader::timestamp_array::TimestampArrayDecoder;
166
167pub use schema::*;
168pub use value_iter::ValueIter;
169
170mod binary_array;
171mod boolean_array;
172mod decimal_array;
173mod list_array;
174mod map_array;
175mod null_array;
176mod primitive_array;
177mod run_end_array;
178mod schema;
179mod serializer;
180mod string_array;
181mod string_view_array;
182mod struct_array;
183mod tape;
184mod timestamp_array;
185mod value_iter;
186
187/// A builder for [`Reader`] and [`Decoder`]
188pub struct ReaderBuilder {
189    batch_size: usize,
190    coerce_primitive: bool,
191    strict_mode: bool,
192    ignore_type_conflicts: bool,
193    is_field: bool,
194    struct_mode: StructMode,
195
196    schema: SchemaRef,
197}
198
199impl ReaderBuilder {
200    /// Create a new [`ReaderBuilder`] with the provided [`SchemaRef`]
201    ///
202    /// This could be obtained using [`infer_json_schema`] if not known
203    ///
204    /// Any columns not present in `schema` will be ignored, unless `strict_mode` is set to true.
205    /// In this case, an error is returned when a column is missing from `schema`.
206    ///
207    /// [`infer_json_schema`]: crate::reader::infer_json_schema
208    pub fn new(schema: SchemaRef) -> Self {
209        Self {
210            batch_size: 1024,
211            coerce_primitive: false,
212            strict_mode: false,
213            ignore_type_conflicts: false,
214            is_field: false,
215            struct_mode: Default::default(),
216            schema,
217        }
218    }
219
220    /// Create a new [`ReaderBuilder`] that will parse JSON values of `field.data_type()`
221    ///
222    /// Unlike [`ReaderBuilder::new`] this does not require the root of the JSON data
223    /// to be an object, i.e. `{..}`, allowing for parsing of any valid JSON value(s)
224    ///
225    /// ```
226    /// # use std::sync::Arc;
227    /// # use arrow_array::cast::AsArray;
228    /// # use arrow_array::types::Int32Type;
229    /// # use arrow_json::ReaderBuilder;
230    /// # use arrow_schema::{DataType, Field};
231    /// // Root of JSON schema is a numeric type
232    /// let data = "1\n2\n3\n";
233    /// let field = Arc::new(Field::new("int", DataType::Int32, true));
234    /// let mut reader = ReaderBuilder::new_with_field(field.clone()).build(data.as_bytes()).unwrap();
235    /// let b = reader.next().unwrap().unwrap();
236    /// let values = b.column(0).as_primitive::<Int32Type>().values();
237    /// assert_eq!(values, &[1, 2, 3]);
238    ///
239    /// // Root of JSON schema is a list type
240    /// let data = "[1, 2, 3, 4, 5, 6, 7]\n[1, 2, 3]";
241    /// let field = Field::new_list("int", field.clone(), true);
242    /// let mut reader = ReaderBuilder::new_with_field(field).build(data.as_bytes()).unwrap();
243    /// let b = reader.next().unwrap().unwrap();
244    /// let list = b.column(0).as_list::<i32>();
245    ///
246    /// assert_eq!(list.offsets().as_ref(), &[0, 7, 10]);
247    /// let list_values = list.values().as_primitive::<Int32Type>();
248    /// assert_eq!(list_values.values(), &[1, 2, 3, 4, 5, 6, 7, 1, 2, 3]);
249    /// ```
250    pub fn new_with_field(field: impl Into<FieldRef>) -> Self {
251        Self {
252            batch_size: 1024,
253            coerce_primitive: false,
254            strict_mode: false,
255            ignore_type_conflicts: false,
256            is_field: true,
257            struct_mode: Default::default(),
258            schema: Arc::new(Schema::new([field.into()])),
259        }
260    }
261
262    /// Sets the batch size in rows to read
263    pub fn with_batch_size(self, batch_size: usize) -> Self {
264        Self { batch_size, ..self }
265    }
266
267    /// Sets if the decoder should coerce primitive values (bool and number) into string
268    /// when the Schema's column is Utf8 or LargeUtf8.
269    pub fn with_coerce_primitive(self, coerce_primitive: bool) -> Self {
270        Self {
271            coerce_primitive,
272            ..self
273        }
274    }
275
276    /// Sets if the decoder should return an error if it encounters a column not
277    /// present in `schema`. If `struct_mode` is `ListOnly` the value of
278    /// `strict_mode` is effectively `true`. It is required for all fields of
279    /// the struct to be in the list: without field names, there is no way to
280    /// determine which field is missing.
281    pub fn with_strict_mode(self, strict_mode: bool) -> Self {
282        Self {
283            strict_mode,
284            ..self
285        }
286    }
287
288    /// Set the [`StructMode`] for the reader, which determines whether structs
289    /// can be decoded from JSON as objects or lists. For more details refer to
290    /// the enum documentation. Default is to use `ObjectOnly`.
291    pub fn with_struct_mode(self, struct_mode: StructMode) -> Self {
292        Self {
293            struct_mode,
294            ..self
295        }
296    }
297
298    /// Sets whether the decoder should produce NULL instead of returning an error if it encounters
299    /// value that can not be parsed into the specified column type.
300    ///
301    /// For example, if the type is declared to be a nullable array of `DataType::Int32` but the
302    /// reader encounters a string value `"foo"` and the value `ignore_type_conflicts` is:
303    ///
304    /// * `false` (the default): The reader will return an error.
305    ///
306    /// * `true`: The reader will fill in NULL value for that array element.
307    ///
308    /// NOTE: An inferred NULL due to a type conflict will still produce parsing errors for
309    /// non-nullable fields, the same as any other NULL or missing value.
310    pub fn with_ignore_type_conflicts(self, ignore_type_conflicts: bool) -> Self {
311        Self {
312            ignore_type_conflicts,
313            ..self
314        }
315    }
316
317    /// Create a [`Reader`] with the provided [`BufRead`]
318    pub fn build<R: BufRead>(self, reader: R) -> Result<Reader<R>, ArrowError> {
319        Ok(Reader {
320            reader,
321            decoder: self.build_decoder()?,
322        })
323    }
324
325    /// Create a [`Decoder`]
326    pub fn build_decoder(self) -> Result<Decoder, ArrowError> {
327        let (data_type, nullable) = if self.is_field {
328            let field = &self.schema.fields[0];
329            let data_type = Cow::Borrowed(field.data_type());
330            (data_type, field.is_nullable())
331        } else {
332            let data_type = Cow::Owned(DataType::Struct(self.schema.fields.clone()));
333            (data_type, false)
334        };
335
336        let ctx = DecoderContext {
337            coerce_primitive: self.coerce_primitive,
338            strict_mode: self.strict_mode,
339            struct_mode: self.struct_mode,
340            ignore_type_conflicts: self.ignore_type_conflicts,
341        };
342        let decoder = ctx.make_decoder(data_type.as_ref(), nullable)?;
343
344        let num_fields = self.schema.flattened_fields().len();
345
346        Ok(Decoder {
347            decoder,
348            is_field: self.is_field,
349            tape_decoder: TapeDecoder::new(self.batch_size, num_fields),
350            batch_size: self.batch_size,
351            schema: self.schema,
352        })
353    }
354}
355
356/// Reads JSON data with a known schema directly into arrow [`RecordBatch`]
357///
358/// Lines consisting solely of ASCII whitespace are ignored
359pub struct Reader<R> {
360    reader: R,
361    decoder: Decoder,
362}
363
364impl<R> std::fmt::Debug for Reader<R> {
365    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
366        f.debug_struct("Reader")
367            .field("decoder", &self.decoder)
368            .finish()
369    }
370}
371
372impl<R: BufRead> Reader<R> {
373    /// Reads the next [`RecordBatch`] returning `Ok(None)` if EOF
374    fn read(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
375        loop {
376            let buf = self.reader.fill_buf()?;
377            if buf.is_empty() {
378                break;
379            }
380            let read = buf.len();
381
382            let decoded = self.decoder.decode(buf)?;
383            self.reader.consume(decoded);
384            if decoded != read {
385                break;
386            }
387        }
388        self.decoder.flush()
389    }
390}
391
392impl<R: BufRead> Iterator for Reader<R> {
393    type Item = Result<RecordBatch, ArrowError>;
394
395    fn next(&mut self) -> Option<Self::Item> {
396        self.read().transpose()
397    }
398}
399
400impl<R: BufRead> RecordBatchReader for Reader<R> {
401    fn schema(&self) -> SchemaRef {
402        self.decoder.schema.clone()
403    }
404}
405
406/// A low-level interface for reading JSON data from a byte stream
407///
408/// See [`Reader`] for a higher-level interface for interface with [`BufRead`]
409///
410/// The push-based interface facilitates integration with sources that yield arbitrarily
411/// delimited bytes ranges, such as [`BufRead`], or a chunked byte stream received from
412/// object storage
413///
414/// ```
415/// # use std::io::BufRead;
416/// # use arrow_array::RecordBatch;
417/// # use arrow_json::reader::{Decoder, ReaderBuilder};
418/// # use arrow_schema::{ArrowError, SchemaRef};
419/// #
420/// fn read_from_json<R: BufRead>(
421///     mut reader: R,
422///     schema: SchemaRef,
423/// ) -> Result<impl Iterator<Item = Result<RecordBatch, ArrowError>>, ArrowError> {
424///     let mut decoder = ReaderBuilder::new(schema).build_decoder()?;
425///     let mut next = move || {
426///         loop {
427///             // Decoder is agnostic that buf doesn't contain whole records
428///             let buf = reader.fill_buf()?;
429///             if buf.is_empty() {
430///                 break; // Input exhausted
431///             }
432///             let read = buf.len();
433///             let decoded = decoder.decode(buf)?;
434///
435///             // Consume the number of bytes read
436///             reader.consume(decoded);
437///             if decoded != read {
438///                 break; // Read batch size
439///             }
440///         }
441///         decoder.flush()
442///     };
443///     Ok(std::iter::from_fn(move || next().transpose()))
444/// }
445/// ```
446pub struct Decoder {
447    tape_decoder: TapeDecoder,
448    decoder: Box<dyn ArrayDecoder>,
449    batch_size: usize,
450    is_field: bool,
451    schema: SchemaRef,
452}
453
454impl std::fmt::Debug for Decoder {
455    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
456        f.debug_struct("Decoder")
457            .field("schema", &self.schema)
458            .field("batch_size", &self.batch_size)
459            .finish()
460    }
461}
462
463impl Decoder {
464    /// Read JSON objects from `buf`, returning the number of bytes read
465    ///
466    /// This method returns once `batch_size` objects have been parsed since the
467    /// last call to [`Self::flush`], or `buf` is exhausted. Any remaining bytes
468    /// should be included in the next call to [`Self::decode`]
469    ///
470    /// There is no requirement that `buf` contains a whole number of records, facilitating
471    /// integration with arbitrary byte streams, such as those yielded by [`BufRead`]
472    pub fn decode(&mut self, buf: &[u8]) -> Result<usize, ArrowError> {
473        self.tape_decoder.decode(buf)
474    }
475
476    /// Serialize `rows` to this [`Decoder`]
477    ///
478    /// This provides a simple way to convert [serde]-compatible datastructures into arrow
479    /// [`RecordBatch`].
480    ///
481    /// Custom conversion logic as described in [arrow_array::builder] will likely outperform this,
482    /// especially where the schema is known at compile-time, however, this provides a mechanism
483    /// to get something up and running quickly
484    ///
485    /// It can be used with [`serde_json::Value`]
486    ///
487    /// ```
488    /// # use std::sync::Arc;
489    /// # use serde_json::{Value, json};
490    /// # use arrow_array::cast::AsArray;
491    /// # use arrow_array::types::Float32Type;
492    /// # use arrow_json::ReaderBuilder;
493    /// # use arrow_schema::{DataType, Field, Schema};
494    /// let json = vec![json!({"float": 2.3}), json!({"float": 5.7})];
495    ///
496    /// let schema = Schema::new(vec![Field::new("float", DataType::Float32, true)]);
497    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
498    ///
499    /// decoder.serialize(&json).unwrap();
500    /// let batch = decoder.flush().unwrap().unwrap();
501    /// assert_eq!(batch.num_rows(), 2);
502    /// assert_eq!(batch.num_columns(), 1);
503    /// let values = batch.column(0).as_primitive::<Float32Type>().values();
504    /// assert_eq!(values, &[2.3, 5.7])
505    /// ```
506    ///
507    /// Or with arbitrary [`Serialize`] types
508    ///
509    /// ```
510    /// # use std::sync::Arc;
511    /// # use arrow_json::ReaderBuilder;
512    /// # use arrow_schema::{DataType, Field, Schema};
513    /// # use serde::Serialize;
514    /// # use arrow_array::cast::AsArray;
515    /// # use arrow_array::types::{Float32Type, Int32Type};
516    /// #
517    /// #[derive(Serialize)]
518    /// struct MyStruct {
519    ///     int32: i32,
520    ///     float: f32,
521    /// }
522    ///
523    /// let schema = Schema::new(vec![
524    ///     Field::new("int32", DataType::Int32, false),
525    ///     Field::new("float", DataType::Float32, false),
526    /// ]);
527    ///
528    /// let rows = vec![
529    ///     MyStruct{ int32: 0, float: 3. },
530    ///     MyStruct{ int32: 4, float: 67.53 },
531    /// ];
532    ///
533    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
534    /// decoder.serialize(&rows).unwrap();
535    ///
536    /// let batch = decoder.flush().unwrap().unwrap();
537    ///
538    /// // Expect batch containing two columns
539    /// let int32 = batch.column(0).as_primitive::<Int32Type>();
540    /// assert_eq!(int32.values(), &[0, 4]);
541    ///
542    /// let float = batch.column(1).as_primitive::<Float32Type>();
543    /// assert_eq!(float.values(), &[3., 67.53]);
544    /// ```
545    ///
546    /// Or even complex nested types
547    ///
548    /// ```
549    /// # use std::collections::BTreeMap;
550    /// # use std::sync::Arc;
551    /// # use arrow_array::StructArray;
552    /// # use arrow_cast::display::{ArrayFormatter, FormatOptions};
553    /// # use arrow_json::ReaderBuilder;
554    /// # use arrow_schema::{DataType, Field, Fields, Schema};
555    /// # use serde::Serialize;
556    /// #
557    /// #[derive(Serialize)]
558    /// struct MyStruct {
559    ///     int32: i32,
560    ///     list: Vec<f64>,
561    ///     nested: Vec<Option<Nested>>,
562    /// }
563    ///
564    /// impl MyStruct {
565    ///     /// Returns the [`Fields`] for [`MyStruct`]
566    ///     fn fields() -> Fields {
567    ///         let nested = DataType::Struct(Nested::fields());
568    ///         Fields::from([
569    ///             Arc::new(Field::new("int32", DataType::Int32, false)),
570    ///             Arc::new(Field::new_list(
571    ///                 "list",
572    ///                 Field::new("element", DataType::Float64, false),
573    ///                 false,
574    ///             )),
575    ///             Arc::new(Field::new_list(
576    ///                 "nested",
577    ///                 Field::new("element", nested, true),
578    ///                 true,
579    ///             )),
580    ///         ])
581    ///     }
582    /// }
583    ///
584    /// #[derive(Serialize)]
585    /// struct Nested {
586    ///     map: BTreeMap<String, Vec<String>>
587    /// }
588    ///
589    /// impl Nested {
590    ///     /// Returns the [`Fields`] for [`Nested`]
591    ///     fn fields() -> Fields {
592    ///         let element = Field::new("element", DataType::Utf8, false);
593    ///         Fields::from([
594    ///             Arc::new(Field::new_map(
595    ///                 "map",
596    ///                 "entries",
597    ///                 Field::new("key", DataType::Utf8, false),
598    ///                 Field::new_list("value", element, false),
599    ///                 false, // sorted
600    ///                 false, // nullable
601    ///             ))
602    ///         ])
603    ///     }
604    /// }
605    ///
606    /// let data = vec![
607    ///     MyStruct {
608    ///         int32: 34,
609    ///         list: vec![1., 2., 34.],
610    ///         nested: vec![
611    ///             None,
612    ///             Some(Nested {
613    ///                 map: vec![
614    ///                     ("key1".to_string(), vec!["foo".to_string(), "bar".to_string()]),
615    ///                     ("key2".to_string(), vec!["baz".to_string()])
616    ///                 ].into_iter().collect()
617    ///             })
618    ///         ]
619    ///     },
620    ///     MyStruct {
621    ///         int32: 56,
622    ///         list: vec![],
623    ///         nested: vec![]
624    ///     },
625    ///     MyStruct {
626    ///         int32: 24,
627    ///         list: vec![-1., 245.],
628    ///         nested: vec![None]
629    ///     }
630    /// ];
631    ///
632    /// let schema = Schema::new(MyStruct::fields());
633    /// let mut decoder = ReaderBuilder::new(Arc::new(schema)).build_decoder().unwrap();
634    /// decoder.serialize(&data).unwrap();
635    /// let batch = decoder.flush().unwrap().unwrap();
636    /// assert_eq!(batch.num_rows(), 3);
637    /// assert_eq!(batch.num_columns(), 3);
638    ///
639    /// // Convert to StructArray to format
640    /// let s = StructArray::from(batch);
641    /// let options = FormatOptions::default().with_null("null");
642    /// let formatter = ArrayFormatter::try_new(&s, &options).unwrap();
643    ///
644    /// assert_eq!(&formatter.value(0).to_string(), "{int32: 34, list: [1.0, 2.0, 34.0], nested: [null, {map: {key1: [foo, bar], key2: [baz]}}]}");
645    /// assert_eq!(&formatter.value(1).to_string(), "{int32: 56, list: [], nested: []}");
646    /// assert_eq!(&formatter.value(2).to_string(), "{int32: 24, list: [-1.0, 245.0], nested: [null]}");
647    /// ```
648    ///
649    /// Note: this ignores any batch size setting, and always decodes all rows
650    ///
651    /// [serde]: https://docs.rs/serde/latest/serde/
652    pub fn serialize<S: Serialize>(&mut self, rows: &[S]) -> Result<(), ArrowError> {
653        self.tape_decoder.serialize(rows)
654    }
655
656    /// True if the decoder is currently part way through decoding a record.
657    pub fn has_partial_record(&self) -> bool {
658        self.tape_decoder.has_partial_row()
659    }
660
661    /// The number of unflushed records, including the partially decoded record (if any).
662    pub fn len(&self) -> usize {
663        self.tape_decoder.num_buffered_rows()
664    }
665
666    /// True if there are no records to flush, i.e. [`Self::len`] is zero.
667    pub fn is_empty(&self) -> bool {
668        self.len() == 0
669    }
670
671    /// Flushes the currently buffered data to a [`RecordBatch`]
672    ///
673    /// Returns `Ok(None)` if no buffered data, i.e. [`Self::is_empty`] is true.
674    ///
675    /// Note: This will return an error if called part way through decoding a record,
676    /// i.e. [`Self::has_partial_record`] is true.
677    pub fn flush(&mut self) -> Result<Option<RecordBatch>, ArrowError> {
678        let tape = self.tape_decoder.finish()?;
679
680        if tape.num_rows() == 0 {
681            return Ok(None);
682        }
683
684        // First offset is null sentinel
685        let mut next_object = 1;
686        let pos: Vec<_> = (0..tape.num_rows())
687            .map(|_| {
688                let next = tape.next(next_object, "row").unwrap();
689                std::mem::replace(&mut next_object, next)
690            })
691            .collect();
692
693        let decoded = self.decoder.decode(&tape, &pos)?;
694        self.tape_decoder.clear();
695
696        let batch = match self.is_field {
697            true => RecordBatch::try_new(self.schema.clone(), vec![decoded])?,
698            false => {
699                RecordBatch::from(decoded.as_struct().clone()).with_schema(self.schema.clone())?
700            }
701        };
702
703        Ok(Some(batch))
704    }
705}
706
707trait ArrayDecoder: Send {
708    /// Decode elements from `tape` starting at the indexes contained in `pos`
709    fn decode(&mut self, tape: &Tape<'_>, pos: &[u32]) -> Result<ArrayRef, ArrowError>;
710}
711
712/// Context for decoder creation, containing configuration.
713///
714/// This context is passed through the decoder creation process and contains
715/// all the configuration needed to create decoders recursively.
716pub struct DecoderContext {
717    /// Whether to coerce primitives to strings
718    coerce_primitive: bool,
719    /// Whether to validate struct fields strictly
720    strict_mode: bool,
721    /// How to decode struct fields
722    struct_mode: StructMode,
723    /// Whether to treat columns with incompatible types as missing (i.e. NULL)
724    ignore_type_conflicts: bool,
725}
726
727impl DecoderContext {
728    /// Returns whether to coerce primitive types (e.g., number to string)
729    pub fn coerce_primitive(&self) -> bool {
730        self.coerce_primitive
731    }
732
733    /// Returns whether to validate struct fields strictly
734    pub fn strict_mode(&self) -> bool {
735        self.strict_mode
736    }
737
738    /// Returns how to decode struct fields
739    pub fn struct_mode(&self) -> StructMode {
740        self.struct_mode
741    }
742
743    /// Returns whether to treat columns with incompatible types as missing (i.e. NULL)
744    pub fn ignore_type_conflicts(&self) -> bool {
745        self.ignore_type_conflicts
746    }
747
748    /// Create a decoder for a type.
749    ///
750    /// This is the standard way to create child decoders from within a decoder
751    /// implementation.
752    fn make_decoder(
753        &self,
754        data_type: &DataType,
755        is_nullable: bool,
756    ) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
757        make_decoder(self, data_type, is_nullable)
758    }
759}
760
761fn make_decoder(
762    ctx: &DecoderContext,
763    data_type: &DataType,
764    is_nullable: bool,
765) -> Result<Box<dyn ArrayDecoder>, ArrowError> {
766    macro_rules! primitive_decoder {
767        ($t:ty, $data_type:expr) => {
768            Ok(Box::new(PrimitiveArrayDecoder::<$t>::new(ctx, $data_type)))
769        };
770    }
771    macro_rules! timestamp_decoder {
772        ($t:ty, $data_type:expr, $tz:expr) => {{
773            Ok(Box::new(TimestampArrayDecoder::<$t, _>::new(
774                ctx, $data_type, $tz,
775            )))
776        }};
777    }
778    macro_rules! decimal_decoder {
779        ($t:ty, $p:expr, $s:expr) => {
780            Ok(Box::new(DecimalArrayDecoder::<$t>::new(ctx, $p, $s)))
781        };
782    }
783
784    downcast_integer! {
785        *data_type => (primitive_decoder, data_type),
786        DataType::Null => Ok(Box::new(NullArrayDecoder::new(ctx))),
787        DataType::Float16 => primitive_decoder!(Float16Type, data_type),
788        DataType::Float32 => primitive_decoder!(Float32Type, data_type),
789        DataType::Float64 => primitive_decoder!(Float64Type, data_type),
790        DataType::Timestamp(TimeUnit::Second, None) => {
791            timestamp_decoder!(TimestampSecondType, data_type, Utc)
792        },
793        DataType::Timestamp(TimeUnit::Millisecond, None) => {
794            timestamp_decoder!(TimestampMillisecondType, data_type, Utc)
795        },
796        DataType::Timestamp(TimeUnit::Microsecond, None) => {
797            timestamp_decoder!(TimestampMicrosecondType, data_type, Utc)
798        },
799        DataType::Timestamp(TimeUnit::Nanosecond, None) => {
800            timestamp_decoder!(TimestampNanosecondType, data_type, Utc)
801        },
802        DataType::Timestamp(TimeUnit::Second, Some(ref tz)) => {
803            let tz: Tz = tz.parse()?;
804            timestamp_decoder!(TimestampSecondType, data_type, tz)
805        },
806        DataType::Timestamp(TimeUnit::Millisecond, Some(ref tz)) => {
807            let tz: Tz = tz.parse()?;
808            timestamp_decoder!(TimestampMillisecondType, data_type, tz)
809        },
810        DataType::Timestamp(TimeUnit::Microsecond, Some(ref tz)) => {
811            let tz: Tz = tz.parse()?;
812            timestamp_decoder!(TimestampMicrosecondType, data_type, tz)
813        },
814        DataType::Timestamp(TimeUnit::Nanosecond, Some(ref tz)) => {
815            let tz: Tz = tz.parse()?;
816            timestamp_decoder!(TimestampNanosecondType, data_type, tz)
817        },
818        DataType::Date32 => primitive_decoder!(Date32Type, data_type),
819        DataType::Date64 => primitive_decoder!(Date64Type, data_type),
820        DataType::Time32(TimeUnit::Second) => primitive_decoder!(Time32SecondType, data_type),
821        DataType::Time32(TimeUnit::Millisecond) => primitive_decoder!(Time32MillisecondType, data_type),
822        DataType::Time64(TimeUnit::Microsecond) => primitive_decoder!(Time64MicrosecondType, data_type),
823        DataType::Time64(TimeUnit::Nanosecond) => primitive_decoder!(Time64NanosecondType, data_type),
824        DataType::Duration(TimeUnit::Nanosecond) => primitive_decoder!(DurationNanosecondType, data_type),
825        DataType::Duration(TimeUnit::Microsecond) => primitive_decoder!(DurationMicrosecondType, data_type),
826        DataType::Duration(TimeUnit::Millisecond) => primitive_decoder!(DurationMillisecondType, data_type),
827        DataType::Duration(TimeUnit::Second) => primitive_decoder!(DurationSecondType, data_type),
828        DataType::Decimal32(p, s) => decimal_decoder!(Decimal32Type, p, s),
829        DataType::Decimal64(p, s) => decimal_decoder!(Decimal64Type, p, s),
830        DataType::Decimal128(p, s) => decimal_decoder!(Decimal128Type, p, s),
831        DataType::Decimal256(p, s) => decimal_decoder!(Decimal256Type, p, s),
832        DataType::Boolean => Ok(Box::new(BooleanArrayDecoder::new(ctx))),
833        DataType::Utf8 => Ok(Box::new(StringArrayDecoder::<i32>::new(ctx))),
834        DataType::Utf8View => Ok(Box::new(StringViewArrayDecoder::new(ctx))),
835        DataType::LargeUtf8 => Ok(Box::new(StringArrayDecoder::<i64>::new(ctx))),
836        DataType::List(_) => Ok(Box::new(ListArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
837        DataType::LargeList(_) => Ok(Box::new(ListArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
838        DataType::ListView(_) => Ok(Box::new(ListViewArrayDecoder::<i32>::new(ctx, data_type, is_nullable)?)),
839        DataType::LargeListView(_) => Ok(Box::new(ListViewArrayDecoder::<i64>::new(ctx, data_type, is_nullable)?)),
840        DataType::FixedSizeList(_, _) => Ok(Box::new(FixedSizeListArrayDecoder::new(ctx, data_type, is_nullable)?)),
841        DataType::Struct(_) => Ok(Box::new(StructArrayDecoder::new(ctx, data_type, is_nullable)?)),
842        DataType::Binary => Ok(Box::new(BinaryArrayDecoder::<i32>::default())),
843        DataType::LargeBinary => Ok(Box::new(BinaryArrayDecoder::<i64>::default())),
844        DataType::FixedSizeBinary(len) => Ok(Box::new(FixedSizeBinaryArrayDecoder::new(len))),
845        DataType::BinaryView => Ok(Box::new(BinaryViewDecoder::default())),
846        DataType::Map(_, _) => Ok(Box::new(MapArrayDecoder::new(ctx, data_type, is_nullable)?)),
847        DataType::RunEndEncoded(ref r, _) => match r.data_type() {
848            DataType::Int16 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int16Type>::new(ctx, data_type, is_nullable)?)),
849            DataType::Int32 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int32Type>::new(ctx, data_type, is_nullable)?)),
850            DataType::Int64 => Ok(Box::new(RunEndEncodedArrayDecoder::<Int64Type>::new(ctx, data_type, is_nullable)?)),
851            d => unreachable!("unsupported run end index type: {d}"),
852        },
853        _ => Err(ArrowError::NotYetImplemented(format!("Support for {data_type} in JSON reader")))
854    }
855}
856
857#[cfg(test)]
858mod tests {
859    use arrow_array::cast::AsArray;
860    use arrow_array::{
861        Array, BooleanArray, Float64Array, GenericListViewArray, Int32Array, ListArray, MapArray,
862        NullArray, OffsetSizeTrait, StringArray, StringViewArray, StructArray,
863    };
864    use arrow_buffer::{ArrowNativeType, NullBuffer, OffsetBuffer, ScalarBuffer};
865    use arrow_cast::display::{ArrayFormatter, FormatOptions};
866    use arrow_schema::{Field, Fields};
867    use serde_json::json;
868    use std::fs::File;
869    use std::io::{BufReader, Cursor, Seek};
870
871    use super::*;
872
873    fn do_read(
874        buf: &str,
875        batch_size: usize,
876        coerce_primitive: bool,
877        strict_mode: bool,
878        schema: SchemaRef,
879    ) -> Vec<RecordBatch> {
880        let mut unbuffered = vec![];
881
882        // Test with different batch sizes to test for boundary conditions
883        for batch_size in [1, 3, 100, batch_size] {
884            unbuffered = ReaderBuilder::new(schema.clone())
885                .with_batch_size(batch_size)
886                .with_coerce_primitive(coerce_primitive)
887                .build(Cursor::new(buf.as_bytes()))
888                .unwrap()
889                .collect::<Result<Vec<_>, _>>()
890                .unwrap();
891
892            for b in unbuffered.iter().take(unbuffered.len() - 1) {
893                assert_eq!(b.num_rows(), batch_size)
894            }
895
896            // Test with different buffer sizes to test for boundary conditions
897            for b in [1, 3, 5] {
898                let buffered = ReaderBuilder::new(schema.clone())
899                    .with_batch_size(batch_size)
900                    .with_coerce_primitive(coerce_primitive)
901                    .with_strict_mode(strict_mode)
902                    .build(BufReader::with_capacity(b, Cursor::new(buf.as_bytes())))
903                    .unwrap()
904                    .collect::<Result<Vec<_>, _>>()
905                    .unwrap();
906                assert_eq!(unbuffered, buffered);
907            }
908        }
909
910        unbuffered
911    }
912
913    #[test]
914    fn test_basic() {
915        let buf = r#"
916        {"a": 1, "b": 2, "c": true, "d": 1}
917        {"a": 2E0, "b": 4, "c": false, "d": 2, "e": 254}
918
919        {"b": 6, "a": 2.0, "d": 45}
920        {"b": "5", "a": 2}
921        {"b": 4e0}
922        {"b": 7, "a": null}
923        "#;
924
925        let schema = Arc::new(Schema::new(vec![
926            Field::new("a", DataType::Int64, true),
927            Field::new("b", DataType::Int32, true),
928            Field::new("c", DataType::Boolean, true),
929            Field::new("d", DataType::Date32, true),
930            Field::new("e", DataType::Date64, true),
931        ]));
932
933        let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
934        assert!(decoder.is_empty());
935        assert_eq!(decoder.len(), 0);
936        assert!(!decoder.has_partial_record());
937        assert_eq!(decoder.decode(buf.as_bytes()).unwrap(), 221);
938        assert!(!decoder.is_empty());
939        assert_eq!(decoder.len(), 6);
940        assert!(!decoder.has_partial_record());
941        let batch = decoder.flush().unwrap().unwrap();
942        assert_eq!(batch.num_rows(), 6);
943        assert!(decoder.is_empty());
944        assert_eq!(decoder.len(), 0);
945        assert!(!decoder.has_partial_record());
946
947        let batches = do_read(buf, 1024, false, false, schema);
948        assert_eq!(batches.len(), 1);
949
950        let col1 = batches[0].column(0).as_primitive::<Int64Type>();
951        assert_eq!(col1.null_count(), 2);
952        assert_eq!(col1.values(), &[1, 2, 2, 2, 0, 0]);
953        assert!(col1.is_null(4));
954        assert!(col1.is_null(5));
955
956        let col2 = batches[0].column(1).as_primitive::<Int32Type>();
957        assert_eq!(col2.null_count(), 0);
958        assert_eq!(col2.values(), &[2, 4, 6, 5, 4, 7]);
959
960        let col3 = batches[0].column(2).as_boolean();
961        assert_eq!(col3.null_count(), 4);
962        assert!(col3.value(0));
963        assert!(!col3.is_null(0));
964        assert!(!col3.value(1));
965        assert!(!col3.is_null(1));
966
967        let col4 = batches[0].column(3).as_primitive::<Date32Type>();
968        assert_eq!(col4.null_count(), 3);
969        assert!(col4.is_null(3));
970        assert_eq!(col4.values(), &[1, 2, 45, 0, 0, 0]);
971
972        let col5 = batches[0].column(4).as_primitive::<Date64Type>();
973        assert_eq!(col5.null_count(), 5);
974        assert!(col5.is_null(0));
975        assert!(col5.is_null(2));
976        assert!(col5.is_null(3));
977        assert_eq!(col5.values(), &[0, 254, 0, 0, 0, 0]);
978    }
979
980    #[test]
981    fn test_string() {
982        let buf = r#"
983        {"a": "1", "b": "2"}
984        {"a": "hello", "b": "shoo"}
985        {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
986
987        {"b": null}
988        {"b": "", "a": null}
989
990        "#;
991        let schema = Arc::new(Schema::new(vec![
992            Field::new("a", DataType::Utf8, true),
993            Field::new("b", DataType::LargeUtf8, true),
994        ]));
995
996        let batches = do_read(buf, 1024, false, false, schema);
997        assert_eq!(batches.len(), 1);
998
999        let col1 = batches[0].column(0).as_string::<i32>();
1000        assert_eq!(col1.null_count(), 2);
1001        assert_eq!(col1.value(0), "1");
1002        assert_eq!(col1.value(1), "hello");
1003        assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1004        assert!(col1.is_null(3));
1005        assert!(col1.is_null(4));
1006
1007        let col2 = batches[0].column(1).as_string::<i64>();
1008        assert_eq!(col2.null_count(), 1);
1009        assert_eq!(col2.value(0), "2");
1010        assert_eq!(col2.value(1), "shoo");
1011        assert_eq!(col2.value(2), "\t😁foo");
1012        assert!(col2.is_null(3));
1013        assert_eq!(col2.value(4), "");
1014    }
1015
1016    #[test]
1017    fn test_long_string_view_allocation() {
1018        // The JSON input contains field "a" with different string lengths.
1019        // According to the implementation in the decoder:
1020        // - For a string, capacity is only increased if its length > 12 bytes.
1021        // Therefore, for:
1022        // Row 1: "short" (5 bytes) -> capacity += 0
1023        // Row 2: "this is definitely long" (24 bytes) -> capacity += 24
1024        // Row 3: "hello" (5 bytes) -> capacity += 0
1025        // Row 4: "\nfoobar😀asfgÿ" (17 bytes) -> capacity += 17
1026        // Expected total capacity = 24 + 17 = 41
1027        let expected_capacity: usize = 41;
1028
1029        let buf = r#"
1030        {"a": "short", "b": "dummy"}
1031        {"a": "this is definitely long", "b": "dummy"}
1032        {"a": "hello", "b": "dummy"}
1033        {"a": "\nfoobar😀asfgÿ", "b": "dummy"}
1034        "#;
1035
1036        let schema = Arc::new(Schema::new(vec![
1037            Field::new("a", DataType::Utf8View, true),
1038            Field::new("b", DataType::LargeUtf8, true),
1039        ]));
1040
1041        let batches = do_read(buf, 1024, false, false, schema);
1042        assert_eq!(batches.len(), 1, "Expected one record batch");
1043
1044        // Get the first column ("a") as a StringViewArray.
1045        let col_a = batches[0].column(0);
1046        let string_view_array = col_a
1047            .as_any()
1048            .downcast_ref::<StringViewArray>()
1049            .expect("Column should be a StringViewArray");
1050
1051        // Retrieve the underlying data buffer from the array.
1052        // The builder pre-allocates capacity based on the sum of lengths for long strings.
1053        let data_buffer = string_view_array.to_data().buffers()[0].len();
1054
1055        // Check that the allocated capacity is at least what we expected.
1056        // (The actual buffer may be larger than expected due to rounding or internal allocation strategies.)
1057        assert!(
1058            data_buffer >= expected_capacity,
1059            "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1060        );
1061
1062        // Additionally, verify that the decoded values are correct.
1063        assert_eq!(string_view_array.value(0), "short");
1064        assert_eq!(string_view_array.value(1), "this is definitely long");
1065        assert_eq!(string_view_array.value(2), "hello");
1066        assert_eq!(string_view_array.value(3), "\nfoobar😀asfgÿ");
1067    }
1068
1069    /// Test the memory capacity allocation logic when converting numeric types to strings.
1070    #[test]
1071    fn test_numeric_view_allocation() {
1072        // For numeric types, the expected capacity calculation is as follows:
1073        // Row 1: 123456789  -> Number converts to the string "123456789" (length 9), 9 <= 12, so no capacity is added.
1074        // Row 2: 1000000000000 -> Treated as an I64 number; its string is "1000000000000" (length 13),
1075        //                        which is >12 and its absolute value is > 999_999_999_999, so 13 bytes are added.
1076        // Row 3: 3.1415 -> F32 number, a fixed estimate of 10 bytes is added.
1077        // Row 4: 2.718281828459045 -> F64 number, a fixed estimate of 10 bytes is added.
1078        // Total expected capacity = 13 + 10 + 10 = 33 bytes.
1079        let expected_capacity: usize = 33;
1080
1081        let buf = r#"
1082    {"n": 123456789}
1083    {"n": 1000000000000}
1084    {"n": 3.1415}
1085    {"n": 2.718281828459045}
1086    "#;
1087
1088        let schema = Arc::new(Schema::new(vec![Field::new("n", DataType::Utf8View, true)]));
1089
1090        let batches = do_read(buf, 1024, true, false, schema);
1091        assert_eq!(batches.len(), 1, "Expected one record batch");
1092
1093        let col_n = batches[0].column(0);
1094        let string_view_array = col_n
1095            .as_any()
1096            .downcast_ref::<StringViewArray>()
1097            .expect("Column should be a StringViewArray");
1098
1099        // Check that the underlying data buffer capacity is at least the expected value.
1100        let data_buffer = string_view_array.to_data().buffers()[0].len();
1101        assert!(
1102            data_buffer >= expected_capacity,
1103            "Data buffer length ({data_buffer}) should be at least {expected_capacity}",
1104        );
1105
1106        // Verify that the converted string values are correct.
1107        // Note: The format of the number converted to a string should match the actual implementation.
1108        assert_eq!(string_view_array.value(0), "123456789");
1109        assert_eq!(string_view_array.value(1), "1000000000000");
1110        assert_eq!(string_view_array.value(2), "3.1415");
1111        assert_eq!(string_view_array.value(3), "2.718281828459045");
1112    }
1113
1114    #[test]
1115    fn test_string_with_uft8view() {
1116        let buf = r#"
1117        {"a": "1", "b": "2"}
1118        {"a": "hello", "b": "shoo"}
1119        {"b": "\t😁foo", "a": "\nfoobar\ud83d\ude00\u0061\u0073\u0066\u0067\u00FF"}
1120
1121        {"b": null}
1122        {"b": "", "a": null}
1123
1124        "#;
1125        let schema = Arc::new(Schema::new(vec![
1126            Field::new("a", DataType::Utf8View, true),
1127            Field::new("b", DataType::LargeUtf8, true),
1128        ]));
1129
1130        let batches = do_read(buf, 1024, false, false, schema);
1131        assert_eq!(batches.len(), 1);
1132
1133        let col1 = batches[0].column(0).as_string_view();
1134        assert_eq!(col1.null_count(), 2);
1135        assert_eq!(col1.value(0), "1");
1136        assert_eq!(col1.value(1), "hello");
1137        assert_eq!(col1.value(2), "\nfoobar😀asfgÿ");
1138        assert!(col1.is_null(3));
1139        assert!(col1.is_null(4));
1140        assert_eq!(col1.data_type(), &DataType::Utf8View);
1141
1142        let col2 = batches[0].column(1).as_string::<i64>();
1143        assert_eq!(col2.null_count(), 1);
1144        assert_eq!(col2.value(0), "2");
1145        assert_eq!(col2.value(1), "shoo");
1146        assert_eq!(col2.value(2), "\t😁foo");
1147        assert!(col2.is_null(3));
1148        assert_eq!(col2.value(4), "");
1149    }
1150
1151    #[test]
1152    fn test_complex() {
1153        let buf = r#"
1154           {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3}, {"c": 4}]}}
1155           {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1156           {"list": null, "nested": {"a": null}}
1157        "#;
1158
1159        let schema = Arc::new(Schema::new(vec![
1160            Field::new_list("list", Field::new("element", DataType::Int32, false), true),
1161            Field::new_struct(
1162                "nested",
1163                vec![
1164                    Field::new("a", DataType::Int32, true),
1165                    Field::new("b", DataType::Int32, true),
1166                ],
1167                true,
1168            ),
1169            Field::new_struct(
1170                "nested_list",
1171                vec![Field::new_list(
1172                    "list2",
1173                    Field::new_struct(
1174                        "element",
1175                        vec![Field::new("c", DataType::Int32, false)],
1176                        false,
1177                    ),
1178                    true,
1179                )],
1180                true,
1181            ),
1182        ]));
1183
1184        let batches = do_read(buf, 1024, false, false, schema);
1185        assert_eq!(batches.len(), 1);
1186
1187        let list = batches[0].column(0).as_list::<i32>();
1188        assert_eq!(list.len(), 3);
1189        assert_eq!(list.value_offsets(), &[0, 0, 2, 2]);
1190        assert_eq!(list.null_count(), 1);
1191        assert!(list.is_null(2));
1192        let list_values = list.values().as_primitive::<Int32Type>();
1193        assert_eq!(list_values.values(), &[5, 6]);
1194
1195        let nested = batches[0].column(1).as_struct();
1196        let a = nested.column(0).as_primitive::<Int32Type>();
1197        assert_eq!(list.null_count(), 1);
1198        assert_eq!(a.values(), &[1, 7, 0]);
1199        assert!(list.is_null(2));
1200
1201        let b = nested.column(1).as_primitive::<Int32Type>();
1202        assert_eq!(b.null_count(), 2);
1203        assert_eq!(b.len(), 3);
1204        assert_eq!(b.value(0), 2);
1205        assert!(b.is_null(1));
1206        assert!(b.is_null(2));
1207
1208        let nested_list = batches[0].column(2).as_struct();
1209        assert_eq!(nested_list.len(), 3);
1210        assert_eq!(nested_list.null_count(), 1);
1211        assert!(nested_list.is_null(2));
1212
1213        let list2 = nested_list.column(0).as_list::<i32>();
1214        assert_eq!(list2.len(), 3);
1215        assert_eq!(list2.null_count(), 1);
1216        assert_eq!(list2.value_offsets(), &[0, 2, 2, 2]);
1217        assert!(list2.is_null(2));
1218
1219        let list2_values = list2.values().as_struct();
1220
1221        let c = list2_values.column(0).as_primitive::<Int32Type>();
1222        assert_eq!(c.values(), &[3, 4]);
1223    }
1224
1225    #[test]
1226    fn test_projection() {
1227        let buf = r#"
1228           {"list": [], "nested": {"a": 1, "b": 2}, "nested_list": {"list2": [{"c": 3, "d": 5}, {"c": 4}]}}
1229           {"list": [5, 6], "nested": {"a": 7}, "nested_list": {"list2": []}}
1230        "#;
1231
1232        let schema = Arc::new(Schema::new(vec![
1233            Field::new_struct(
1234                "nested",
1235                vec![Field::new("a", DataType::Int32, false)],
1236                true,
1237            ),
1238            Field::new_struct(
1239                "nested_list",
1240                vec![Field::new_list(
1241                    "list2",
1242                    Field::new_struct(
1243                        "element",
1244                        vec![Field::new("d", DataType::Int32, true)],
1245                        false,
1246                    ),
1247                    true,
1248                )],
1249                true,
1250            ),
1251        ]));
1252
1253        let batches = do_read(buf, 1024, false, false, schema);
1254        assert_eq!(batches.len(), 1);
1255
1256        let nested = batches[0].column(0).as_struct();
1257        assert_eq!(nested.num_columns(), 1);
1258        let a = nested.column(0).as_primitive::<Int32Type>();
1259        assert_eq!(a.null_count(), 0);
1260        assert_eq!(a.values(), &[1, 7]);
1261
1262        let nested_list = batches[0].column(1).as_struct();
1263        assert_eq!(nested_list.num_columns(), 1);
1264        assert_eq!(nested_list.null_count(), 0);
1265
1266        let list2 = nested_list.column(0).as_list::<i32>();
1267        assert_eq!(list2.value_offsets(), &[0, 2, 2]);
1268        assert_eq!(list2.null_count(), 0);
1269
1270        let child = list2.values().as_struct();
1271        assert_eq!(child.num_columns(), 1);
1272        assert_eq!(child.len(), 2);
1273        assert_eq!(child.null_count(), 0);
1274
1275        let c = child.column(0).as_primitive::<Int32Type>();
1276        assert_eq!(c.values(), &[5, 0]);
1277        assert_eq!(c.null_count(), 1);
1278        assert!(c.is_null(1));
1279    }
1280
1281    #[test]
1282    fn test_map() {
1283        let buf = r#"
1284           {"map": {"a": ["foo", null]}}
1285           {"map": {"a": [null], "b": []}}
1286           {"map": {"c": null, "a": ["baz"]}}
1287        "#;
1288        let map = Field::new_map(
1289            "map",
1290            Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1291            Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1292            Field::new_list(
1293                Field::MAP_VALUE_FIELD_DEFAULT_NAME,
1294                Field::new("element", DataType::Utf8, true),
1295                true,
1296            ),
1297            false,
1298            true,
1299        );
1300
1301        let schema = Arc::new(Schema::new(vec![map]));
1302
1303        let batches = do_read(buf, 1024, false, false, schema);
1304        assert_eq!(batches.len(), 1);
1305
1306        let map = batches[0].column(0).as_map();
1307        let map_keys = map.keys().as_string::<i32>();
1308        let map_values = map.values().as_list::<i32>();
1309        assert_eq!(map.value_offsets(), &[0, 1, 3, 5]);
1310
1311        let k: Vec<_> = map_keys.iter().flatten().collect();
1312        assert_eq!(&k, &["a", "a", "b", "c", "a"]);
1313
1314        let list_values = map_values.values().as_string::<i32>();
1315        let lv: Vec<_> = list_values.iter().collect();
1316        assert_eq!(&lv, &[Some("foo"), None, None, Some("baz")]);
1317        assert_eq!(map_values.value_offsets(), &[0, 2, 3, 3, 3, 4]);
1318        assert_eq!(map_values.null_count(), 1);
1319        assert!(map_values.is_null(3));
1320
1321        let options = FormatOptions::default().with_null("null");
1322        let formatter = ArrayFormatter::try_new(map, &options).unwrap();
1323        assert_eq!(formatter.value(0).to_string(), "{a: [foo, null]}");
1324        assert_eq!(formatter.value(1).to_string(), "{a: [null], b: []}");
1325        assert_eq!(formatter.value(2).to_string(), "{c: null, a: [baz]}");
1326    }
1327
1328    #[test]
1329    fn test_map_non_nullable_value() {
1330        let map = Field::new_map(
1331            "map",
1332            Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
1333            Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1334            Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, false),
1335            false,
1336            false,
1337        );
1338        let schema = Arc::new(Schema::new(vec![map]));
1339        let buf = r#"{"map": {"key": null}}"#;
1340
1341        let err = ReaderBuilder::new(schema)
1342            .build(Cursor::new(buf.as_bytes()))
1343            .unwrap()
1344            .read()
1345            .unwrap_err();
1346
1347        assert_eq!(
1348            err.to_string(),
1349            "Invalid argument error: Found unmasked nulls for non-nullable StructArray field \"value\""
1350        );
1351    }
1352
1353    #[test]
1354    fn test_not_coercing_primitive_into_string_without_flag() {
1355        let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, true)]));
1356
1357        let buf = r#"{"a": 1}"#;
1358        let err = ReaderBuilder::new(schema.clone())
1359            .with_batch_size(1024)
1360            .build(Cursor::new(buf.as_bytes()))
1361            .unwrap()
1362            .read()
1363            .unwrap_err();
1364
1365        assert_eq!(
1366            err.to_string(),
1367            "Json error: whilst decoding field 'a': expected string got 1"
1368        );
1369
1370        let buf = r#"{"a": true}"#;
1371        let err = ReaderBuilder::new(schema)
1372            .with_batch_size(1024)
1373            .build(Cursor::new(buf.as_bytes()))
1374            .unwrap()
1375            .read()
1376            .unwrap_err();
1377
1378        assert_eq!(
1379            err.to_string(),
1380            "Json error: whilst decoding field 'a': expected string got true"
1381        );
1382    }
1383
1384    #[test]
1385    fn test_coercing_primitive_into_string() {
1386        let buf = r#"
1387        {"a": 1, "b": 2, "c": true}
1388        {"a": 2E0, "b": 4, "c": false}
1389
1390        {"b": 6, "a": 2.0}
1391        {"b": "5", "a": 2}
1392        {"b": 4e0}
1393        {"b": 7, "a": null}
1394        "#;
1395
1396        let schema = Arc::new(Schema::new(vec![
1397            Field::new("a", DataType::Utf8, true),
1398            Field::new("b", DataType::Utf8, true),
1399            Field::new("c", DataType::Utf8, true),
1400        ]));
1401
1402        let batches = do_read(buf, 1024, true, false, schema);
1403        assert_eq!(batches.len(), 1);
1404
1405        let col1 = batches[0].column(0).as_string::<i32>();
1406        assert_eq!(col1.null_count(), 2);
1407        assert_eq!(col1.value(0), "1");
1408        assert_eq!(col1.value(1), "2E0");
1409        assert_eq!(col1.value(2), "2.0");
1410        assert_eq!(col1.value(3), "2");
1411        assert!(col1.is_null(4));
1412        assert!(col1.is_null(5));
1413
1414        let col2 = batches[0].column(1).as_string::<i32>();
1415        assert_eq!(col2.null_count(), 0);
1416        assert_eq!(col2.value(0), "2");
1417        assert_eq!(col2.value(1), "4");
1418        assert_eq!(col2.value(2), "6");
1419        assert_eq!(col2.value(3), "5");
1420        assert_eq!(col2.value(4), "4e0");
1421        assert_eq!(col2.value(5), "7");
1422
1423        let col3 = batches[0].column(2).as_string::<i32>();
1424        assert_eq!(col3.null_count(), 4);
1425        assert_eq!(col3.value(0), "true");
1426        assert_eq!(col3.value(1), "false");
1427        assert!(col3.is_null(2));
1428        assert!(col3.is_null(3));
1429        assert!(col3.is_null(4));
1430        assert!(col3.is_null(5));
1431    }
1432
1433    fn test_decimal<T: DecimalType>(data_type: DataType) {
1434        let buf = r#"
1435        {"a": 1, "b": 2, "c": 38.30}
1436        {"a": 2, "b": 4, "c": 123.456}
1437
1438        {"b": 1337, "a": "2.0452"}
1439        {"b": "5", "a": "11034.2"}
1440        {"b": 40}
1441        {"b": 1234, "a": null}
1442        "#;
1443
1444        let schema = Arc::new(Schema::new(vec![
1445            Field::new("a", data_type.clone(), true),
1446            Field::new("b", data_type.clone(), true),
1447            Field::new("c", data_type, true),
1448        ]));
1449
1450        let batches = do_read(buf, 1024, true, false, schema);
1451        assert_eq!(batches.len(), 1);
1452
1453        let col1 = batches[0].column(0).as_primitive::<T>();
1454        assert_eq!(col1.null_count(), 2);
1455        assert!(col1.is_null(4));
1456        assert!(col1.is_null(5));
1457        assert_eq!(
1458            col1.values(),
1459            &[100, 200, 204, 1103420, 0, 0].map(T::Native::usize_as)
1460        );
1461
1462        let col2 = batches[0].column(1).as_primitive::<T>();
1463        assert_eq!(col2.null_count(), 0);
1464        assert_eq!(
1465            col2.values(),
1466            &[200, 400, 133700, 500, 4000, 123400].map(T::Native::usize_as)
1467        );
1468
1469        let col3 = batches[0].column(2).as_primitive::<T>();
1470        assert_eq!(col3.null_count(), 4);
1471        assert!(!col3.is_null(0));
1472        assert!(!col3.is_null(1));
1473        assert!(col3.is_null(2));
1474        assert!(col3.is_null(3));
1475        assert!(col3.is_null(4));
1476        assert!(col3.is_null(5));
1477        assert_eq!(
1478            col3.values(),
1479            &[3830, 12345, 0, 0, 0, 0].map(T::Native::usize_as)
1480        );
1481    }
1482
1483    #[test]
1484    fn test_decimals() {
1485        test_decimal::<Decimal32Type>(DataType::Decimal32(8, 2));
1486        test_decimal::<Decimal64Type>(DataType::Decimal64(10, 2));
1487        test_decimal::<Decimal128Type>(DataType::Decimal128(10, 2));
1488        test_decimal::<Decimal256Type>(DataType::Decimal256(10, 2));
1489    }
1490
1491    fn test_timestamp<T: ArrowTimestampType>() {
1492        let buf = r#"
1493        {"a": 1, "b": "2020-09-08T13:42:29.190855+00:00", "c": 38.30, "d": "1997-01-31T09:26:56.123"}
1494        {"a": 2, "b": "2020-09-08T13:42:29.190855Z", "c": 123.456, "d": 123.456}
1495
1496        {"b": 1337, "b": "2020-09-08T13:42:29Z", "c": "1997-01-31T09:26:56.123", "d": "1997-01-31T09:26:56.123Z"}
1497        {"b": 40, "c": "2020-09-08T13:42:29.190855+00:00", "d": "1997-01-31 09:26:56.123-05:00"}
1498        {"b": 1234, "a": null, "c": "1997-01-31 09:26:56.123Z", "d": "1997-01-31 092656"}
1499        {"c": "1997-01-31T14:26:56.123-05:00", "d": "1997-01-31"}
1500        "#;
1501
1502        let with_timezone = DataType::Timestamp(T::UNIT, Some("+08:00".into()));
1503        let schema = Arc::new(Schema::new(vec![
1504            Field::new("a", T::DATA_TYPE, true),
1505            Field::new("b", T::DATA_TYPE, true),
1506            Field::new("c", T::DATA_TYPE, true),
1507            Field::new("d", with_timezone, true),
1508        ]));
1509
1510        let batches = do_read(buf, 1024, true, false, schema);
1511        assert_eq!(batches.len(), 1);
1512
1513        let unit_in_nanos: i64 = match T::UNIT {
1514            TimeUnit::Second => 1_000_000_000,
1515            TimeUnit::Millisecond => 1_000_000,
1516            TimeUnit::Microsecond => 1_000,
1517            TimeUnit::Nanosecond => 1,
1518        };
1519
1520        let col1 = batches[0].column(0).as_primitive::<T>();
1521        assert_eq!(col1.null_count(), 4);
1522        assert!(col1.is_null(2));
1523        assert!(col1.is_null(3));
1524        assert!(col1.is_null(4));
1525        assert!(col1.is_null(5));
1526        assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1527
1528        let col2 = batches[0].column(1).as_primitive::<T>();
1529        assert_eq!(col2.null_count(), 1);
1530        assert!(col2.is_null(5));
1531        assert_eq!(
1532            col2.values(),
1533            &[
1534                1599572549190855000 / unit_in_nanos,
1535                1599572549190855000 / unit_in_nanos,
1536                1599572549000000000 / unit_in_nanos,
1537                40,
1538                1234,
1539                0
1540            ]
1541        );
1542
1543        let col3 = batches[0].column(2).as_primitive::<T>();
1544        assert_eq!(col3.null_count(), 0);
1545        assert_eq!(
1546            col3.values(),
1547            &[
1548                38,
1549                123,
1550                854702816123000000 / unit_in_nanos,
1551                1599572549190855000 / unit_in_nanos,
1552                854702816123000000 / unit_in_nanos,
1553                854738816123000000 / unit_in_nanos
1554            ]
1555        );
1556
1557        let col4 = batches[0].column(3).as_primitive::<T>();
1558
1559        assert_eq!(col4.null_count(), 0);
1560        assert_eq!(
1561            col4.values(),
1562            &[
1563                854674016123000000 / unit_in_nanos,
1564                123,
1565                854702816123000000 / unit_in_nanos,
1566                854720816123000000 / unit_in_nanos,
1567                854674016000000000 / unit_in_nanos,
1568                854640000000000000 / unit_in_nanos
1569            ]
1570        );
1571    }
1572
1573    #[test]
1574    #[cfg_attr(miri, ignore)] // Takes too long
1575    fn test_timestamps() {
1576        test_timestamp::<TimestampSecondType>();
1577        test_timestamp::<TimestampMillisecondType>();
1578        test_timestamp::<TimestampMicrosecondType>();
1579        test_timestamp::<TimestampNanosecondType>();
1580    }
1581
1582    fn test_time<T: ArrowTemporalType>() {
1583        let buf = r#"
1584        {"a": 1, "b": "09:26:56.123 AM", "c": 38.30}
1585        {"a": 2, "b": "23:59:59", "c": 123.456}
1586
1587        {"b": 1337, "b": "6:00 pm", "c": "09:26:56.123"}
1588        {"b": 40, "c": "13:42:29.190855"}
1589        {"b": 1234, "a": null, "c": "09:26:56.123"}
1590        {"c": "14:26:56.123"}
1591        "#;
1592
1593        let unit = match T::DATA_TYPE {
1594            DataType::Time32(unit) | DataType::Time64(unit) => unit,
1595            _ => unreachable!(),
1596        };
1597
1598        let unit_in_nanos = match unit {
1599            TimeUnit::Second => 1_000_000_000,
1600            TimeUnit::Millisecond => 1_000_000,
1601            TimeUnit::Microsecond => 1_000,
1602            TimeUnit::Nanosecond => 1,
1603        };
1604
1605        let schema = Arc::new(Schema::new(vec![
1606            Field::new("a", T::DATA_TYPE, true),
1607            Field::new("b", T::DATA_TYPE, true),
1608            Field::new("c", T::DATA_TYPE, true),
1609        ]));
1610
1611        let batches = do_read(buf, 1024, true, false, schema);
1612        assert_eq!(batches.len(), 1);
1613
1614        let col1 = batches[0].column(0).as_primitive::<T>();
1615        assert_eq!(col1.null_count(), 4);
1616        assert!(col1.is_null(2));
1617        assert!(col1.is_null(3));
1618        assert!(col1.is_null(4));
1619        assert!(col1.is_null(5));
1620        assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1621
1622        let col2 = batches[0].column(1).as_primitive::<T>();
1623        assert_eq!(col2.null_count(), 1);
1624        assert!(col2.is_null(5));
1625        assert_eq!(
1626            col2.values(),
1627            &[
1628                34016123000000 / unit_in_nanos,
1629                86399000000000 / unit_in_nanos,
1630                64800000000000 / unit_in_nanos,
1631                40,
1632                1234,
1633                0
1634            ]
1635            .map(T::Native::usize_as)
1636        );
1637
1638        let col3 = batches[0].column(2).as_primitive::<T>();
1639        assert_eq!(col3.null_count(), 0);
1640        assert_eq!(
1641            col3.values(),
1642            &[
1643                38,
1644                123,
1645                34016123000000 / unit_in_nanos,
1646                49349190855000 / unit_in_nanos,
1647                34016123000000 / unit_in_nanos,
1648                52016123000000 / unit_in_nanos
1649            ]
1650            .map(T::Native::usize_as)
1651        );
1652    }
1653
1654    #[test]
1655    fn test_times() {
1656        test_time::<Time32MillisecondType>();
1657        test_time::<Time32SecondType>();
1658        test_time::<Time64MicrosecondType>();
1659        test_time::<Time64NanosecondType>();
1660    }
1661
1662    fn test_duration<T: ArrowTemporalType>() {
1663        let buf = r#"
1664        {"a": 1, "b": "2"}
1665        {"a": 3, "b": null}
1666        "#;
1667
1668        let schema = Arc::new(Schema::new(vec![
1669            Field::new("a", T::DATA_TYPE, true),
1670            Field::new("b", T::DATA_TYPE, true),
1671        ]));
1672
1673        let batches = do_read(buf, 1024, true, false, schema);
1674        assert_eq!(batches.len(), 1);
1675
1676        let col_a = batches[0].column_by_name("a").unwrap().as_primitive::<T>();
1677        assert_eq!(col_a.null_count(), 0);
1678        assert_eq!(col_a.values(), &[1, 3].map(T::Native::usize_as));
1679
1680        let col2 = batches[0].column_by_name("b").unwrap().as_primitive::<T>();
1681        assert_eq!(col2.null_count(), 1);
1682        assert_eq!(col2.values(), &[2, 0].map(T::Native::usize_as));
1683    }
1684
1685    #[test]
1686    fn test_durations() {
1687        test_duration::<DurationNanosecondType>();
1688        test_duration::<DurationMicrosecondType>();
1689        test_duration::<DurationMillisecondType>();
1690        test_duration::<DurationSecondType>();
1691    }
1692
1693    #[test]
1694    fn test_delta_checkpoint() {
1695        let json = "{\"protocol\":{\"minReaderVersion\":1,\"minWriterVersion\":2}}";
1696        let schema = Arc::new(Schema::new(vec![
1697            Field::new_struct(
1698                "protocol",
1699                vec![
1700                    Field::new("minReaderVersion", DataType::Int32, true),
1701                    Field::new("minWriterVersion", DataType::Int32, true),
1702                ],
1703                true,
1704            ),
1705            Field::new_struct(
1706                "add",
1707                vec![Field::new_map(
1708                    "partitionValues",
1709                    "key_value",
1710                    Field::new("key", DataType::Utf8, false),
1711                    Field::new("value", DataType::Utf8, true),
1712                    false,
1713                    false,
1714                )],
1715                true,
1716            ),
1717        ]));
1718
1719        let batches = do_read(json, 1024, true, false, schema);
1720        assert_eq!(batches.len(), 1);
1721
1722        let s: StructArray = batches.into_iter().next().unwrap().into();
1723        let opts = FormatOptions::default().with_null("null");
1724        let formatter = ArrayFormatter::try_new(&s, &opts).unwrap();
1725        assert_eq!(
1726            formatter.value(0).to_string(),
1727            "{protocol: {minReaderVersion: 1, minWriterVersion: 2}, add: null}"
1728        );
1729    }
1730
1731    #[test]
1732    fn struct_nullability() {
1733        let do_test = |child: DataType| {
1734            // Test correctly enforced nullability
1735            let non_null = r#"{"foo": {}}"#;
1736            let schema = Arc::new(Schema::new(vec![Field::new_struct(
1737                "foo",
1738                vec![Field::new("bar", child, false)],
1739                true,
1740            )]));
1741            let mut reader = ReaderBuilder::new(schema.clone())
1742                .build(Cursor::new(non_null.as_bytes()))
1743                .unwrap();
1744            assert!(reader.next().unwrap().is_err()); // Should error as not nullable
1745
1746            let null = r#"{"foo": {bar: null}}"#;
1747            let mut reader = ReaderBuilder::new(schema.clone())
1748                .build(Cursor::new(null.as_bytes()))
1749                .unwrap();
1750            assert!(reader.next().unwrap().is_err()); // Should error as not nullable
1751
1752            // Test nulls in nullable parent can mask nulls in non-nullable child
1753            let null = r#"{"foo": null}"#;
1754            let mut reader = ReaderBuilder::new(schema)
1755                .build(Cursor::new(null.as_bytes()))
1756                .unwrap();
1757            let batch = reader.next().unwrap().unwrap();
1758            assert_eq!(batch.num_columns(), 1);
1759            let foo = batch.column(0).as_struct();
1760            assert_eq!(foo.len(), 1);
1761            assert!(foo.is_null(0));
1762            assert_eq!(foo.num_columns(), 1);
1763
1764            let bar = foo.column(0);
1765            assert_eq!(bar.len(), 1);
1766            // Non-nullable child can still contain null as masked by parent
1767            assert!(bar.is_null(0));
1768        };
1769
1770        do_test(DataType::Boolean);
1771        do_test(DataType::Int32);
1772        do_test(DataType::Utf8);
1773        do_test(DataType::Decimal128(2, 1));
1774        do_test(DataType::Timestamp(
1775            TimeUnit::Microsecond,
1776            Some("+00:00".into()),
1777        ));
1778    }
1779
1780    #[test]
1781    fn test_truncation() {
1782        let buf = r#"
1783        {"i64": 9223372036854775807, "u64": 18446744073709551615 }
1784        {"i64": "9223372036854775807", "u64": "18446744073709551615" }
1785        {"i64": -9223372036854775808, "u64": 0 }
1786        {"i64": "-9223372036854775808", "u64": 0 }
1787        "#;
1788
1789        let schema = Arc::new(Schema::new(vec![
1790            Field::new("i64", DataType::Int64, true),
1791            Field::new("u64", DataType::UInt64, true),
1792        ]));
1793
1794        let batches = do_read(buf, 1024, true, false, schema);
1795        assert_eq!(batches.len(), 1);
1796
1797        let i64 = batches[0].column(0).as_primitive::<Int64Type>();
1798        assert_eq!(i64.values(), &[i64::MAX, i64::MAX, i64::MIN, i64::MIN]);
1799
1800        let u64 = batches[0].column(1).as_primitive::<UInt64Type>();
1801        assert_eq!(u64.values(), &[u64::MAX, u64::MAX, u64::MIN, u64::MIN]);
1802    }
1803
1804    #[test]
1805    fn test_timestamp_truncation() {
1806        let buf = r#"
1807        {"time": 9223372036854775807 }
1808        {"time": -9223372036854775808 }
1809        {"time": 9e5 }
1810        "#;
1811
1812        let schema = Arc::new(Schema::new(vec![Field::new(
1813            "time",
1814            DataType::Timestamp(TimeUnit::Nanosecond, None),
1815            true,
1816        )]));
1817
1818        let batches = do_read(buf, 1024, true, false, schema);
1819        assert_eq!(batches.len(), 1);
1820
1821        let i64 = batches[0]
1822            .column(0)
1823            .as_primitive::<TimestampNanosecondType>();
1824        assert_eq!(i64.values(), &[i64::MAX, i64::MIN, 900000]);
1825    }
1826
1827    #[test]
1828    fn test_strict_mode_no_missing_columns_in_schema() {
1829        let buf = r#"
1830        {"a": 1, "b": "2", "c": true}
1831        {"a": 2E0, "b": "4", "c": false}
1832        "#;
1833
1834        let schema = Arc::new(Schema::new(vec![
1835            Field::new("a", DataType::Int16, false),
1836            Field::new("b", DataType::Utf8, false),
1837            Field::new("c", DataType::Boolean, false),
1838        ]));
1839
1840        let batches = do_read(buf, 1024, true, true, schema);
1841        assert_eq!(batches.len(), 1);
1842
1843        let buf = r#"
1844        {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1845        {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1846        "#;
1847
1848        let schema = Arc::new(Schema::new(vec![
1849            Field::new("a", DataType::Int16, false),
1850            Field::new("b", DataType::Utf8, false),
1851            Field::new_struct(
1852                "c",
1853                vec![
1854                    Field::new("a", DataType::Boolean, false),
1855                    Field::new("b", DataType::Int16, false),
1856                ],
1857                false,
1858            ),
1859        ]));
1860
1861        let batches = do_read(buf, 1024, true, true, schema);
1862        assert_eq!(batches.len(), 1);
1863    }
1864
1865    #[test]
1866    fn test_strict_mode_missing_columns_in_schema() {
1867        let buf = r#"
1868        {"a": 1, "b": "2", "c": true}
1869        {"a": 2E0, "b": "4", "c": false}
1870        "#;
1871
1872        let schema = Arc::new(Schema::new(vec![
1873            Field::new("a", DataType::Int16, true),
1874            Field::new("c", DataType::Boolean, true),
1875        ]));
1876
1877        let err = ReaderBuilder::new(schema)
1878            .with_batch_size(1024)
1879            .with_strict_mode(true)
1880            .build(Cursor::new(buf.as_bytes()))
1881            .unwrap()
1882            .read()
1883            .unwrap_err();
1884
1885        assert_eq!(
1886            err.to_string(),
1887            "Json error: column 'b' missing from schema"
1888        );
1889
1890        let buf = r#"
1891        {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1892        {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1893        "#;
1894
1895        let schema = Arc::new(Schema::new(vec![
1896            Field::new("a", DataType::Int16, false),
1897            Field::new("b", DataType::Utf8, false),
1898            Field::new_struct("c", vec![Field::new("a", DataType::Boolean, false)], false),
1899        ]));
1900
1901        let err = ReaderBuilder::new(schema)
1902            .with_batch_size(1024)
1903            .with_strict_mode(true)
1904            .build(Cursor::new(buf.as_bytes()))
1905            .unwrap()
1906            .read()
1907            .unwrap_err();
1908
1909        assert_eq!(
1910            err.to_string(),
1911            "Json error: whilst decoding field 'c': column 'b' missing from schema"
1912        );
1913    }
1914
1915    fn read_file(path: &str, schema: Option<Schema>) -> Reader<BufReader<File>> {
1916        let file = File::open(path).unwrap();
1917        let mut reader = BufReader::new(file);
1918        let schema = schema.unwrap_or_else(|| {
1919            let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
1920            reader.rewind().unwrap();
1921            schema
1922        });
1923        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
1924        builder.build(reader).unwrap()
1925    }
1926
1927    #[test]
1928    fn test_json_basic() {
1929        let mut reader = read_file("test/data/basic.json", None);
1930        let batch = reader.next().unwrap().unwrap();
1931
1932        assert_eq!(8, batch.num_columns());
1933        assert_eq!(12, batch.num_rows());
1934
1935        let schema = reader.schema();
1936        let batch_schema = batch.schema();
1937        assert_eq!(schema, batch_schema);
1938
1939        let a = schema.column_with_name("a").unwrap();
1940        assert_eq!(0, a.0);
1941        assert_eq!(&DataType::Int64, a.1.data_type());
1942        let b = schema.column_with_name("b").unwrap();
1943        assert_eq!(1, b.0);
1944        assert_eq!(&DataType::Float64, b.1.data_type());
1945        let c = schema.column_with_name("c").unwrap();
1946        assert_eq!(2, c.0);
1947        assert_eq!(&DataType::Boolean, c.1.data_type());
1948        let d = schema.column_with_name("d").unwrap();
1949        assert_eq!(3, d.0);
1950        assert_eq!(&DataType::Utf8, d.1.data_type());
1951
1952        let aa = batch.column(a.0).as_primitive::<Int64Type>();
1953        assert_eq!(1, aa.value(0));
1954        assert_eq!(-10, aa.value(1));
1955        let bb = batch.column(b.0).as_primitive::<Float64Type>();
1956        assert_eq!(2.0, bb.value(0));
1957        assert_eq!(-3.5, bb.value(1));
1958        let cc = batch.column(c.0).as_boolean();
1959        assert!(!cc.value(0));
1960        assert!(cc.value(10));
1961        let dd = batch.column(d.0).as_string::<i32>();
1962        assert_eq!("4", dd.value(0));
1963        assert_eq!("text", dd.value(8));
1964    }
1965
1966    #[test]
1967    fn test_json_empty_projection() {
1968        let mut reader = read_file("test/data/basic.json", Some(Schema::empty()));
1969        let batch = reader.next().unwrap().unwrap();
1970
1971        assert_eq!(0, batch.num_columns());
1972        assert_eq!(12, batch.num_rows());
1973    }
1974
1975    #[test]
1976    fn test_json_basic_with_nulls() {
1977        let mut reader = read_file("test/data/basic_nulls.json", None);
1978        let batch = reader.next().unwrap().unwrap();
1979
1980        assert_eq!(4, batch.num_columns());
1981        assert_eq!(12, batch.num_rows());
1982
1983        let schema = reader.schema();
1984        let batch_schema = batch.schema();
1985        assert_eq!(schema, batch_schema);
1986
1987        let a = schema.column_with_name("a").unwrap();
1988        assert_eq!(&DataType::Int64, a.1.data_type());
1989        let b = schema.column_with_name("b").unwrap();
1990        assert_eq!(&DataType::Float64, b.1.data_type());
1991        let c = schema.column_with_name("c").unwrap();
1992        assert_eq!(&DataType::Boolean, c.1.data_type());
1993        let d = schema.column_with_name("d").unwrap();
1994        assert_eq!(&DataType::Utf8, d.1.data_type());
1995
1996        let aa = batch.column(a.0).as_primitive::<Int64Type>();
1997        assert!(aa.is_valid(0));
1998        assert!(!aa.is_valid(1));
1999        assert!(!aa.is_valid(11));
2000        let bb = batch.column(b.0).as_primitive::<Float64Type>();
2001        assert!(bb.is_valid(0));
2002        assert!(!bb.is_valid(2));
2003        assert!(!bb.is_valid(11));
2004        let cc = batch.column(c.0).as_boolean();
2005        assert!(cc.is_valid(0));
2006        assert!(!cc.is_valid(4));
2007        assert!(!cc.is_valid(11));
2008        let dd = batch.column(d.0).as_string::<i32>();
2009        assert!(!dd.is_valid(0));
2010        assert!(dd.is_valid(1));
2011        assert!(!dd.is_valid(4));
2012        assert!(!dd.is_valid(11));
2013    }
2014
2015    #[test]
2016    fn test_json_basic_schema() {
2017        let schema = Schema::new(vec![
2018            Field::new("a", DataType::Int64, true),
2019            Field::new("b", DataType::Float32, false),
2020            Field::new("c", DataType::Boolean, false),
2021            Field::new("d", DataType::Utf8, false),
2022        ]);
2023
2024        let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2025        let reader_schema = reader.schema();
2026        assert_eq!(reader_schema.as_ref(), &schema);
2027        let batch = reader.next().unwrap().unwrap();
2028
2029        assert_eq!(4, batch.num_columns());
2030        assert_eq!(12, batch.num_rows());
2031
2032        let schema = batch.schema();
2033
2034        let a = schema.column_with_name("a").unwrap();
2035        assert_eq!(&DataType::Int64, a.1.data_type());
2036        let b = schema.column_with_name("b").unwrap();
2037        assert_eq!(&DataType::Float32, b.1.data_type());
2038        let c = schema.column_with_name("c").unwrap();
2039        assert_eq!(&DataType::Boolean, c.1.data_type());
2040        let d = schema.column_with_name("d").unwrap();
2041        assert_eq!(&DataType::Utf8, d.1.data_type());
2042
2043        let aa = batch.column(a.0).as_primitive::<Int64Type>();
2044        assert_eq!(1, aa.value(0));
2045        assert_eq!(100000000000000, aa.value(11));
2046        let bb = batch.column(b.0).as_primitive::<Float32Type>();
2047        assert_eq!(2.0, bb.value(0));
2048        assert_eq!(-3.5, bb.value(1));
2049    }
2050
2051    #[test]
2052    fn test_json_basic_schema_projection() {
2053        let schema = Schema::new(vec![
2054            Field::new("a", DataType::Int64, true),
2055            Field::new("c", DataType::Boolean, false),
2056        ]);
2057
2058        let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2059        let batch = reader.next().unwrap().unwrap();
2060
2061        assert_eq!(2, batch.num_columns());
2062        assert_eq!(2, batch.schema().fields().len());
2063        assert_eq!(12, batch.num_rows());
2064
2065        assert_eq!(batch.schema().as_ref(), &schema);
2066
2067        let a = schema.column_with_name("a").unwrap();
2068        assert_eq!(0, a.0);
2069        assert_eq!(&DataType::Int64, a.1.data_type());
2070        let c = schema.column_with_name("c").unwrap();
2071        assert_eq!(1, c.0);
2072        assert_eq!(&DataType::Boolean, c.1.data_type());
2073    }
2074
2075    #[test]
2076    fn test_json_arrays() {
2077        let mut reader = read_file("test/data/arrays.json", None);
2078        let batch = reader.next().unwrap().unwrap();
2079
2080        assert_eq!(4, batch.num_columns());
2081        assert_eq!(3, batch.num_rows());
2082
2083        let schema = batch.schema();
2084
2085        let a = schema.column_with_name("a").unwrap();
2086        assert_eq!(&DataType::Int64, a.1.data_type());
2087        let b = schema.column_with_name("b").unwrap();
2088        assert_eq!(
2089            &DataType::List(Arc::new(Field::new_list_field(DataType::Float64, true))),
2090            b.1.data_type()
2091        );
2092        let c = schema.column_with_name("c").unwrap();
2093        assert_eq!(
2094            &DataType::List(Arc::new(Field::new_list_field(DataType::Boolean, true))),
2095            c.1.data_type()
2096        );
2097        let d = schema.column_with_name("d").unwrap();
2098        assert_eq!(&DataType::Utf8, d.1.data_type());
2099
2100        let aa = batch.column(a.0).as_primitive::<Int64Type>();
2101        assert_eq!(1, aa.value(0));
2102        assert_eq!(-10, aa.value(1));
2103        assert_eq!(1627668684594000000, aa.value(2));
2104        let bb = batch.column(b.0).as_list::<i32>();
2105        let bb = bb.values().as_primitive::<Float64Type>();
2106        assert_eq!(9, bb.len());
2107        assert_eq!(2.0, bb.value(0));
2108        assert_eq!(-6.1, bb.value(5));
2109        assert!(!bb.is_valid(7));
2110
2111        let cc = batch
2112            .column(c.0)
2113            .as_any()
2114            .downcast_ref::<ListArray>()
2115            .unwrap();
2116        let cc = cc.values().as_boolean();
2117        assert_eq!(6, cc.len());
2118        assert!(!cc.value(0));
2119        assert!(!cc.value(4));
2120        assert!(!cc.is_valid(5));
2121    }
2122
2123    #[test]
2124    fn test_empty_json_arrays() {
2125        let json_content = r#"
2126            {"items": []}
2127            {"items": null}
2128            {}
2129            "#;
2130
2131        let schema = Arc::new(Schema::new(vec![Field::new(
2132            "items",
2133            DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2134            true,
2135        )]));
2136
2137        let batches = do_read(json_content, 1024, false, false, schema);
2138        assert_eq!(batches.len(), 1);
2139
2140        let col1 = batches[0].column(0).as_list::<i32>();
2141        assert_eq!(col1.null_count(), 2);
2142        assert!(col1.value(0).is_empty());
2143        assert_eq!(col1.value(0).data_type(), &DataType::Null);
2144        assert!(col1.is_null(1));
2145        assert!(col1.is_null(2));
2146    }
2147
2148    #[test]
2149    fn test_nested_empty_json_arrays() {
2150        let json_content = r#"
2151            {"items": [[],[]]}
2152            {"items": [[null, null],[null]]}
2153            "#;
2154
2155        let schema = Arc::new(Schema::new(vec![Field::new(
2156            "items",
2157            DataType::List(FieldRef::new(Field::new_list_field(
2158                DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2159                true,
2160            ))),
2161            true,
2162        )]));
2163
2164        let batches = do_read(json_content, 1024, false, false, schema);
2165        assert_eq!(batches.len(), 1);
2166
2167        let col1 = batches[0].column(0).as_list::<i32>();
2168        assert_eq!(col1.null_count(), 0);
2169        assert_eq!(col1.value(0).len(), 2);
2170        assert!(col1.value(0).as_list::<i32>().value(0).is_empty());
2171        assert!(col1.value(0).as_list::<i32>().value(1).is_empty());
2172
2173        assert_eq!(col1.value(1).len(), 2);
2174        assert_eq!(col1.value(1).as_list::<i32>().value(0).len(), 2);
2175        assert_eq!(col1.value(1).as_list::<i32>().value(1).len(), 1);
2176    }
2177
2178    #[test]
2179    fn test_nested_list_json_arrays() {
2180        let c_field = Field::new_struct("c", vec![Field::new("d", DataType::Utf8, true)], true);
2181        let a_struct_field = Field::new_struct(
2182            "a",
2183            vec![Field::new("b", DataType::Boolean, true), c_field.clone()],
2184            true,
2185        );
2186        let a_field = Field::new("a", DataType::List(Arc::new(a_struct_field.clone())), true);
2187        let schema = Arc::new(Schema::new(vec![a_field.clone()]));
2188        let builder = ReaderBuilder::new(schema).with_batch_size(64);
2189        let json_content = r#"
2190        {"a": [{"b": true, "c": {"d": "a_text"}}, {"b": false, "c": {"d": "b_text"}}]}
2191        {"a": [{"b": false, "c": null}]}
2192        {"a": [{"b": true, "c": {"d": "c_text"}}, {"b": null, "c": {"d": "d_text"}}, {"b": true, "c": {"d": null}}]}
2193        {"a": null}
2194        {"a": []}
2195        {"a": [null]}
2196        "#;
2197        let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2198
2199        // build expected output
2200        let d = StringArray::from(vec![
2201            Some("a_text"),
2202            Some("b_text"),
2203            None,
2204            Some("c_text"),
2205            Some("d_text"),
2206            None,
2207            None,
2208        ]);
2209        let c = StructArray::new(
2210            vec![Field::new("d", DataType::Utf8, true)].into(),
2211            vec![Arc::new(d.clone()) as ArrayRef],
2212            Some(NullBuffer::from(vec![
2213                true, true, false, true, true, true, false,
2214            ])),
2215        );
2216        let b = BooleanArray::from(vec![
2217            Some(true),
2218            Some(false),
2219            Some(false),
2220            Some(true),
2221            None,
2222            Some(true),
2223            None,
2224        ]);
2225        let a = StructArray::new(
2226            vec![Field::new("b", DataType::Boolean, true), c_field.clone()].into(),
2227            vec![
2228                Arc::new(b.clone()) as ArrayRef,
2229                Arc::new(c.clone()) as ArrayRef,
2230            ],
2231            Some(NullBuffer::from(vec![
2232                true, true, true, true, true, true, false,
2233            ])),
2234        );
2235        let a_list = ListArray::new(
2236            Arc::new(a_struct_field.clone()),
2237            OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 6, 6, 6, 7])),
2238            Arc::new(a),
2239            Some(NullBuffer::from(vec![true, true, true, false, true, true])),
2240        );
2241
2242        // compare `a` with result from json reader
2243        let batch = reader.next().unwrap().unwrap();
2244        let read = batch.column(0);
2245        assert_eq!(read.len(), 6);
2246        // compare the arrays the long way around, to better detect differences
2247        let read: &ListArray = read.as_list::<i32>();
2248        let expected = &a_list;
2249        assert_eq!(read.value_offsets(), &[0, 2, 3, 6, 6, 6, 7]);
2250        // compare list null buffers
2251        assert_eq!(read.nulls(), expected.nulls());
2252        // build struct from list
2253        let struct_array = read.values().as_struct();
2254        let expected_struct_array = expected.values().as_struct();
2255
2256        assert_eq!(7, struct_array.len());
2257        assert_eq!(1, struct_array.null_count());
2258        assert_eq!(7, expected_struct_array.len());
2259        assert_eq!(1, expected_struct_array.null_count());
2260        // test struct's nulls
2261        assert_eq!(struct_array.nulls(), expected_struct_array.nulls());
2262        // test struct's fields
2263        let read_b = struct_array.column(0);
2264        assert_eq!(read_b.as_ref(), &b);
2265        let read_c = struct_array.column(1);
2266        assert_eq!(read_c.as_struct(), &c);
2267        let read_c = read_c.as_struct();
2268        let read_d = read_c.column(0);
2269        assert_eq!(read_d.as_ref(), &d);
2270
2271        assert_eq!(read, expected);
2272    }
2273
2274    fn assert_read_list_view<O: OffsetSizeTrait>() {
2275        let field = Arc::new(Field::new("item", DataType::Int32, true));
2276        let data_type = GenericListViewArray::<O>::DATA_TYPE_CONSTRUCTOR(field.clone());
2277        let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2278
2279        let buf = r#"
2280        {"lv": [1, 2, 3]}
2281        {"lv": [4, null]}
2282        {"lv": null}
2283        {"lv": [6]}
2284        {"lv": []}
2285        "#;
2286
2287        let batches = do_read(buf, 1024, false, false, schema);
2288        assert_eq!(batches.len(), 1);
2289        let batch = &batches[0];
2290        let col = batch.column(0);
2291        let list_view = col
2292            .as_any()
2293            .downcast_ref::<GenericListViewArray<O>>()
2294            .unwrap();
2295
2296        assert_eq!(list_view.len(), 5);
2297
2298        // Check offsets and sizes
2299        let expected_offsets: Vec<O> = vec![0, 3, 5, 5, 6]
2300            .into_iter()
2301            .map(|v| O::usize_as(v))
2302            .collect();
2303        let expected_sizes: Vec<O> = vec![3, 2, 0, 1, 0]
2304            .into_iter()
2305            .map(|v| O::usize_as(v))
2306            .collect();
2307        assert_eq!(list_view.value_offsets(), &expected_offsets);
2308        assert_eq!(list_view.value_sizes(), &expected_sizes);
2309
2310        // Row 0: [1, 2, 3]
2311        assert!(list_view.is_valid(0));
2312        let vals = list_view.value(0);
2313        let ints = vals.as_primitive::<Int32Type>();
2314        assert_eq!(ints.values(), &[1, 2, 3]);
2315
2316        // Row 1: [4, null]
2317        assert!(list_view.is_valid(1));
2318        let vals = list_view.value(1);
2319        let ints = vals.as_primitive::<Int32Type>();
2320        assert_eq!(ints.len(), 2);
2321        assert_eq!(ints.value(0), 4);
2322        assert!(ints.is_null(1));
2323
2324        // Row 2: null
2325        assert!(list_view.is_null(2));
2326
2327        // Row 3: [6]
2328        assert!(list_view.is_valid(3));
2329        let vals = list_view.value(3);
2330        let ints = vals.as_primitive::<Int32Type>();
2331        assert_eq!(ints.values(), &[6]);
2332
2333        // Row 4: []
2334        assert!(list_view.is_valid(4));
2335        let vals = list_view.value(4);
2336        assert_eq!(vals.len(), 0);
2337    }
2338
2339    #[test]
2340    fn test_read_list_view() {
2341        assert_read_list_view::<i32>();
2342        assert_read_list_view::<i64>();
2343    }
2344
2345    #[test]
2346    fn test_read_list_view_rejects_null_non_nullable_child() {
2347        let field = Arc::new(Field::new("item", DataType::Int32, false));
2348        for (data_type, array_type) in [
2349            (DataType::ListView(field.clone()), "ListViewArray"),
2350            (DataType::LargeListView(field.clone()), "LargeListViewArray"),
2351        ] {
2352            let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2353            let buf = r#"
2354            {"lv": [1, 2, 3]}
2355            {"lv": [4, null]}
2356            "#;
2357
2358            let error = ReaderBuilder::new(schema)
2359                .build(Cursor::new(buf.as_bytes()))
2360                .unwrap()
2361                .collect::<Result<Vec<_>, _>>()
2362                .unwrap_err();
2363
2364            assert_eq!(
2365                error.to_string(),
2366                format!(
2367                    "Invalid argument error: Non-nullable field of {array_type} \"item\" cannot contain nulls"
2368                )
2369            );
2370        }
2371    }
2372
2373    #[test]
2374    fn test_fixed_size_list() {
2375        let buf = r#"
2376        {"a": [1, 2, 3]}
2377        {"a": [4, 5, 6]}
2378        {"a": [7, 8, 9]}
2379        "#;
2380
2381        let field = Field::new_list_field(DataType::Int32, true);
2382        let schema = Arc::new(Schema::new(vec![Field::new(
2383            "a",
2384            DataType::FixedSizeList(Arc::new(field), 3),
2385            false,
2386        )]));
2387
2388        let batches = do_read(buf, 1024, false, false, schema);
2389        assert_eq!(batches.len(), 1);
2390
2391        let col = batches[0].column(0).as_fixed_size_list();
2392        assert_eq!(col.len(), 3);
2393        assert_eq!(col.value_length(), 3);
2394
2395        let values = col.values().as_primitive::<Int32Type>();
2396        assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8, 9]);
2397    }
2398
2399    #[test]
2400    fn test_fixed_size_list_nullable() {
2401        let buf = r#"
2402        {"a": [1, 2]}
2403        {"a": null}
2404        {"a": [3, null]}
2405        "#;
2406
2407        let field = Field::new_list_field(DataType::Int32, true);
2408        let schema = Arc::new(Schema::new(vec![Field::new(
2409            "a",
2410            DataType::FixedSizeList(Arc::new(field), 2),
2411            true,
2412        )]));
2413
2414        let batches = do_read(buf, 1024, false, false, schema);
2415        assert_eq!(batches.len(), 1);
2416
2417        let col = batches[0].column(0).as_fixed_size_list();
2418        assert_eq!(col.len(), 3);
2419        assert!(col.is_valid(0));
2420        assert!(col.is_null(1));
2421        assert!(col.is_valid(2));
2422
2423        let values = col.values().as_primitive::<Int32Type>();
2424        assert_eq!(values.value(0), 1);
2425        assert_eq!(values.value(1), 2);
2426        assert_eq!(values.value(4), 3);
2427        assert!(values.is_null(5));
2428    }
2429
2430    #[test]
2431    fn test_fixed_size_list_zero_size_non_nullable() {
2432        let buf = r#"
2433        {"a": []}
2434        {"a": []}
2435        {"a": []}
2436        "#;
2437
2438        let field = Field::new_list_field(DataType::Int32, true);
2439        let schema = Arc::new(Schema::new(vec![Field::new(
2440            "a",
2441            DataType::FixedSizeList(Arc::new(field), 0),
2442            false,
2443        )]));
2444
2445        let batches = do_read(buf, 1024, false, false, schema);
2446        assert_eq!(batches.len(), 1);
2447
2448        let col = batches[0].column(0).as_fixed_size_list();
2449        assert_eq!(col.len(), 3);
2450        assert_eq!(col.value_length(), 0);
2451
2452        let values = col.values().as_primitive::<Int32Type>();
2453        assert!(values.values().is_empty());
2454    }
2455
2456    #[test]
2457    fn test_fixed_size_list_wrong_size() {
2458        let buf = r#"{"a": [1, 2, 3]}"#;
2459
2460        let field = Field::new_list_field(DataType::Int32, true);
2461        let schema = Arc::new(Schema::new(vec![Field::new(
2462            "a",
2463            DataType::FixedSizeList(Arc::new(field), 2),
2464            false,
2465        )]));
2466
2467        let err = ReaderBuilder::new(schema)
2468            .build(Cursor::new(buf.as_bytes()))
2469            .unwrap()
2470            .next()
2471            .unwrap()
2472            .unwrap_err();
2473
2474        assert!(err.to_string().contains("expected 2 but got 3"), "{}", err);
2475    }
2476
2477    #[test]
2478    fn test_fixed_size_list_nested() {
2479        let buf = r#"
2480        {"a": [[1, 2], [3, 4]]}
2481        {"a": [[5, 6], [7, 8]]}
2482        "#;
2483
2484        let inner_field = Field::new_list_field(DataType::Int32, true);
2485        let inner_type = DataType::FixedSizeList(Arc::new(inner_field), 2);
2486        let outer_field = Arc::new(Field::new_list_field(inner_type.clone(), true));
2487        let schema = Arc::new(Schema::new(vec![Field::new(
2488            "a",
2489            DataType::FixedSizeList(outer_field, 2),
2490            false,
2491        )]));
2492
2493        let batches = do_read(buf, 1024, false, false, schema);
2494        assert_eq!(batches.len(), 1);
2495
2496        let col = batches[0].column(0).as_fixed_size_list();
2497        assert_eq!(col.len(), 2);
2498        assert_eq!(col.value_length(), 2);
2499
2500        let inner = col.values().as_fixed_size_list();
2501        assert_eq!(inner.len(), 4);
2502        assert_eq!(inner.value_length(), 2);
2503
2504        let values = inner.values().as_primitive::<Int32Type>();
2505        assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8]);
2506    }
2507
2508    #[test]
2509    fn test_fixed_size_list_ignore_type_conflicts() {
2510        let field = Field::new("item", DataType::Int32, true);
2511        let schema = Arc::new(Schema::new(vec![Field::new(
2512            "a",
2513            DataType::FixedSizeList(Arc::new(field), 2),
2514            true,
2515        )]));
2516
2517        let json = vec![
2518            json!({"a": [1, 2]}),
2519            json!({"a": "not a list"}),
2520            json!({"a": 42}),
2521            json!({"a": [6, 7]}),
2522        ];
2523
2524        let mut decoder = ReaderBuilder::new(schema)
2525            .with_ignore_type_conflicts(true)
2526            .build_decoder()
2527            .unwrap();
2528        decoder.serialize(&json).unwrap();
2529        let batch = decoder.flush().unwrap().unwrap();
2530
2531        let col = batch.column(0).as_fixed_size_list();
2532        assert_eq!(col.len(), 4);
2533        assert!(col.is_valid(0));
2534        assert!(col.is_null(1)); // string -> null
2535        assert!(col.is_null(2)); // number -> null
2536        assert!(col.is_valid(3));
2537
2538        let values = col.values().as_primitive::<Int32Type>();
2539        assert_eq!(values.value(0), 1);
2540        assert_eq!(values.value(1), 2);
2541        assert_eq!(values.value(6), 6);
2542        assert_eq!(values.value(7), 7);
2543    }
2544
2545    #[test]
2546    fn test_skip_empty_lines() {
2547        let schema = Schema::new(vec![Field::new("a", DataType::Int64, true)]);
2548        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
2549        let json_content = "
2550        {\"a\": 1}
2551        {\"a\": 2}
2552        {\"a\": 3}";
2553        let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2554        let batch = reader.next().unwrap().unwrap();
2555
2556        assert_eq!(1, batch.num_columns());
2557        assert_eq!(3, batch.num_rows());
2558
2559        let schema = reader.schema();
2560        let c = schema.column_with_name("a").unwrap();
2561        assert_eq!(&DataType::Int64, c.1.data_type());
2562    }
2563
2564    #[test]
2565    fn test_with_multiple_batches() {
2566        let file = File::open("test/data/basic_nulls.json").unwrap();
2567        let mut reader = BufReader::new(file);
2568        let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2569        reader.rewind().unwrap();
2570
2571        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2572        let mut reader = builder.build(reader).unwrap();
2573
2574        let mut num_records = Vec::new();
2575        while let Some(rb) = reader.next().transpose().unwrap() {
2576            num_records.push(rb.num_rows());
2577        }
2578
2579        assert_eq!(vec![5, 5, 2], num_records);
2580    }
2581
2582    #[test]
2583    fn test_timestamp_from_json_seconds() {
2584        let schema = Schema::new(vec![Field::new(
2585            "a",
2586            DataType::Timestamp(TimeUnit::Second, None),
2587            true,
2588        )]);
2589
2590        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2591        let batch = reader.next().unwrap().unwrap();
2592
2593        assert_eq!(1, batch.num_columns());
2594        assert_eq!(12, batch.num_rows());
2595
2596        let schema = reader.schema();
2597        let batch_schema = batch.schema();
2598        assert_eq!(schema, batch_schema);
2599
2600        let a = schema.column_with_name("a").unwrap();
2601        assert_eq!(
2602            &DataType::Timestamp(TimeUnit::Second, None),
2603            a.1.data_type()
2604        );
2605
2606        let aa = batch.column(a.0).as_primitive::<TimestampSecondType>();
2607        assert!(aa.is_valid(0));
2608        assert!(!aa.is_valid(1));
2609        assert!(!aa.is_valid(2));
2610        assert_eq!(1, aa.value(0));
2611        assert_eq!(1, aa.value(3));
2612        assert_eq!(5, aa.value(7));
2613    }
2614
2615    #[test]
2616    fn test_timestamp_from_json_milliseconds() {
2617        let schema = Schema::new(vec![Field::new(
2618            "a",
2619            DataType::Timestamp(TimeUnit::Millisecond, None),
2620            true,
2621        )]);
2622
2623        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2624        let batch = reader.next().unwrap().unwrap();
2625
2626        assert_eq!(1, batch.num_columns());
2627        assert_eq!(12, batch.num_rows());
2628
2629        let schema = reader.schema();
2630        let batch_schema = batch.schema();
2631        assert_eq!(schema, batch_schema);
2632
2633        let a = schema.column_with_name("a").unwrap();
2634        assert_eq!(
2635            &DataType::Timestamp(TimeUnit::Millisecond, None),
2636            a.1.data_type()
2637        );
2638
2639        let aa = batch.column(a.0).as_primitive::<TimestampMillisecondType>();
2640        assert!(aa.is_valid(0));
2641        assert!(!aa.is_valid(1));
2642        assert!(!aa.is_valid(2));
2643        assert_eq!(1, aa.value(0));
2644        assert_eq!(1, aa.value(3));
2645        assert_eq!(5, aa.value(7));
2646    }
2647
2648    #[test]
2649    fn test_date_from_json_milliseconds() {
2650        let schema = Schema::new(vec![Field::new("a", DataType::Date64, true)]);
2651
2652        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2653        let batch = reader.next().unwrap().unwrap();
2654
2655        assert_eq!(1, batch.num_columns());
2656        assert_eq!(12, batch.num_rows());
2657
2658        let schema = reader.schema();
2659        let batch_schema = batch.schema();
2660        assert_eq!(schema, batch_schema);
2661
2662        let a = schema.column_with_name("a").unwrap();
2663        assert_eq!(&DataType::Date64, a.1.data_type());
2664
2665        let aa = batch.column(a.0).as_primitive::<Date64Type>();
2666        assert!(aa.is_valid(0));
2667        assert!(!aa.is_valid(1));
2668        assert!(!aa.is_valid(2));
2669        assert_eq!(1, aa.value(0));
2670        assert_eq!(1, aa.value(3));
2671        assert_eq!(5, aa.value(7));
2672    }
2673
2674    #[test]
2675    fn test_time_from_json_nanoseconds() {
2676        let schema = Schema::new(vec![Field::new(
2677            "a",
2678            DataType::Time64(TimeUnit::Nanosecond),
2679            true,
2680        )]);
2681
2682        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2683        let batch = reader.next().unwrap().unwrap();
2684
2685        assert_eq!(1, batch.num_columns());
2686        assert_eq!(12, batch.num_rows());
2687
2688        let schema = reader.schema();
2689        let batch_schema = batch.schema();
2690        assert_eq!(schema, batch_schema);
2691
2692        let a = schema.column_with_name("a").unwrap();
2693        assert_eq!(&DataType::Time64(TimeUnit::Nanosecond), a.1.data_type());
2694
2695        let aa = batch.column(a.0).as_primitive::<Time64NanosecondType>();
2696        assert!(aa.is_valid(0));
2697        assert!(!aa.is_valid(1));
2698        assert!(!aa.is_valid(2));
2699        assert_eq!(1, aa.value(0));
2700        assert_eq!(1, aa.value(3));
2701        assert_eq!(5, aa.value(7));
2702    }
2703
2704    #[test]
2705    fn test_json_iterator() {
2706        let file = File::open("test/data/basic.json").unwrap();
2707        let mut reader = BufReader::new(file);
2708        let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2709        reader.rewind().unwrap();
2710
2711        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2712        let reader = builder.build(reader).unwrap();
2713        let schema = reader.schema();
2714        let (col_a_index, _) = schema.column_with_name("a").unwrap();
2715
2716        let mut sum_num_rows = 0;
2717        let mut num_batches = 0;
2718        let mut sum_a = 0;
2719        for batch in reader {
2720            let batch = batch.unwrap();
2721            assert_eq!(8, batch.num_columns());
2722            sum_num_rows += batch.num_rows();
2723            num_batches += 1;
2724            let batch_schema = batch.schema();
2725            assert_eq!(schema, batch_schema);
2726            let a_array = batch.column(col_a_index).as_primitive::<Int64Type>();
2727            sum_a += (0..a_array.len()).map(|i| a_array.value(i)).sum::<i64>();
2728        }
2729        assert_eq!(12, sum_num_rows);
2730        assert_eq!(3, num_batches);
2731        assert_eq!(100000000000011, sum_a);
2732    }
2733
2734    #[test]
2735    fn test_decoder_error() {
2736        let schema = Arc::new(Schema::new(vec![Field::new_struct(
2737            "a",
2738            vec![Field::new("child", DataType::Int32, false)],
2739            true,
2740        )]));
2741
2742        let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
2743        let _ = decoder.decode(r#"{"a": { "child":"#.as_bytes()).unwrap();
2744        assert!(decoder.tape_decoder.has_partial_row());
2745        assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2746        let _ = decoder.flush().unwrap_err();
2747        assert!(decoder.tape_decoder.has_partial_row());
2748        assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2749
2750        let parse_err = |s: &str| {
2751            ReaderBuilder::new(schema.clone())
2752                .build(Cursor::new(s.as_bytes()))
2753                .unwrap()
2754                .next()
2755                .unwrap()
2756                .unwrap_err()
2757                .to_string()
2758        };
2759
2760        let err = parse_err(r#"{"a": 123}"#);
2761        assert_eq!(
2762            err,
2763            "Json error: whilst decoding field 'a': expected { got 123"
2764        );
2765
2766        let err = parse_err(r#"{"a": ["bar"]}"#);
2767        assert_eq!(
2768            err,
2769            r#"Json error: whilst decoding field 'a': expected { got ["bar"]"#
2770        );
2771
2772        let err = parse_err(r#"{"a": []}"#);
2773        assert_eq!(
2774            err,
2775            "Json error: whilst decoding field 'a': expected { got []"
2776        );
2777
2778        let err = parse_err(r#"{"a": [{"child": 234}]}"#);
2779        assert_eq!(
2780            err,
2781            r#"Json error: whilst decoding field 'a': expected { got [{"child": 234}]"#
2782        );
2783
2784        let err = parse_err(r#"{"a": [{"child": {"foo": [{"foo": ["bar"]}]}}]}"#);
2785        assert_eq!(
2786            err,
2787            r#"Json error: whilst decoding field 'a': expected { got [{"child": {"foo": [{"foo": ["bar"]}]}}]"#
2788        );
2789
2790        let err = parse_err(r#"{"a": true}"#);
2791        assert_eq!(
2792            err,
2793            "Json error: whilst decoding field 'a': expected { got true"
2794        );
2795
2796        let err = parse_err(r#"{"a": false}"#);
2797        assert_eq!(
2798            err,
2799            "Json error: whilst decoding field 'a': expected { got false"
2800        );
2801
2802        let err = parse_err(r#"{"a": "foo"}"#);
2803        assert_eq!(
2804            err,
2805            "Json error: whilst decoding field 'a': expected { got \"foo\""
2806        );
2807
2808        let err = parse_err(r#"{"a": {"child": false}}"#);
2809        assert_eq!(
2810            err,
2811            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got false"
2812        );
2813
2814        let err = parse_err(r#"{"a": {"child": []}}"#);
2815        assert_eq!(
2816            err,
2817            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got []"
2818        );
2819
2820        let err = parse_err(r#"{"a": {"child": [123]}}"#);
2821        assert_eq!(
2822            err,
2823            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123]"
2824        );
2825
2826        let err = parse_err(r#"{"a": {"child": [123, 3465346]}}"#);
2827        assert_eq!(
2828            err,
2829            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123, 3465346]"
2830        );
2831    }
2832
2833    #[test]
2834    fn test_serialize_timestamp() {
2835        let json = vec![
2836            json!({"timestamp": 1681319393}),
2837            json!({"timestamp": "1970-01-01T00:00:00+02:00"}),
2838        ];
2839        let schema = Schema::new(vec![Field::new(
2840            "timestamp",
2841            DataType::Timestamp(TimeUnit::Second, None),
2842            true,
2843        )]);
2844        let mut decoder = ReaderBuilder::new(Arc::new(schema))
2845            .build_decoder()
2846            .unwrap();
2847        decoder.serialize(&json).unwrap();
2848        let batch = decoder.flush().unwrap().unwrap();
2849        assert_eq!(batch.num_rows(), 2);
2850        assert_eq!(batch.num_columns(), 1);
2851        let values = batch.column(0).as_primitive::<TimestampSecondType>();
2852        assert_eq!(values.values(), &[1681319393, -7200]);
2853    }
2854
2855    #[test]
2856    fn test_serialize_decimal() {
2857        let json = vec![
2858            json!({"decimal": 1.234}),
2859            json!({"decimal": "1.234"}),
2860            json!({"decimal": 1234}),
2861            json!({"decimal": "1234"}),
2862        ];
2863        let schema = Schema::new(vec![Field::new(
2864            "decimal",
2865            DataType::Decimal128(10, 3),
2866            true,
2867        )]);
2868        let mut decoder = ReaderBuilder::new(Arc::new(schema))
2869            .build_decoder()
2870            .unwrap();
2871        decoder.serialize(&json).unwrap();
2872        let batch = decoder.flush().unwrap().unwrap();
2873        assert_eq!(batch.num_rows(), 4);
2874        assert_eq!(batch.num_columns(), 1);
2875        let values = batch.column(0).as_primitive::<Decimal128Type>();
2876        assert_eq!(values.values(), &[1234, 1234, 1234000, 1234000]);
2877    }
2878
2879    #[test]
2880    fn test_serde_field() {
2881        let field = Field::new("int", DataType::Int32, true);
2882        let mut decoder = ReaderBuilder::new_with_field(field)
2883            .build_decoder()
2884            .unwrap();
2885        decoder.serialize(&[1_i32, 2, 3, 4]).unwrap();
2886        let b = decoder.flush().unwrap().unwrap();
2887        let values = b.column(0).as_primitive::<Int32Type>().values();
2888        assert_eq!(values, &[1, 2, 3, 4]);
2889    }
2890
2891    #[test]
2892    fn test_serde_large_numbers() {
2893        let field = Field::new("int", DataType::Int64, true);
2894        let mut decoder = ReaderBuilder::new_with_field(field)
2895            .build_decoder()
2896            .unwrap();
2897
2898        decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2899        let b = decoder.flush().unwrap().unwrap();
2900        let values = b.column(0).as_primitive::<Int64Type>().values();
2901        assert_eq!(values, &[1699148028689, 2, 3, 4]);
2902
2903        let field = Field::new(
2904            "int",
2905            DataType::Timestamp(TimeUnit::Microsecond, None),
2906            true,
2907        );
2908        let mut decoder = ReaderBuilder::new_with_field(field)
2909            .build_decoder()
2910            .unwrap();
2911
2912        decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2913        let b = decoder.flush().unwrap().unwrap();
2914        let values = b
2915            .column(0)
2916            .as_primitive::<TimestampMicrosecondType>()
2917            .values();
2918        assert_eq!(values, &[1699148028689, 2, 3, 4]);
2919    }
2920
2921    #[test]
2922    fn test_coercing_primitive_into_string_decoder() {
2923        let buf = &format!(
2924            r#"[{{"a": 1, "b": "A", "c": "T"}}, {{"a": 2, "b": "BB", "c": "F"}}, {{"a": {}, "b": 123, "c": false}}, {{"a": {}, "b": 789, "c": true}}]"#,
2925            (i32::MAX as i64 + 10),
2926            i64::MAX - 10
2927        );
2928        let schema = Schema::new(vec![
2929            Field::new("a", DataType::Float64, true),
2930            Field::new("b", DataType::Utf8, true),
2931            Field::new("c", DataType::Utf8, true),
2932        ]);
2933        let json_array: Vec<serde_json::Value> = serde_json::from_str(buf).unwrap();
2934        let schema_ref = Arc::new(schema);
2935
2936        // read record batches
2937        let reader = ReaderBuilder::new(schema_ref.clone()).with_coerce_primitive(true);
2938        let mut decoder = reader.build_decoder().unwrap();
2939        decoder.serialize(json_array.as_slice()).unwrap();
2940        let batch = decoder.flush().unwrap().unwrap();
2941        assert_eq!(
2942            batch,
2943            RecordBatch::try_new(
2944                schema_ref,
2945                vec![
2946                    Arc::new(Float64Array::from(vec![
2947                        1.0,
2948                        2.0,
2949                        (i32::MAX as i64 + 10) as f64,
2950                        (i64::MAX - 10) as f64
2951                    ])),
2952                    Arc::new(StringArray::from(vec!["A", "BB", "123", "789"])),
2953                    Arc::new(StringArray::from(vec!["T", "F", "false", "true"])),
2954                ]
2955            )
2956            .unwrap()
2957        );
2958    }
2959
2960    #[test]
2961    fn test_serialize_f32_into_string() {
2962        // Coercing an f32 into a string column must render the value, not its raw bit pattern.
2963        let field = Field::new("f", DataType::Utf8, true);
2964        let mut decoder = ReaderBuilder::new_with_field(field)
2965            .with_coerce_primitive(true)
2966            .build_decoder()
2967            .unwrap();
2968        decoder.serialize(&[1.5_f32, -2.25_f32]).unwrap();
2969        let batch = decoder.flush().unwrap().unwrap();
2970        let values = batch.column(0).as_string::<i32>();
2971        assert_eq!(values.value(0), "1.5");
2972        assert_eq!(values.value(1), "-2.25");
2973    }
2974
2975    // Parse the given `row` in `struct_mode` as a type given by fields.
2976    //
2977    // If as_struct == true, wrap the fields in a Struct field with name "r".
2978    // If as_struct == false, wrap the fields in a Schema.
2979    fn _parse_structs(
2980        row: &str,
2981        struct_mode: StructMode,
2982        fields: Fields,
2983        as_struct: bool,
2984    ) -> Result<RecordBatch, ArrowError> {
2985        let builder = if as_struct {
2986            ReaderBuilder::new_with_field(Field::new("r", DataType::Struct(fields), true))
2987        } else {
2988            ReaderBuilder::new(Arc::new(Schema::new(fields)))
2989        };
2990        builder
2991            .with_struct_mode(struct_mode)
2992            .build(Cursor::new(row.as_bytes()))
2993            .unwrap()
2994            .next()
2995            .unwrap()
2996    }
2997
2998    #[test]
2999    fn test_struct_decoding_list_length() {
3000        use arrow_array::array;
3001
3002        let row = "[1, 2]";
3003
3004        let mut fields = vec![Field::new("a", DataType::Int32, true)];
3005        let too_few_fields = Fields::from(fields.clone());
3006        fields.push(Field::new("b", DataType::Int32, true));
3007        let correct_fields = Fields::from(fields.clone());
3008        fields.push(Field::new("c", DataType::Int32, true));
3009        let too_many_fields = Fields::from(fields.clone());
3010
3011        let parse = |fields: Fields, as_struct: bool| {
3012            _parse_structs(row, StructMode::ListOnly, fields, as_struct)
3013        };
3014
3015        let expected_row = StructArray::new(
3016            correct_fields.clone(),
3017            vec![
3018                Arc::new(array::Int32Array::from(vec![1])),
3019                Arc::new(array::Int32Array::from(vec![2])),
3020            ],
3021            None,
3022        );
3023        let row_field = Field::new("r", DataType::Struct(correct_fields.clone()), true);
3024
3025        assert_eq!(
3026            parse(too_few_fields.clone(), true).unwrap_err().to_string(),
3027            "Json error: found extra columns for 1 fields".to_string()
3028        );
3029        assert_eq!(
3030            parse(too_few_fields, false).unwrap_err().to_string(),
3031            "Json error: found extra columns for 1 fields".to_string()
3032        );
3033        assert_eq!(
3034            parse(correct_fields.clone(), true).unwrap(),
3035            RecordBatch::try_new(
3036                Arc::new(Schema::new(vec![row_field])),
3037                vec![Arc::new(expected_row.clone())]
3038            )
3039            .unwrap()
3040        );
3041        assert_eq!(
3042            parse(correct_fields, false).unwrap(),
3043            RecordBatch::from(expected_row)
3044        );
3045        assert_eq!(
3046            parse(too_many_fields.clone(), true)
3047                .unwrap_err()
3048                .to_string(),
3049            "Json error: found 2 columns for 3 fields".to_string()
3050        );
3051        assert_eq!(
3052            parse(too_many_fields, false).unwrap_err().to_string(),
3053            "Json error: found 2 columns for 3 fields".to_string()
3054        );
3055    }
3056
3057    #[test]
3058    fn test_struct_decoding() {
3059        use arrow_array::builder;
3060
3061        let nested_object_json = r#"{"a": {"b": [1, 2], "c": {"d": 3}}}"#;
3062        let nested_list_json = r#"[[[1, 2], {"d": 3}]]"#;
3063        let nested_mixed_json = r#"{"a": [[1, 2], {"d": 3}]}"#;
3064
3065        let struct_fields = Fields::from(vec![
3066            Field::new("b", DataType::new_list(DataType::Int32, true), true),
3067            Field::new_map(
3068                "c",
3069                Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3070                Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3071                Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Int32, true),
3072                false,
3073                false,
3074            ),
3075        ]);
3076
3077        let list_array =
3078            ListArray::from_iter_primitive::<Int32Type, _, _>(vec![Some(vec![Some(1), Some(2)])]);
3079
3080        let map_array = {
3081            let mut map_builder = builder::MapBuilder::new(
3082                None,
3083                builder::StringBuilder::new(),
3084                builder::Int32Builder::new(),
3085            );
3086            map_builder.keys().append_value("d");
3087            map_builder.values().append_value(3);
3088            map_builder.append(true).unwrap();
3089            map_builder.finish()
3090        };
3091
3092        let struct_array = StructArray::new(
3093            struct_fields.clone(),
3094            vec![Arc::new(list_array), Arc::new(map_array)],
3095            None,
3096        );
3097
3098        let fields = Fields::from(vec![Field::new("a", DataType::Struct(struct_fields), true)]);
3099        let schema = Arc::new(Schema::new(fields.clone()));
3100        let expected = RecordBatch::try_new(schema.clone(), vec![Arc::new(struct_array)]).unwrap();
3101
3102        let parse = |row: &str, struct_mode: StructMode| {
3103            _parse_structs(row, struct_mode, fields.clone(), false)
3104        };
3105
3106        assert_eq!(
3107            parse(nested_object_json, StructMode::ObjectOnly).unwrap(),
3108            expected
3109        );
3110        assert_eq!(
3111            parse(nested_list_json, StructMode::ObjectOnly)
3112                .unwrap_err()
3113                .to_string(),
3114            "Json error: expected { got [[[1, 2], {\"d\": 3}]]".to_owned()
3115        );
3116        assert_eq!(
3117            parse(nested_mixed_json, StructMode::ObjectOnly)
3118                .unwrap_err()
3119                .to_string(),
3120            "Json error: whilst decoding field 'a': expected { got [[1, 2], {\"d\": 3}]".to_owned()
3121        );
3122
3123        assert_eq!(
3124            parse(nested_list_json, StructMode::ListOnly).unwrap(),
3125            expected
3126        );
3127        assert_eq!(
3128            parse(nested_object_json, StructMode::ListOnly)
3129                .unwrap_err()
3130                .to_string(),
3131            "Json error: expected [ got {\"a\": {\"b\": [1, 2]\"c\": {\"d\": 3}}}".to_owned()
3132        );
3133        assert_eq!(
3134            parse(nested_mixed_json, StructMode::ListOnly)
3135                .unwrap_err()
3136                .to_string(),
3137            "Json error: expected [ got {\"a\": [[1, 2], {\"d\": 3}]}".to_owned()
3138        );
3139    }
3140
3141    // Test cases:
3142    // [] -> RecordBatch row with no entries.  Schema = [('a', Int32)] -> Error
3143    // [] -> RecordBatch row with no entries. Schema = [('r', [('a', Int32)])] -> Error
3144    // [] -> StructArray row with no entries. Fields [('a', Int32')] -> Error
3145    // [[]] -> RecordBatch row with empty struct entry. Schema = [('r', [('a', Int32)])] -> Error
3146    #[test]
3147    fn test_struct_decoding_empty_list() {
3148        let int_field = Field::new("a", DataType::Int32, true);
3149        let struct_field = Field::new(
3150            "r",
3151            DataType::Struct(Fields::from(vec![int_field.clone()])),
3152            true,
3153        );
3154
3155        let parse = |row: &str, as_struct: bool, field: Field| {
3156            _parse_structs(
3157                row,
3158                StructMode::ListOnly,
3159                Fields::from(vec![field]),
3160                as_struct,
3161            )
3162        };
3163
3164        // Missing fields
3165        assert_eq!(
3166            parse("[]", true, struct_field.clone())
3167                .unwrap_err()
3168                .to_string(),
3169            "Json error: found 0 columns for 1 fields".to_owned()
3170        );
3171        assert_eq!(
3172            parse("[]", false, int_field.clone())
3173                .unwrap_err()
3174                .to_string(),
3175            "Json error: found 0 columns for 1 fields".to_owned()
3176        );
3177        assert_eq!(
3178            parse("[]", false, struct_field.clone())
3179                .unwrap_err()
3180                .to_string(),
3181            "Json error: found 0 columns for 1 fields".to_owned()
3182        );
3183        assert_eq!(
3184            parse("[[]]", false, struct_field.clone())
3185                .unwrap_err()
3186                .to_string(),
3187            "Json error: whilst decoding field 'r': found 0 columns for 1 fields".to_owned()
3188        );
3189    }
3190
3191    #[test]
3192    fn test_decode_list_struct_with_wrong_types() {
3193        let int_field = Field::new("a", DataType::Int32, true);
3194        let struct_field = Field::new(
3195            "r",
3196            DataType::Struct(Fields::from(vec![int_field.clone()])),
3197            true,
3198        );
3199
3200        let parse = |row: &str, as_struct: bool, field: Field| {
3201            _parse_structs(
3202                row,
3203                StructMode::ListOnly,
3204                Fields::from(vec![field]),
3205                as_struct,
3206            )
3207        };
3208
3209        // Wrong values
3210        assert_eq!(
3211            parse(r#"[["a"]]"#, false, struct_field.clone())
3212                .unwrap_err()
3213                .to_string(),
3214            "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3215        );
3216        assert_eq!(
3217            parse(r#"[["a"]]"#, true, struct_field.clone())
3218                .unwrap_err()
3219                .to_string(),
3220            "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3221        );
3222        assert_eq!(
3223            parse(r#"["a"]"#, true, int_field.clone())
3224                .unwrap_err()
3225                .to_string(),
3226            "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3227        );
3228        assert_eq!(
3229            parse(r#"["a"]"#, false, int_field.clone())
3230                .unwrap_err()
3231                .to_string(),
3232            "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3233        );
3234    }
3235
3236    #[test]
3237    fn test_type_conflict_nulls() {
3238        let schema = Schema::new(vec![
3239            Field::new("null", DataType::Null, true),
3240            Field::new("bool", DataType::Boolean, true),
3241            Field::new("primitive", DataType::Int32, true),
3242            Field::new("numeric", DataType::Decimal128(10, 3), true),
3243            Field::new("string", DataType::Utf8, true),
3244            Field::new("string_view", DataType::Utf8View, true),
3245            Field::new(
3246                "timestamp",
3247                DataType::Timestamp(TimeUnit::Second, None),
3248                true,
3249            ),
3250            Field::new(
3251                "array",
3252                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3253                true,
3254            ),
3255            Field::new(
3256                "map",
3257                DataType::Map(
3258                    Arc::new(Field::new(
3259                        Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3260                        DataType::Struct(Fields::from(vec![
3261                            Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3262                            Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3263                        ])),
3264                        false, // not nullable
3265                    )),
3266                    false, // not sorted
3267                ),
3268                true, // nullable
3269            ),
3270            Field::new(
3271                "struct",
3272                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3273                true,
3274            ),
3275        ]);
3276
3277        // A compatible value for each schema field above, in schema order
3278        let json_values = vec![
3279            json!(null),
3280            json!(true),
3281            json!(42),
3282            json!(1.234),
3283            json!("hi"),
3284            json!("ho"),
3285            json!("1970-01-01T00:00:00+02:00"),
3286            json!([1, "ho", 3]),
3287            json!({"k": "value"}),
3288            json!({"a": 1}),
3289        ];
3290
3291        // Create a set of JSON rows that rotates each value past every field
3292        let json: Vec<_> = (0..json_values.len())
3293            .map(|i| {
3294                let pairs = json_values[i..]
3295                    .iter()
3296                    .chain(json_values[..i].iter())
3297                    .zip(&schema.fields)
3298                    .map(|(v, f)| (f.name().to_string(), v.clone()))
3299                    .collect();
3300                serde_json::Value::Object(pairs)
3301            })
3302            .collect();
3303        let mut decoder = ReaderBuilder::new(Arc::new(schema))
3304            .with_ignore_type_conflicts(true)
3305            .with_coerce_primitive(true)
3306            .build_decoder()
3307            .unwrap();
3308        decoder.serialize(&json).unwrap();
3309        let batch = decoder.flush().unwrap().unwrap();
3310        assert_eq!(batch.num_rows(), 10);
3311        assert_eq!(batch.num_columns(), 10);
3312
3313        // NOTE: NullArray doesn't materialize any values (they're all NULL by definition)
3314        let _ = batch
3315            .column(0)
3316            .as_any()
3317            .downcast_ref::<NullArray>()
3318            .unwrap();
3319
3320        assert!(
3321            batch
3322                .column(1)
3323                .as_any()
3324                .downcast_ref::<BooleanArray>()
3325                .unwrap()
3326                .iter()
3327                .eq([
3328                    Some(true),
3329                    None,
3330                    None,
3331                    None,
3332                    None,
3333                    None,
3334                    None,
3335                    None,
3336                    None,
3337                    None
3338                ])
3339        );
3340
3341        assert!(batch.column(2).as_primitive::<Int32Type>().iter().eq([
3342            Some(42),
3343            Some(1),
3344            None,
3345            None,
3346            None,
3347            None,
3348            None,
3349            None,
3350            None,
3351            None
3352        ]));
3353
3354        assert!(batch.column(3).as_primitive::<Decimal128Type>().iter().eq([
3355            Some(1234),
3356            None,
3357            None,
3358            None,
3359            None,
3360            None,
3361            None,
3362            None,
3363            None,
3364            Some(42000)
3365        ]));
3366
3367        assert!(
3368            batch
3369                .column(4)
3370                .as_any()
3371                .downcast_ref::<StringArray>()
3372                .unwrap()
3373                .iter()
3374                .eq([
3375                    Some("hi"),
3376                    Some("ho"),
3377                    Some("1970-01-01T00:00:00+02:00"),
3378                    None,
3379                    None,
3380                    None,
3381                    None,
3382                    Some("true"),
3383                    Some("42"),
3384                    Some("1.234"),
3385                ])
3386        );
3387
3388        assert!(
3389            batch
3390                .column(5)
3391                .as_any()
3392                .downcast_ref::<StringViewArray>()
3393                .unwrap()
3394                .iter()
3395                .eq([
3396                    Some("ho"),
3397                    Some("1970-01-01T00:00:00+02:00"),
3398                    None,
3399                    None,
3400                    None,
3401                    None,
3402                    Some("true"),
3403                    Some("42"),
3404                    Some("1.234"),
3405                    Some("hi"),
3406                ])
3407        );
3408
3409        assert!(
3410            batch
3411                .column(6)
3412                .as_primitive::<TimestampSecondType>()
3413                .iter()
3414                .eq([
3415                    Some(-7200),
3416                    None,
3417                    None,
3418                    None,
3419                    None,
3420                    None,
3421                    Some(42),
3422                    None,
3423                    None,
3424                    None,
3425                ])
3426        );
3427
3428        let arrays = batch
3429            .column(7)
3430            .as_any()
3431            .downcast_ref::<ListArray>()
3432            .unwrap();
3433        assert_eq!(
3434            arrays.nulls(),
3435            Some(&NullBuffer::from(
3436                &[
3437                    true, false, false, false, false, false, false, false, false, false
3438                ][..]
3439            ))
3440        );
3441        assert_eq!(arrays.offsets()[1], 3);
3442        let array_values = arrays
3443            .values()
3444            .as_any()
3445            .downcast_ref::<Int32Array>()
3446            .unwrap();
3447        assert!(array_values.iter().eq([Some(1), None, Some(3)]));
3448
3449        let maps = batch.column(8).as_any().downcast_ref::<MapArray>().unwrap();
3450        assert_eq!(
3451            maps.nulls(),
3452            Some(&NullBuffer::from(
3453                // Both map and struct can parse
3454                &[
3455                    true, true, false, false, false, false, false, false, false, false
3456                ][..]
3457            ))
3458        );
3459        let map_keys = maps.keys().as_any().downcast_ref::<StringArray>().unwrap();
3460        assert!(map_keys.iter().eq([Some("k"), Some("a")]));
3461        let map_values = maps
3462            .values()
3463            .as_any()
3464            .downcast_ref::<StringArray>()
3465            .unwrap();
3466        assert!(map_values.iter().eq([Some("value"), Some("1")]));
3467
3468        let structs = batch
3469            .column(9)
3470            .as_any()
3471            .downcast_ref::<StructArray>()
3472            .unwrap();
3473        assert_eq!(
3474            structs.nulls(),
3475            Some(&NullBuffer::from(
3476                // Both map and struct can parse
3477                &[
3478                    true, false, false, false, false, false, false, false, false, true
3479                ][..]
3480            ))
3481        );
3482        let struct_fields = structs
3483            .column(0)
3484            .as_any()
3485            .downcast_ref::<Int32Array>()
3486            .unwrap();
3487        assert!(struct_fields.slice(0, 2).iter().eq([Some(1), None]));
3488    }
3489
3490    #[test]
3491    fn test_type_conflict_non_nullable() {
3492        let fields = [
3493            Field::new("bool", DataType::Boolean, false),
3494            Field::new("primitive", DataType::Int32, false),
3495            Field::new("numeric", DataType::Decimal128(10, 3), false),
3496            Field::new("string", DataType::Utf8, false),
3497            Field::new("string_view", DataType::Utf8View, false),
3498            Field::new(
3499                "timestamp",
3500                DataType::Timestamp(TimeUnit::Second, None),
3501                false,
3502            ),
3503            Field::new(
3504                "array",
3505                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3506                false,
3507            ),
3508            Field::new(
3509                "fixed_size_list",
3510                DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3511                false,
3512            ),
3513            Field::new(
3514                "map",
3515                DataType::Map(
3516                    Arc::new(Field::new(
3517                        Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3518                        DataType::Struct(Fields::from(vec![
3519                            Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3520                            Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3521                        ])),
3522                        false, // not nullable
3523                    )),
3524                    false, // not sorted
3525                ),
3526                false, // not nullable
3527            ),
3528            Field::new(
3529                "struct",
3530                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3531                false,
3532            ),
3533        ];
3534
3535        // Every field above will have a type conflict with at least one of these values
3536        let json_values = vec![json!(true), json!({"a": 1})];
3537
3538        for field in fields {
3539            let mut decoder = ReaderBuilder::new_with_field(field)
3540                .with_ignore_type_conflicts(true)
3541                .build_decoder()
3542                .unwrap();
3543            decoder.serialize(&json_values).unwrap();
3544            decoder
3545                .flush()
3546                .expect_err("type conflict on non-nullable type");
3547        }
3548    }
3549
3550    #[test]
3551    fn test_ignore_type_conflicts_disabled() {
3552        let fields = [
3553            Field::new("null", DataType::Null, true),
3554            Field::new("bool", DataType::Boolean, true),
3555            Field::new("primitive", DataType::Int32, true),
3556            Field::new("numeric", DataType::Decimal128(10, 3), true),
3557            Field::new("string", DataType::Utf8, true),
3558            Field::new("string_view", DataType::Utf8View, true),
3559            Field::new(
3560                "timestamp",
3561                DataType::Timestamp(TimeUnit::Second, None),
3562                true,
3563            ),
3564            Field::new(
3565                "array",
3566                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3567                true,
3568            ),
3569            Field::new(
3570                "fixed_size_list",
3571                DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3572                true,
3573            ),
3574            Field::new(
3575                "map",
3576                DataType::Map(
3577                    Arc::new(Field::new(
3578                        Field::MAP_ENTRIES_FIELD_DEFAULT_NAME,
3579                        DataType::Struct(Fields::from(vec![
3580                            Field::new(Field::MAP_KEY_FIELD_DEFAULT_NAME, DataType::Utf8, false),
3581                            Field::new(Field::MAP_VALUE_FIELD_DEFAULT_NAME, DataType::Utf8, true),
3582                        ])),
3583                        false, // not nullable
3584                    )),
3585                    false, // not sorted
3586                ),
3587                true, // not nullable
3588            ),
3589            Field::new(
3590                "struct",
3591                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3592                true,
3593            ),
3594        ];
3595
3596        // Every field above will have a type conflict with at least one of these values
3597        let json_values = vec![json!(true), json!({"a": 1})];
3598
3599        for field in fields {
3600            let mut decoder = ReaderBuilder::new_with_field(field)
3601                .build_decoder()
3602                .unwrap();
3603            decoder.serialize(&json_values).unwrap();
3604            decoder
3605                .flush()
3606                .expect_err("type conflict on non-nullable type");
3607        }
3608    }
3609
3610    #[test]
3611    fn test_read_run_end_encoded() {
3612        let buf = r#"
3613        {"a": "x"}
3614        {"a": "x"}
3615        {"a": "y"}
3616        {"a": "y"}
3617        {"a": "y"}
3618        "#;
3619
3620        let ree_type = DataType::RunEndEncoded(
3621            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3622            Arc::new(Field::new("values", DataType::Utf8, true)),
3623        );
3624        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3625        let batches = do_read(buf, 1024, false, false, schema);
3626        assert_eq!(batches.len(), 1);
3627
3628        let col = batches[0].column(0);
3629        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3630
3631        // 5 logical values compressed into 2 runs
3632        assert_eq!(run_array.len(), 5);
3633        assert_eq!(run_array.run_ends().values(), &[2, 5]);
3634
3635        let values = run_array.values().as_string::<i32>();
3636        assert_eq!(values.len(), 2);
3637        assert_eq!(values.value(0), "x");
3638        assert_eq!(values.value(1), "y");
3639    }
3640
3641    #[test]
3642    fn test_read_run_end_encoded_consecutive_nulls() {
3643        let buf = r#"
3644        {"a": "x"}
3645        {}
3646        {}
3647        {}
3648        {"a": "y"}
3649        "#;
3650
3651        let ree_type = DataType::RunEndEncoded(
3652            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3653            Arc::new(Field::new("values", DataType::Utf8, true)),
3654        );
3655        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3656        let batches = do_read(buf, 1024, false, false, schema);
3657        assert_eq!(batches.len(), 1);
3658
3659        let col = batches[0].column(0);
3660        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3661
3662        // 5 logical values: "x", null, null, null, "y" → 3 runs
3663        assert_eq!(run_array.len(), 5);
3664        assert_eq!(run_array.run_ends().values(), &[1, 4, 5]);
3665
3666        let values = run_array.values().as_string::<i32>();
3667        assert_eq!(values.len(), 3);
3668        assert_eq!(values.value(0), "x");
3669        assert!(values.is_null(1));
3670        assert_eq!(values.value(2), "y");
3671    }
3672
3673    #[test]
3674    fn test_read_run_end_encoded_all_unique() {
3675        let buf = r#"
3676        {"a": 1}
3677        {"a": 2}
3678        {"a": 3}
3679        "#;
3680
3681        let ree_type = DataType::RunEndEncoded(
3682            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3683            Arc::new(Field::new("values", DataType::Int32, true)),
3684        );
3685        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3686        let batches = do_read(buf, 1024, false, false, schema);
3687        assert_eq!(batches.len(), 1);
3688
3689        let col = batches[0].column(0);
3690        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3691
3692        // No compression: 3 unique values → 3 runs
3693        assert_eq!(run_array.len(), 3);
3694        assert_eq!(run_array.run_ends().values(), &[1, 2, 3]);
3695    }
3696
3697    #[test]
3698    fn test_read_run_end_encoded_int16_run_ends() {
3699        let buf = r#"
3700        {"a": "x"}
3701        {"a": "x"}
3702        {"a": "y"}
3703        "#;
3704
3705        let ree_type = DataType::RunEndEncoded(
3706            Arc::new(Field::new("run_ends", DataType::Int16, false)),
3707            Arc::new(Field::new("values", DataType::Utf8, true)),
3708        );
3709        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3710        let batches = do_read(buf, 1024, false, false, schema);
3711        assert_eq!(batches.len(), 1);
3712
3713        let col = batches[0].column(0);
3714        let run_array = col.as_run::<arrow_array::types::Int16Type>();
3715
3716        assert_eq!(run_array.len(), 3);
3717        assert_eq!(run_array.run_ends().values(), &[2i16, 3]);
3718    }
3719
3720    #[test]
3721    fn test_read_nested_run_end_encoded() {
3722        let buf = r#"
3723        {"a": "x"}
3724        {"a": "x"}
3725        {"a": "y"}
3726        "#;
3727
3728        // The outer REE compresses whole rows, while the inner REE compresses the
3729        // repeated string values produced by decoding those rows.
3730        let inner_type = DataType::RunEndEncoded(
3731            Arc::new(Field::new("run_ends", DataType::Int64, false)),
3732            Arc::new(Field::new("values", DataType::Utf8, true)),
3733        );
3734        let outer_type = DataType::RunEndEncoded(
3735            Arc::new(Field::new("run_ends", DataType::Int64, false)),
3736            Arc::new(Field::new("values", inner_type, true)),
3737        );
3738        let schema = Arc::new(Schema::new(vec![Field::new("a", outer_type, true)]));
3739        let batches = do_read(buf, 1024, false, false, schema);
3740        assert_eq!(batches.len(), 1);
3741
3742        let col = batches[0].column(0);
3743        let outer = col.as_run::<arrow_array::types::Int64Type>();
3744        // Three logical rows compress to two outer runs: ["x", "x"] and ["y"].
3745        assert_eq!(outer.len(), 3);
3746        assert_eq!(outer.run_ends().values(), &[2, 3]);
3747
3748        let nested = outer.values().as_run::<arrow_array::types::Int64Type>();
3749        // The physical values of the outer REE are themselves a two-element REE.
3750        assert_eq!(nested.len(), 2);
3751        assert_eq!(nested.run_ends().values(), &[1, 2]);
3752
3753        let nested_values = nested.values().as_string::<i32>();
3754        assert_eq!(nested_values.len(), 2);
3755        assert_eq!(nested_values.value(0), "x");
3756        assert_eq!(nested_values.value(1), "y");
3757    }
3758}