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            "entries",
1291            Field::new("key", DataType::Utf8, false),
1292            Field::new_list("value", Field::new("element", DataType::Utf8, true), true),
1293            false,
1294            true,
1295        );
1296
1297        let schema = Arc::new(Schema::new(vec![map]));
1298
1299        let batches = do_read(buf, 1024, false, false, schema);
1300        assert_eq!(batches.len(), 1);
1301
1302        let map = batches[0].column(0).as_map();
1303        let map_keys = map.keys().as_string::<i32>();
1304        let map_values = map.values().as_list::<i32>();
1305        assert_eq!(map.value_offsets(), &[0, 1, 3, 5]);
1306
1307        let k: Vec<_> = map_keys.iter().flatten().collect();
1308        assert_eq!(&k, &["a", "a", "b", "c", "a"]);
1309
1310        let list_values = map_values.values().as_string::<i32>();
1311        let lv: Vec<_> = list_values.iter().collect();
1312        assert_eq!(&lv, &[Some("foo"), None, None, Some("baz")]);
1313        assert_eq!(map_values.value_offsets(), &[0, 2, 3, 3, 3, 4]);
1314        assert_eq!(map_values.null_count(), 1);
1315        assert!(map_values.is_null(3));
1316
1317        let options = FormatOptions::default().with_null("null");
1318        let formatter = ArrayFormatter::try_new(map, &options).unwrap();
1319        assert_eq!(formatter.value(0).to_string(), "{a: [foo, null]}");
1320        assert_eq!(formatter.value(1).to_string(), "{a: [null], b: []}");
1321        assert_eq!(formatter.value(2).to_string(), "{c: null, a: [baz]}");
1322    }
1323
1324    #[test]
1325    fn test_map_non_nullable_value() {
1326        let map = Field::new_map(
1327            "map",
1328            "entries",
1329            Field::new("keys", DataType::Utf8, false),
1330            Field::new("values", DataType::Utf8, false),
1331            false,
1332            false,
1333        );
1334        let schema = Arc::new(Schema::new(vec![map]));
1335        let buf = r#"{"map": {"key": null}}"#;
1336
1337        let err = ReaderBuilder::new(schema)
1338            .build(Cursor::new(buf.as_bytes()))
1339            .unwrap()
1340            .read()
1341            .unwrap_err();
1342
1343        assert_eq!(
1344            err.to_string(),
1345            "Invalid argument error: Found unmasked nulls for non-nullable StructArray field \"values\""
1346        );
1347    }
1348
1349    #[test]
1350    fn test_not_coercing_primitive_into_string_without_flag() {
1351        let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Utf8, true)]));
1352
1353        let buf = r#"{"a": 1}"#;
1354        let err = ReaderBuilder::new(schema.clone())
1355            .with_batch_size(1024)
1356            .build(Cursor::new(buf.as_bytes()))
1357            .unwrap()
1358            .read()
1359            .unwrap_err();
1360
1361        assert_eq!(
1362            err.to_string(),
1363            "Json error: whilst decoding field 'a': expected string got 1"
1364        );
1365
1366        let buf = r#"{"a": true}"#;
1367        let err = ReaderBuilder::new(schema)
1368            .with_batch_size(1024)
1369            .build(Cursor::new(buf.as_bytes()))
1370            .unwrap()
1371            .read()
1372            .unwrap_err();
1373
1374        assert_eq!(
1375            err.to_string(),
1376            "Json error: whilst decoding field 'a': expected string got true"
1377        );
1378    }
1379
1380    #[test]
1381    fn test_coercing_primitive_into_string() {
1382        let buf = r#"
1383        {"a": 1, "b": 2, "c": true}
1384        {"a": 2E0, "b": 4, "c": false}
1385
1386        {"b": 6, "a": 2.0}
1387        {"b": "5", "a": 2}
1388        {"b": 4e0}
1389        {"b": 7, "a": null}
1390        "#;
1391
1392        let schema = Arc::new(Schema::new(vec![
1393            Field::new("a", DataType::Utf8, true),
1394            Field::new("b", DataType::Utf8, true),
1395            Field::new("c", DataType::Utf8, true),
1396        ]));
1397
1398        let batches = do_read(buf, 1024, true, false, schema);
1399        assert_eq!(batches.len(), 1);
1400
1401        let col1 = batches[0].column(0).as_string::<i32>();
1402        assert_eq!(col1.null_count(), 2);
1403        assert_eq!(col1.value(0), "1");
1404        assert_eq!(col1.value(1), "2E0");
1405        assert_eq!(col1.value(2), "2.0");
1406        assert_eq!(col1.value(3), "2");
1407        assert!(col1.is_null(4));
1408        assert!(col1.is_null(5));
1409
1410        let col2 = batches[0].column(1).as_string::<i32>();
1411        assert_eq!(col2.null_count(), 0);
1412        assert_eq!(col2.value(0), "2");
1413        assert_eq!(col2.value(1), "4");
1414        assert_eq!(col2.value(2), "6");
1415        assert_eq!(col2.value(3), "5");
1416        assert_eq!(col2.value(4), "4e0");
1417        assert_eq!(col2.value(5), "7");
1418
1419        let col3 = batches[0].column(2).as_string::<i32>();
1420        assert_eq!(col3.null_count(), 4);
1421        assert_eq!(col3.value(0), "true");
1422        assert_eq!(col3.value(1), "false");
1423        assert!(col3.is_null(2));
1424        assert!(col3.is_null(3));
1425        assert!(col3.is_null(4));
1426        assert!(col3.is_null(5));
1427    }
1428
1429    fn test_decimal<T: DecimalType>(data_type: DataType) {
1430        let buf = r#"
1431        {"a": 1, "b": 2, "c": 38.30}
1432        {"a": 2, "b": 4, "c": 123.456}
1433
1434        {"b": 1337, "a": "2.0452"}
1435        {"b": "5", "a": "11034.2"}
1436        {"b": 40}
1437        {"b": 1234, "a": null}
1438        "#;
1439
1440        let schema = Arc::new(Schema::new(vec![
1441            Field::new("a", data_type.clone(), true),
1442            Field::new("b", data_type.clone(), true),
1443            Field::new("c", data_type, true),
1444        ]));
1445
1446        let batches = do_read(buf, 1024, true, false, schema);
1447        assert_eq!(batches.len(), 1);
1448
1449        let col1 = batches[0].column(0).as_primitive::<T>();
1450        assert_eq!(col1.null_count(), 2);
1451        assert!(col1.is_null(4));
1452        assert!(col1.is_null(5));
1453        assert_eq!(
1454            col1.values(),
1455            &[100, 200, 204, 1103420, 0, 0].map(T::Native::usize_as)
1456        );
1457
1458        let col2 = batches[0].column(1).as_primitive::<T>();
1459        assert_eq!(col2.null_count(), 0);
1460        assert_eq!(
1461            col2.values(),
1462            &[200, 400, 133700, 500, 4000, 123400].map(T::Native::usize_as)
1463        );
1464
1465        let col3 = batches[0].column(2).as_primitive::<T>();
1466        assert_eq!(col3.null_count(), 4);
1467        assert!(!col3.is_null(0));
1468        assert!(!col3.is_null(1));
1469        assert!(col3.is_null(2));
1470        assert!(col3.is_null(3));
1471        assert!(col3.is_null(4));
1472        assert!(col3.is_null(5));
1473        assert_eq!(
1474            col3.values(),
1475            &[3830, 12345, 0, 0, 0, 0].map(T::Native::usize_as)
1476        );
1477    }
1478
1479    #[test]
1480    fn test_decimals() {
1481        test_decimal::<Decimal32Type>(DataType::Decimal32(8, 2));
1482        test_decimal::<Decimal64Type>(DataType::Decimal64(10, 2));
1483        test_decimal::<Decimal128Type>(DataType::Decimal128(10, 2));
1484        test_decimal::<Decimal256Type>(DataType::Decimal256(10, 2));
1485    }
1486
1487    fn test_timestamp<T: ArrowTimestampType>() {
1488        let buf = r#"
1489        {"a": 1, "b": "2020-09-08T13:42:29.190855+00:00", "c": 38.30, "d": "1997-01-31T09:26:56.123"}
1490        {"a": 2, "b": "2020-09-08T13:42:29.190855Z", "c": 123.456, "d": 123.456}
1491
1492        {"b": 1337, "b": "2020-09-08T13:42:29Z", "c": "1997-01-31T09:26:56.123", "d": "1997-01-31T09:26:56.123Z"}
1493        {"b": 40, "c": "2020-09-08T13:42:29.190855+00:00", "d": "1997-01-31 09:26:56.123-05:00"}
1494        {"b": 1234, "a": null, "c": "1997-01-31 09:26:56.123Z", "d": "1997-01-31 092656"}
1495        {"c": "1997-01-31T14:26:56.123-05:00", "d": "1997-01-31"}
1496        "#;
1497
1498        let with_timezone = DataType::Timestamp(T::UNIT, Some("+08:00".into()));
1499        let schema = Arc::new(Schema::new(vec![
1500            Field::new("a", T::DATA_TYPE, true),
1501            Field::new("b", T::DATA_TYPE, true),
1502            Field::new("c", T::DATA_TYPE, true),
1503            Field::new("d", with_timezone, true),
1504        ]));
1505
1506        let batches = do_read(buf, 1024, true, false, schema);
1507        assert_eq!(batches.len(), 1);
1508
1509        let unit_in_nanos: i64 = match T::UNIT {
1510            TimeUnit::Second => 1_000_000_000,
1511            TimeUnit::Millisecond => 1_000_000,
1512            TimeUnit::Microsecond => 1_000,
1513            TimeUnit::Nanosecond => 1,
1514        };
1515
1516        let col1 = batches[0].column(0).as_primitive::<T>();
1517        assert_eq!(col1.null_count(), 4);
1518        assert!(col1.is_null(2));
1519        assert!(col1.is_null(3));
1520        assert!(col1.is_null(4));
1521        assert!(col1.is_null(5));
1522        assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1523
1524        let col2 = batches[0].column(1).as_primitive::<T>();
1525        assert_eq!(col2.null_count(), 1);
1526        assert!(col2.is_null(5));
1527        assert_eq!(
1528            col2.values(),
1529            &[
1530                1599572549190855000 / unit_in_nanos,
1531                1599572549190855000 / unit_in_nanos,
1532                1599572549000000000 / unit_in_nanos,
1533                40,
1534                1234,
1535                0
1536            ]
1537        );
1538
1539        let col3 = batches[0].column(2).as_primitive::<T>();
1540        assert_eq!(col3.null_count(), 0);
1541        assert_eq!(
1542            col3.values(),
1543            &[
1544                38,
1545                123,
1546                854702816123000000 / unit_in_nanos,
1547                1599572549190855000 / unit_in_nanos,
1548                854702816123000000 / unit_in_nanos,
1549                854738816123000000 / unit_in_nanos
1550            ]
1551        );
1552
1553        let col4 = batches[0].column(3).as_primitive::<T>();
1554
1555        assert_eq!(col4.null_count(), 0);
1556        assert_eq!(
1557            col4.values(),
1558            &[
1559                854674016123000000 / unit_in_nanos,
1560                123,
1561                854702816123000000 / unit_in_nanos,
1562                854720816123000000 / unit_in_nanos,
1563                854674016000000000 / unit_in_nanos,
1564                854640000000000000 / unit_in_nanos
1565            ]
1566        );
1567    }
1568
1569    #[test]
1570    fn test_timestamps() {
1571        test_timestamp::<TimestampSecondType>();
1572        test_timestamp::<TimestampMillisecondType>();
1573        test_timestamp::<TimestampMicrosecondType>();
1574        test_timestamp::<TimestampNanosecondType>();
1575    }
1576
1577    fn test_time<T: ArrowTemporalType>() {
1578        let buf = r#"
1579        {"a": 1, "b": "09:26:56.123 AM", "c": 38.30}
1580        {"a": 2, "b": "23:59:59", "c": 123.456}
1581
1582        {"b": 1337, "b": "6:00 pm", "c": "09:26:56.123"}
1583        {"b": 40, "c": "13:42:29.190855"}
1584        {"b": 1234, "a": null, "c": "09:26:56.123"}
1585        {"c": "14:26:56.123"}
1586        "#;
1587
1588        let unit = match T::DATA_TYPE {
1589            DataType::Time32(unit) | DataType::Time64(unit) => unit,
1590            _ => unreachable!(),
1591        };
1592
1593        let unit_in_nanos = match unit {
1594            TimeUnit::Second => 1_000_000_000,
1595            TimeUnit::Millisecond => 1_000_000,
1596            TimeUnit::Microsecond => 1_000,
1597            TimeUnit::Nanosecond => 1,
1598        };
1599
1600        let schema = Arc::new(Schema::new(vec![
1601            Field::new("a", T::DATA_TYPE, true),
1602            Field::new("b", T::DATA_TYPE, true),
1603            Field::new("c", T::DATA_TYPE, true),
1604        ]));
1605
1606        let batches = do_read(buf, 1024, true, false, schema);
1607        assert_eq!(batches.len(), 1);
1608
1609        let col1 = batches[0].column(0).as_primitive::<T>();
1610        assert_eq!(col1.null_count(), 4);
1611        assert!(col1.is_null(2));
1612        assert!(col1.is_null(3));
1613        assert!(col1.is_null(4));
1614        assert!(col1.is_null(5));
1615        assert_eq!(col1.values(), &[1, 2, 0, 0, 0, 0].map(T::Native::usize_as));
1616
1617        let col2 = batches[0].column(1).as_primitive::<T>();
1618        assert_eq!(col2.null_count(), 1);
1619        assert!(col2.is_null(5));
1620        assert_eq!(
1621            col2.values(),
1622            &[
1623                34016123000000 / unit_in_nanos,
1624                86399000000000 / unit_in_nanos,
1625                64800000000000 / unit_in_nanos,
1626                40,
1627                1234,
1628                0
1629            ]
1630            .map(T::Native::usize_as)
1631        );
1632
1633        let col3 = batches[0].column(2).as_primitive::<T>();
1634        assert_eq!(col3.null_count(), 0);
1635        assert_eq!(
1636            col3.values(),
1637            &[
1638                38,
1639                123,
1640                34016123000000 / unit_in_nanos,
1641                49349190855000 / unit_in_nanos,
1642                34016123000000 / unit_in_nanos,
1643                52016123000000 / unit_in_nanos
1644            ]
1645            .map(T::Native::usize_as)
1646        );
1647    }
1648
1649    #[test]
1650    fn test_times() {
1651        test_time::<Time32MillisecondType>();
1652        test_time::<Time32SecondType>();
1653        test_time::<Time64MicrosecondType>();
1654        test_time::<Time64NanosecondType>();
1655    }
1656
1657    fn test_duration<T: ArrowTemporalType>() {
1658        let buf = r#"
1659        {"a": 1, "b": "2"}
1660        {"a": 3, "b": null}
1661        "#;
1662
1663        let schema = Arc::new(Schema::new(vec![
1664            Field::new("a", T::DATA_TYPE, true),
1665            Field::new("b", T::DATA_TYPE, true),
1666        ]));
1667
1668        let batches = do_read(buf, 1024, true, false, schema);
1669        assert_eq!(batches.len(), 1);
1670
1671        let col_a = batches[0].column_by_name("a").unwrap().as_primitive::<T>();
1672        assert_eq!(col_a.null_count(), 0);
1673        assert_eq!(col_a.values(), &[1, 3].map(T::Native::usize_as));
1674
1675        let col2 = batches[0].column_by_name("b").unwrap().as_primitive::<T>();
1676        assert_eq!(col2.null_count(), 1);
1677        assert_eq!(col2.values(), &[2, 0].map(T::Native::usize_as));
1678    }
1679
1680    #[test]
1681    fn test_durations() {
1682        test_duration::<DurationNanosecondType>();
1683        test_duration::<DurationMicrosecondType>();
1684        test_duration::<DurationMillisecondType>();
1685        test_duration::<DurationSecondType>();
1686    }
1687
1688    #[test]
1689    fn test_delta_checkpoint() {
1690        let json = "{\"protocol\":{\"minReaderVersion\":1,\"minWriterVersion\":2}}";
1691        let schema = Arc::new(Schema::new(vec![
1692            Field::new_struct(
1693                "protocol",
1694                vec![
1695                    Field::new("minReaderVersion", DataType::Int32, true),
1696                    Field::new("minWriterVersion", DataType::Int32, true),
1697                ],
1698                true,
1699            ),
1700            Field::new_struct(
1701                "add",
1702                vec![Field::new_map(
1703                    "partitionValues",
1704                    "key_value",
1705                    Field::new("key", DataType::Utf8, false),
1706                    Field::new("value", DataType::Utf8, true),
1707                    false,
1708                    false,
1709                )],
1710                true,
1711            ),
1712        ]));
1713
1714        let batches = do_read(json, 1024, true, false, schema);
1715        assert_eq!(batches.len(), 1);
1716
1717        let s: StructArray = batches.into_iter().next().unwrap().into();
1718        let opts = FormatOptions::default().with_null("null");
1719        let formatter = ArrayFormatter::try_new(&s, &opts).unwrap();
1720        assert_eq!(
1721            formatter.value(0).to_string(),
1722            "{protocol: {minReaderVersion: 1, minWriterVersion: 2}, add: null}"
1723        );
1724    }
1725
1726    #[test]
1727    fn struct_nullability() {
1728        let do_test = |child: DataType| {
1729            // Test correctly enforced nullability
1730            let non_null = r#"{"foo": {}}"#;
1731            let schema = Arc::new(Schema::new(vec![Field::new_struct(
1732                "foo",
1733                vec![Field::new("bar", child, false)],
1734                true,
1735            )]));
1736            let mut reader = ReaderBuilder::new(schema.clone())
1737                .build(Cursor::new(non_null.as_bytes()))
1738                .unwrap();
1739            assert!(reader.next().unwrap().is_err()); // Should error as not nullable
1740
1741            let null = r#"{"foo": {bar: null}}"#;
1742            let mut reader = ReaderBuilder::new(schema.clone())
1743                .build(Cursor::new(null.as_bytes()))
1744                .unwrap();
1745            assert!(reader.next().unwrap().is_err()); // Should error as not nullable
1746
1747            // Test nulls in nullable parent can mask nulls in non-nullable child
1748            let null = r#"{"foo": null}"#;
1749            let mut reader = ReaderBuilder::new(schema)
1750                .build(Cursor::new(null.as_bytes()))
1751                .unwrap();
1752            let batch = reader.next().unwrap().unwrap();
1753            assert_eq!(batch.num_columns(), 1);
1754            let foo = batch.column(0).as_struct();
1755            assert_eq!(foo.len(), 1);
1756            assert!(foo.is_null(0));
1757            assert_eq!(foo.num_columns(), 1);
1758
1759            let bar = foo.column(0);
1760            assert_eq!(bar.len(), 1);
1761            // Non-nullable child can still contain null as masked by parent
1762            assert!(bar.is_null(0));
1763        };
1764
1765        do_test(DataType::Boolean);
1766        do_test(DataType::Int32);
1767        do_test(DataType::Utf8);
1768        do_test(DataType::Decimal128(2, 1));
1769        do_test(DataType::Timestamp(
1770            TimeUnit::Microsecond,
1771            Some("+00:00".into()),
1772        ));
1773    }
1774
1775    #[test]
1776    fn test_truncation() {
1777        let buf = r#"
1778        {"i64": 9223372036854775807, "u64": 18446744073709551615 }
1779        {"i64": "9223372036854775807", "u64": "18446744073709551615" }
1780        {"i64": -9223372036854775808, "u64": 0 }
1781        {"i64": "-9223372036854775808", "u64": 0 }
1782        "#;
1783
1784        let schema = Arc::new(Schema::new(vec![
1785            Field::new("i64", DataType::Int64, true),
1786            Field::new("u64", DataType::UInt64, true),
1787        ]));
1788
1789        let batches = do_read(buf, 1024, true, false, schema);
1790        assert_eq!(batches.len(), 1);
1791
1792        let i64 = batches[0].column(0).as_primitive::<Int64Type>();
1793        assert_eq!(i64.values(), &[i64::MAX, i64::MAX, i64::MIN, i64::MIN]);
1794
1795        let u64 = batches[0].column(1).as_primitive::<UInt64Type>();
1796        assert_eq!(u64.values(), &[u64::MAX, u64::MAX, u64::MIN, u64::MIN]);
1797    }
1798
1799    #[test]
1800    fn test_timestamp_truncation() {
1801        let buf = r#"
1802        {"time": 9223372036854775807 }
1803        {"time": -9223372036854775808 }
1804        {"time": 9e5 }
1805        "#;
1806
1807        let schema = Arc::new(Schema::new(vec![Field::new(
1808            "time",
1809            DataType::Timestamp(TimeUnit::Nanosecond, None),
1810            true,
1811        )]));
1812
1813        let batches = do_read(buf, 1024, true, false, schema);
1814        assert_eq!(batches.len(), 1);
1815
1816        let i64 = batches[0]
1817            .column(0)
1818            .as_primitive::<TimestampNanosecondType>();
1819        assert_eq!(i64.values(), &[i64::MAX, i64::MIN, 900000]);
1820    }
1821
1822    #[test]
1823    fn test_strict_mode_no_missing_columns_in_schema() {
1824        let buf = r#"
1825        {"a": 1, "b": "2", "c": true}
1826        {"a": 2E0, "b": "4", "c": false}
1827        "#;
1828
1829        let schema = Arc::new(Schema::new(vec![
1830            Field::new("a", DataType::Int16, false),
1831            Field::new("b", DataType::Utf8, false),
1832            Field::new("c", DataType::Boolean, false),
1833        ]));
1834
1835        let batches = do_read(buf, 1024, true, true, schema);
1836        assert_eq!(batches.len(), 1);
1837
1838        let buf = r#"
1839        {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1840        {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1841        "#;
1842
1843        let schema = Arc::new(Schema::new(vec![
1844            Field::new("a", DataType::Int16, false),
1845            Field::new("b", DataType::Utf8, false),
1846            Field::new_struct(
1847                "c",
1848                vec![
1849                    Field::new("a", DataType::Boolean, false),
1850                    Field::new("b", DataType::Int16, false),
1851                ],
1852                false,
1853            ),
1854        ]));
1855
1856        let batches = do_read(buf, 1024, true, true, schema);
1857        assert_eq!(batches.len(), 1);
1858    }
1859
1860    #[test]
1861    fn test_strict_mode_missing_columns_in_schema() {
1862        let buf = r#"
1863        {"a": 1, "b": "2", "c": true}
1864        {"a": 2E0, "b": "4", "c": false}
1865        "#;
1866
1867        let schema = Arc::new(Schema::new(vec![
1868            Field::new("a", DataType::Int16, true),
1869            Field::new("c", DataType::Boolean, true),
1870        ]));
1871
1872        let err = ReaderBuilder::new(schema)
1873            .with_batch_size(1024)
1874            .with_strict_mode(true)
1875            .build(Cursor::new(buf.as_bytes()))
1876            .unwrap()
1877            .read()
1878            .unwrap_err();
1879
1880        assert_eq!(
1881            err.to_string(),
1882            "Json error: column 'b' missing from schema"
1883        );
1884
1885        let buf = r#"
1886        {"a": 1, "b": "2", "c": {"a": true, "b": 1}}
1887        {"a": 2E0, "b": "4", "c": {"a": false, "b": 2}}
1888        "#;
1889
1890        let schema = Arc::new(Schema::new(vec![
1891            Field::new("a", DataType::Int16, false),
1892            Field::new("b", DataType::Utf8, false),
1893            Field::new_struct("c", vec![Field::new("a", DataType::Boolean, false)], false),
1894        ]));
1895
1896        let err = ReaderBuilder::new(schema)
1897            .with_batch_size(1024)
1898            .with_strict_mode(true)
1899            .build(Cursor::new(buf.as_bytes()))
1900            .unwrap()
1901            .read()
1902            .unwrap_err();
1903
1904        assert_eq!(
1905            err.to_string(),
1906            "Json error: whilst decoding field 'c': column 'b' missing from schema"
1907        );
1908    }
1909
1910    fn read_file(path: &str, schema: Option<Schema>) -> Reader<BufReader<File>> {
1911        let file = File::open(path).unwrap();
1912        let mut reader = BufReader::new(file);
1913        let schema = schema.unwrap_or_else(|| {
1914            let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
1915            reader.rewind().unwrap();
1916            schema
1917        });
1918        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
1919        builder.build(reader).unwrap()
1920    }
1921
1922    #[test]
1923    fn test_json_basic() {
1924        let mut reader = read_file("test/data/basic.json", None);
1925        let batch = reader.next().unwrap().unwrap();
1926
1927        assert_eq!(8, batch.num_columns());
1928        assert_eq!(12, batch.num_rows());
1929
1930        let schema = reader.schema();
1931        let batch_schema = batch.schema();
1932        assert_eq!(schema, batch_schema);
1933
1934        let a = schema.column_with_name("a").unwrap();
1935        assert_eq!(0, a.0);
1936        assert_eq!(&DataType::Int64, a.1.data_type());
1937        let b = schema.column_with_name("b").unwrap();
1938        assert_eq!(1, b.0);
1939        assert_eq!(&DataType::Float64, b.1.data_type());
1940        let c = schema.column_with_name("c").unwrap();
1941        assert_eq!(2, c.0);
1942        assert_eq!(&DataType::Boolean, c.1.data_type());
1943        let d = schema.column_with_name("d").unwrap();
1944        assert_eq!(3, d.0);
1945        assert_eq!(&DataType::Utf8, d.1.data_type());
1946
1947        let aa = batch.column(a.0).as_primitive::<Int64Type>();
1948        assert_eq!(1, aa.value(0));
1949        assert_eq!(-10, aa.value(1));
1950        let bb = batch.column(b.0).as_primitive::<Float64Type>();
1951        assert_eq!(2.0, bb.value(0));
1952        assert_eq!(-3.5, bb.value(1));
1953        let cc = batch.column(c.0).as_boolean();
1954        assert!(!cc.value(0));
1955        assert!(cc.value(10));
1956        let dd = batch.column(d.0).as_string::<i32>();
1957        assert_eq!("4", dd.value(0));
1958        assert_eq!("text", dd.value(8));
1959    }
1960
1961    #[test]
1962    fn test_json_empty_projection() {
1963        let mut reader = read_file("test/data/basic.json", Some(Schema::empty()));
1964        let batch = reader.next().unwrap().unwrap();
1965
1966        assert_eq!(0, batch.num_columns());
1967        assert_eq!(12, batch.num_rows());
1968    }
1969
1970    #[test]
1971    fn test_json_basic_with_nulls() {
1972        let mut reader = read_file("test/data/basic_nulls.json", None);
1973        let batch = reader.next().unwrap().unwrap();
1974
1975        assert_eq!(4, batch.num_columns());
1976        assert_eq!(12, batch.num_rows());
1977
1978        let schema = reader.schema();
1979        let batch_schema = batch.schema();
1980        assert_eq!(schema, batch_schema);
1981
1982        let a = schema.column_with_name("a").unwrap();
1983        assert_eq!(&DataType::Int64, a.1.data_type());
1984        let b = schema.column_with_name("b").unwrap();
1985        assert_eq!(&DataType::Float64, b.1.data_type());
1986        let c = schema.column_with_name("c").unwrap();
1987        assert_eq!(&DataType::Boolean, c.1.data_type());
1988        let d = schema.column_with_name("d").unwrap();
1989        assert_eq!(&DataType::Utf8, d.1.data_type());
1990
1991        let aa = batch.column(a.0).as_primitive::<Int64Type>();
1992        assert!(aa.is_valid(0));
1993        assert!(!aa.is_valid(1));
1994        assert!(!aa.is_valid(11));
1995        let bb = batch.column(b.0).as_primitive::<Float64Type>();
1996        assert!(bb.is_valid(0));
1997        assert!(!bb.is_valid(2));
1998        assert!(!bb.is_valid(11));
1999        let cc = batch.column(c.0).as_boolean();
2000        assert!(cc.is_valid(0));
2001        assert!(!cc.is_valid(4));
2002        assert!(!cc.is_valid(11));
2003        let dd = batch.column(d.0).as_string::<i32>();
2004        assert!(!dd.is_valid(0));
2005        assert!(dd.is_valid(1));
2006        assert!(!dd.is_valid(4));
2007        assert!(!dd.is_valid(11));
2008    }
2009
2010    #[test]
2011    fn test_json_basic_schema() {
2012        let schema = Schema::new(vec![
2013            Field::new("a", DataType::Int64, true),
2014            Field::new("b", DataType::Float32, false),
2015            Field::new("c", DataType::Boolean, false),
2016            Field::new("d", DataType::Utf8, false),
2017        ]);
2018
2019        let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2020        let reader_schema = reader.schema();
2021        assert_eq!(reader_schema.as_ref(), &schema);
2022        let batch = reader.next().unwrap().unwrap();
2023
2024        assert_eq!(4, batch.num_columns());
2025        assert_eq!(12, batch.num_rows());
2026
2027        let schema = batch.schema();
2028
2029        let a = schema.column_with_name("a").unwrap();
2030        assert_eq!(&DataType::Int64, a.1.data_type());
2031        let b = schema.column_with_name("b").unwrap();
2032        assert_eq!(&DataType::Float32, b.1.data_type());
2033        let c = schema.column_with_name("c").unwrap();
2034        assert_eq!(&DataType::Boolean, c.1.data_type());
2035        let d = schema.column_with_name("d").unwrap();
2036        assert_eq!(&DataType::Utf8, d.1.data_type());
2037
2038        let aa = batch.column(a.0).as_primitive::<Int64Type>();
2039        assert_eq!(1, aa.value(0));
2040        assert_eq!(100000000000000, aa.value(11));
2041        let bb = batch.column(b.0).as_primitive::<Float32Type>();
2042        assert_eq!(2.0, bb.value(0));
2043        assert_eq!(-3.5, bb.value(1));
2044    }
2045
2046    #[test]
2047    fn test_json_basic_schema_projection() {
2048        let schema = Schema::new(vec![
2049            Field::new("a", DataType::Int64, true),
2050            Field::new("c", DataType::Boolean, false),
2051        ]);
2052
2053        let mut reader = read_file("test/data/basic.json", Some(schema.clone()));
2054        let batch = reader.next().unwrap().unwrap();
2055
2056        assert_eq!(2, batch.num_columns());
2057        assert_eq!(2, batch.schema().fields().len());
2058        assert_eq!(12, batch.num_rows());
2059
2060        assert_eq!(batch.schema().as_ref(), &schema);
2061
2062        let a = schema.column_with_name("a").unwrap();
2063        assert_eq!(0, a.0);
2064        assert_eq!(&DataType::Int64, a.1.data_type());
2065        let c = schema.column_with_name("c").unwrap();
2066        assert_eq!(1, c.0);
2067        assert_eq!(&DataType::Boolean, c.1.data_type());
2068    }
2069
2070    #[test]
2071    fn test_json_arrays() {
2072        let mut reader = read_file("test/data/arrays.json", None);
2073        let batch = reader.next().unwrap().unwrap();
2074
2075        assert_eq!(4, batch.num_columns());
2076        assert_eq!(3, batch.num_rows());
2077
2078        let schema = batch.schema();
2079
2080        let a = schema.column_with_name("a").unwrap();
2081        assert_eq!(&DataType::Int64, a.1.data_type());
2082        let b = schema.column_with_name("b").unwrap();
2083        assert_eq!(
2084            &DataType::List(Arc::new(Field::new_list_field(DataType::Float64, true))),
2085            b.1.data_type()
2086        );
2087        let c = schema.column_with_name("c").unwrap();
2088        assert_eq!(
2089            &DataType::List(Arc::new(Field::new_list_field(DataType::Boolean, true))),
2090            c.1.data_type()
2091        );
2092        let d = schema.column_with_name("d").unwrap();
2093        assert_eq!(&DataType::Utf8, d.1.data_type());
2094
2095        let aa = batch.column(a.0).as_primitive::<Int64Type>();
2096        assert_eq!(1, aa.value(0));
2097        assert_eq!(-10, aa.value(1));
2098        assert_eq!(1627668684594000000, aa.value(2));
2099        let bb = batch.column(b.0).as_list::<i32>();
2100        let bb = bb.values().as_primitive::<Float64Type>();
2101        assert_eq!(9, bb.len());
2102        assert_eq!(2.0, bb.value(0));
2103        assert_eq!(-6.1, bb.value(5));
2104        assert!(!bb.is_valid(7));
2105
2106        let cc = batch
2107            .column(c.0)
2108            .as_any()
2109            .downcast_ref::<ListArray>()
2110            .unwrap();
2111        let cc = cc.values().as_boolean();
2112        assert_eq!(6, cc.len());
2113        assert!(!cc.value(0));
2114        assert!(!cc.value(4));
2115        assert!(!cc.is_valid(5));
2116    }
2117
2118    #[test]
2119    fn test_empty_json_arrays() {
2120        let json_content = r#"
2121            {"items": []}
2122            {"items": null}
2123            {}
2124            "#;
2125
2126        let schema = Arc::new(Schema::new(vec![Field::new(
2127            "items",
2128            DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2129            true,
2130        )]));
2131
2132        let batches = do_read(json_content, 1024, false, false, schema);
2133        assert_eq!(batches.len(), 1);
2134
2135        let col1 = batches[0].column(0).as_list::<i32>();
2136        assert_eq!(col1.null_count(), 2);
2137        assert!(col1.value(0).is_empty());
2138        assert_eq!(col1.value(0).data_type(), &DataType::Null);
2139        assert!(col1.is_null(1));
2140        assert!(col1.is_null(2));
2141    }
2142
2143    #[test]
2144    fn test_nested_empty_json_arrays() {
2145        let json_content = r#"
2146            {"items": [[],[]]}
2147            {"items": [[null, null],[null]]}
2148            "#;
2149
2150        let schema = Arc::new(Schema::new(vec![Field::new(
2151            "items",
2152            DataType::List(FieldRef::new(Field::new_list_field(
2153                DataType::List(FieldRef::new(Field::new_list_field(DataType::Null, true))),
2154                true,
2155            ))),
2156            true,
2157        )]));
2158
2159        let batches = do_read(json_content, 1024, false, false, schema);
2160        assert_eq!(batches.len(), 1);
2161
2162        let col1 = batches[0].column(0).as_list::<i32>();
2163        assert_eq!(col1.null_count(), 0);
2164        assert_eq!(col1.value(0).len(), 2);
2165        assert!(col1.value(0).as_list::<i32>().value(0).is_empty());
2166        assert!(col1.value(0).as_list::<i32>().value(1).is_empty());
2167
2168        assert_eq!(col1.value(1).len(), 2);
2169        assert_eq!(col1.value(1).as_list::<i32>().value(0).len(), 2);
2170        assert_eq!(col1.value(1).as_list::<i32>().value(1).len(), 1);
2171    }
2172
2173    #[test]
2174    fn test_nested_list_json_arrays() {
2175        let c_field = Field::new_struct("c", vec![Field::new("d", DataType::Utf8, true)], true);
2176        let a_struct_field = Field::new_struct(
2177            "a",
2178            vec![Field::new("b", DataType::Boolean, true), c_field.clone()],
2179            true,
2180        );
2181        let a_field = Field::new("a", DataType::List(Arc::new(a_struct_field.clone())), true);
2182        let schema = Arc::new(Schema::new(vec![a_field.clone()]));
2183        let builder = ReaderBuilder::new(schema).with_batch_size(64);
2184        let json_content = r#"
2185        {"a": [{"b": true, "c": {"d": "a_text"}}, {"b": false, "c": {"d": "b_text"}}]}
2186        {"a": [{"b": false, "c": null}]}
2187        {"a": [{"b": true, "c": {"d": "c_text"}}, {"b": null, "c": {"d": "d_text"}}, {"b": true, "c": {"d": null}}]}
2188        {"a": null}
2189        {"a": []}
2190        {"a": [null]}
2191        "#;
2192        let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2193
2194        // build expected output
2195        let d = StringArray::from(vec![
2196            Some("a_text"),
2197            Some("b_text"),
2198            None,
2199            Some("c_text"),
2200            Some("d_text"),
2201            None,
2202            None,
2203        ]);
2204        let c = StructArray::new(
2205            vec![Field::new("d", DataType::Utf8, true)].into(),
2206            vec![Arc::new(d.clone()) as ArrayRef],
2207            Some(NullBuffer::from(vec![
2208                true, true, false, true, true, true, false,
2209            ])),
2210        );
2211        let b = BooleanArray::from(vec![
2212            Some(true),
2213            Some(false),
2214            Some(false),
2215            Some(true),
2216            None,
2217            Some(true),
2218            None,
2219        ]);
2220        let a = StructArray::new(
2221            vec![Field::new("b", DataType::Boolean, true), c_field.clone()].into(),
2222            vec![
2223                Arc::new(b.clone()) as ArrayRef,
2224                Arc::new(c.clone()) as ArrayRef,
2225            ],
2226            Some(NullBuffer::from(vec![
2227                true, true, true, true, true, true, false,
2228            ])),
2229        );
2230        let a_list = ListArray::new(
2231            Arc::new(a_struct_field.clone()),
2232            OffsetBuffer::new(ScalarBuffer::from(vec![0i32, 2, 3, 6, 6, 6, 7])),
2233            Arc::new(a),
2234            Some(NullBuffer::from(vec![true, true, true, false, true, true])),
2235        );
2236
2237        // compare `a` with result from json reader
2238        let batch = reader.next().unwrap().unwrap();
2239        let read = batch.column(0);
2240        assert_eq!(read.len(), 6);
2241        // compare the arrays the long way around, to better detect differences
2242        let read: &ListArray = read.as_list::<i32>();
2243        let expected = &a_list;
2244        assert_eq!(read.value_offsets(), &[0, 2, 3, 6, 6, 6, 7]);
2245        // compare list null buffers
2246        assert_eq!(read.nulls(), expected.nulls());
2247        // build struct from list
2248        let struct_array = read.values().as_struct();
2249        let expected_struct_array = expected.values().as_struct();
2250
2251        assert_eq!(7, struct_array.len());
2252        assert_eq!(1, struct_array.null_count());
2253        assert_eq!(7, expected_struct_array.len());
2254        assert_eq!(1, expected_struct_array.null_count());
2255        // test struct's nulls
2256        assert_eq!(struct_array.nulls(), expected_struct_array.nulls());
2257        // test struct's fields
2258        let read_b = struct_array.column(0);
2259        assert_eq!(read_b.as_ref(), &b);
2260        let read_c = struct_array.column(1);
2261        assert_eq!(read_c.as_struct(), &c);
2262        let read_c = read_c.as_struct();
2263        let read_d = read_c.column(0);
2264        assert_eq!(read_d.as_ref(), &d);
2265
2266        assert_eq!(read, expected);
2267    }
2268
2269    fn assert_read_list_view<O: OffsetSizeTrait>() {
2270        let field = Arc::new(Field::new("item", DataType::Int32, true));
2271        let data_type = GenericListViewArray::<O>::DATA_TYPE_CONSTRUCTOR(field.clone());
2272        let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2273
2274        let buf = r#"
2275        {"lv": [1, 2, 3]}
2276        {"lv": [4, null]}
2277        {"lv": null}
2278        {"lv": [6]}
2279        {"lv": []}
2280        "#;
2281
2282        let batches = do_read(buf, 1024, false, false, schema);
2283        assert_eq!(batches.len(), 1);
2284        let batch = &batches[0];
2285        let col = batch.column(0);
2286        let list_view = col
2287            .as_any()
2288            .downcast_ref::<GenericListViewArray<O>>()
2289            .unwrap();
2290
2291        assert_eq!(list_view.len(), 5);
2292
2293        // Check offsets and sizes
2294        let expected_offsets: Vec<O> = vec![0, 3, 5, 5, 6]
2295            .into_iter()
2296            .map(|v| O::usize_as(v))
2297            .collect();
2298        let expected_sizes: Vec<O> = vec![3, 2, 0, 1, 0]
2299            .into_iter()
2300            .map(|v| O::usize_as(v))
2301            .collect();
2302        assert_eq!(list_view.value_offsets(), &expected_offsets);
2303        assert_eq!(list_view.value_sizes(), &expected_sizes);
2304
2305        // Row 0: [1, 2, 3]
2306        assert!(list_view.is_valid(0));
2307        let vals = list_view.value(0);
2308        let ints = vals.as_primitive::<Int32Type>();
2309        assert_eq!(ints.values(), &[1, 2, 3]);
2310
2311        // Row 1: [4, null]
2312        assert!(list_view.is_valid(1));
2313        let vals = list_view.value(1);
2314        let ints = vals.as_primitive::<Int32Type>();
2315        assert_eq!(ints.len(), 2);
2316        assert_eq!(ints.value(0), 4);
2317        assert!(ints.is_null(1));
2318
2319        // Row 2: null
2320        assert!(list_view.is_null(2));
2321
2322        // Row 3: [6]
2323        assert!(list_view.is_valid(3));
2324        let vals = list_view.value(3);
2325        let ints = vals.as_primitive::<Int32Type>();
2326        assert_eq!(ints.values(), &[6]);
2327
2328        // Row 4: []
2329        assert!(list_view.is_valid(4));
2330        let vals = list_view.value(4);
2331        assert_eq!(vals.len(), 0);
2332    }
2333
2334    #[test]
2335    fn test_read_list_view() {
2336        assert_read_list_view::<i32>();
2337        assert_read_list_view::<i64>();
2338    }
2339
2340    #[test]
2341    fn test_read_list_view_rejects_null_non_nullable_child() {
2342        let field = Arc::new(Field::new("item", DataType::Int32, false));
2343        for (data_type, array_type) in [
2344            (DataType::ListView(field.clone()), "ListViewArray"),
2345            (DataType::LargeListView(field.clone()), "LargeListViewArray"),
2346        ] {
2347            let schema = Arc::new(Schema::new(vec![Field::new("lv", data_type, true)]));
2348            let buf = r#"
2349            {"lv": [1, 2, 3]}
2350            {"lv": [4, null]}
2351            "#;
2352
2353            let error = ReaderBuilder::new(schema)
2354                .build(Cursor::new(buf.as_bytes()))
2355                .unwrap()
2356                .collect::<Result<Vec<_>, _>>()
2357                .unwrap_err();
2358
2359            assert_eq!(
2360                error.to_string(),
2361                format!(
2362                    "Invalid argument error: Non-nullable field of {array_type} \"item\" cannot contain nulls"
2363                )
2364            );
2365        }
2366    }
2367
2368    #[test]
2369    fn test_fixed_size_list() {
2370        let buf = r#"
2371        {"a": [1, 2, 3]}
2372        {"a": [4, 5, 6]}
2373        {"a": [7, 8, 9]}
2374        "#;
2375
2376        let field = Field::new_list_field(DataType::Int32, true);
2377        let schema = Arc::new(Schema::new(vec![Field::new(
2378            "a",
2379            DataType::FixedSizeList(Arc::new(field), 3),
2380            false,
2381        )]));
2382
2383        let batches = do_read(buf, 1024, false, false, schema);
2384        assert_eq!(batches.len(), 1);
2385
2386        let col = batches[0].column(0).as_fixed_size_list();
2387        assert_eq!(col.len(), 3);
2388        assert_eq!(col.value_length(), 3);
2389
2390        let values = col.values().as_primitive::<Int32Type>();
2391        assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8, 9]);
2392    }
2393
2394    #[test]
2395    fn test_fixed_size_list_nullable() {
2396        let buf = r#"
2397        {"a": [1, 2]}
2398        {"a": null}
2399        {"a": [3, null]}
2400        "#;
2401
2402        let field = Field::new_list_field(DataType::Int32, true);
2403        let schema = Arc::new(Schema::new(vec![Field::new(
2404            "a",
2405            DataType::FixedSizeList(Arc::new(field), 2),
2406            true,
2407        )]));
2408
2409        let batches = do_read(buf, 1024, false, false, schema);
2410        assert_eq!(batches.len(), 1);
2411
2412        let col = batches[0].column(0).as_fixed_size_list();
2413        assert_eq!(col.len(), 3);
2414        assert!(col.is_valid(0));
2415        assert!(col.is_null(1));
2416        assert!(col.is_valid(2));
2417
2418        let values = col.values().as_primitive::<Int32Type>();
2419        assert_eq!(values.value(0), 1);
2420        assert_eq!(values.value(1), 2);
2421        assert_eq!(values.value(4), 3);
2422        assert!(values.is_null(5));
2423    }
2424
2425    #[test]
2426    fn test_fixed_size_list_zero_size_non_nullable() {
2427        let buf = r#"
2428        {"a": []}
2429        {"a": []}
2430        {"a": []}
2431        "#;
2432
2433        let field = Field::new_list_field(DataType::Int32, true);
2434        let schema = Arc::new(Schema::new(vec![Field::new(
2435            "a",
2436            DataType::FixedSizeList(Arc::new(field), 0),
2437            false,
2438        )]));
2439
2440        let batches = do_read(buf, 1024, false, false, schema);
2441        assert_eq!(batches.len(), 1);
2442
2443        let col = batches[0].column(0).as_fixed_size_list();
2444        assert_eq!(col.len(), 3);
2445        assert_eq!(col.value_length(), 0);
2446
2447        let values = col.values().as_primitive::<Int32Type>();
2448        assert!(values.values().is_empty());
2449    }
2450
2451    #[test]
2452    fn test_fixed_size_list_wrong_size() {
2453        let buf = r#"{"a": [1, 2, 3]}"#;
2454
2455        let field = Field::new_list_field(DataType::Int32, true);
2456        let schema = Arc::new(Schema::new(vec![Field::new(
2457            "a",
2458            DataType::FixedSizeList(Arc::new(field), 2),
2459            false,
2460        )]));
2461
2462        let err = ReaderBuilder::new(schema)
2463            .build(Cursor::new(buf.as_bytes()))
2464            .unwrap()
2465            .next()
2466            .unwrap()
2467            .unwrap_err();
2468
2469        assert!(err.to_string().contains("expected 2 but got 3"), "{}", err);
2470    }
2471
2472    #[test]
2473    fn test_fixed_size_list_nested() {
2474        let buf = r#"
2475        {"a": [[1, 2], [3, 4]]}
2476        {"a": [[5, 6], [7, 8]]}
2477        "#;
2478
2479        let inner_field = Field::new_list_field(DataType::Int32, true);
2480        let inner_type = DataType::FixedSizeList(Arc::new(inner_field), 2);
2481        let outer_field = Arc::new(Field::new_list_field(inner_type.clone(), true));
2482        let schema = Arc::new(Schema::new(vec![Field::new(
2483            "a",
2484            DataType::FixedSizeList(outer_field, 2),
2485            false,
2486        )]));
2487
2488        let batches = do_read(buf, 1024, false, false, schema);
2489        assert_eq!(batches.len(), 1);
2490
2491        let col = batches[0].column(0).as_fixed_size_list();
2492        assert_eq!(col.len(), 2);
2493        assert_eq!(col.value_length(), 2);
2494
2495        let inner = col.values().as_fixed_size_list();
2496        assert_eq!(inner.len(), 4);
2497        assert_eq!(inner.value_length(), 2);
2498
2499        let values = inner.values().as_primitive::<Int32Type>();
2500        assert_eq!(values.values(), &[1, 2, 3, 4, 5, 6, 7, 8]);
2501    }
2502
2503    #[test]
2504    fn test_fixed_size_list_ignore_type_conflicts() {
2505        let field = Field::new("item", DataType::Int32, true);
2506        let schema = Arc::new(Schema::new(vec![Field::new(
2507            "a",
2508            DataType::FixedSizeList(Arc::new(field), 2),
2509            true,
2510        )]));
2511
2512        let json = vec![
2513            json!({"a": [1, 2]}),
2514            json!({"a": "not a list"}),
2515            json!({"a": 42}),
2516            json!({"a": [6, 7]}),
2517        ];
2518
2519        let mut decoder = ReaderBuilder::new(schema)
2520            .with_ignore_type_conflicts(true)
2521            .build_decoder()
2522            .unwrap();
2523        decoder.serialize(&json).unwrap();
2524        let batch = decoder.flush().unwrap().unwrap();
2525
2526        let col = batch.column(0).as_fixed_size_list();
2527        assert_eq!(col.len(), 4);
2528        assert!(col.is_valid(0));
2529        assert!(col.is_null(1)); // string -> null
2530        assert!(col.is_null(2)); // number -> null
2531        assert!(col.is_valid(3));
2532
2533        let values = col.values().as_primitive::<Int32Type>();
2534        assert_eq!(values.value(0), 1);
2535        assert_eq!(values.value(1), 2);
2536        assert_eq!(values.value(6), 6);
2537        assert_eq!(values.value(7), 7);
2538    }
2539
2540    #[test]
2541    fn test_skip_empty_lines() {
2542        let schema = Schema::new(vec![Field::new("a", DataType::Int64, true)]);
2543        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(64);
2544        let json_content = "
2545        {\"a\": 1}
2546        {\"a\": 2}
2547        {\"a\": 3}";
2548        let mut reader = builder.build(Cursor::new(json_content)).unwrap();
2549        let batch = reader.next().unwrap().unwrap();
2550
2551        assert_eq!(1, batch.num_columns());
2552        assert_eq!(3, batch.num_rows());
2553
2554        let schema = reader.schema();
2555        let c = schema.column_with_name("a").unwrap();
2556        assert_eq!(&DataType::Int64, c.1.data_type());
2557    }
2558
2559    #[test]
2560    fn test_with_multiple_batches() {
2561        let file = File::open("test/data/basic_nulls.json").unwrap();
2562        let mut reader = BufReader::new(file);
2563        let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2564        reader.rewind().unwrap();
2565
2566        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2567        let mut reader = builder.build(reader).unwrap();
2568
2569        let mut num_records = Vec::new();
2570        while let Some(rb) = reader.next().transpose().unwrap() {
2571            num_records.push(rb.num_rows());
2572        }
2573
2574        assert_eq!(vec![5, 5, 2], num_records);
2575    }
2576
2577    #[test]
2578    fn test_timestamp_from_json_seconds() {
2579        let schema = Schema::new(vec![Field::new(
2580            "a",
2581            DataType::Timestamp(TimeUnit::Second, None),
2582            true,
2583        )]);
2584
2585        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2586        let batch = reader.next().unwrap().unwrap();
2587
2588        assert_eq!(1, batch.num_columns());
2589        assert_eq!(12, batch.num_rows());
2590
2591        let schema = reader.schema();
2592        let batch_schema = batch.schema();
2593        assert_eq!(schema, batch_schema);
2594
2595        let a = schema.column_with_name("a").unwrap();
2596        assert_eq!(
2597            &DataType::Timestamp(TimeUnit::Second, None),
2598            a.1.data_type()
2599        );
2600
2601        let aa = batch.column(a.0).as_primitive::<TimestampSecondType>();
2602        assert!(aa.is_valid(0));
2603        assert!(!aa.is_valid(1));
2604        assert!(!aa.is_valid(2));
2605        assert_eq!(1, aa.value(0));
2606        assert_eq!(1, aa.value(3));
2607        assert_eq!(5, aa.value(7));
2608    }
2609
2610    #[test]
2611    fn test_timestamp_from_json_milliseconds() {
2612        let schema = Schema::new(vec![Field::new(
2613            "a",
2614            DataType::Timestamp(TimeUnit::Millisecond, None),
2615            true,
2616        )]);
2617
2618        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2619        let batch = reader.next().unwrap().unwrap();
2620
2621        assert_eq!(1, batch.num_columns());
2622        assert_eq!(12, batch.num_rows());
2623
2624        let schema = reader.schema();
2625        let batch_schema = batch.schema();
2626        assert_eq!(schema, batch_schema);
2627
2628        let a = schema.column_with_name("a").unwrap();
2629        assert_eq!(
2630            &DataType::Timestamp(TimeUnit::Millisecond, None),
2631            a.1.data_type()
2632        );
2633
2634        let aa = batch.column(a.0).as_primitive::<TimestampMillisecondType>();
2635        assert!(aa.is_valid(0));
2636        assert!(!aa.is_valid(1));
2637        assert!(!aa.is_valid(2));
2638        assert_eq!(1, aa.value(0));
2639        assert_eq!(1, aa.value(3));
2640        assert_eq!(5, aa.value(7));
2641    }
2642
2643    #[test]
2644    fn test_date_from_json_milliseconds() {
2645        let schema = Schema::new(vec![Field::new("a", DataType::Date64, true)]);
2646
2647        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2648        let batch = reader.next().unwrap().unwrap();
2649
2650        assert_eq!(1, batch.num_columns());
2651        assert_eq!(12, batch.num_rows());
2652
2653        let schema = reader.schema();
2654        let batch_schema = batch.schema();
2655        assert_eq!(schema, batch_schema);
2656
2657        let a = schema.column_with_name("a").unwrap();
2658        assert_eq!(&DataType::Date64, a.1.data_type());
2659
2660        let aa = batch.column(a.0).as_primitive::<Date64Type>();
2661        assert!(aa.is_valid(0));
2662        assert!(!aa.is_valid(1));
2663        assert!(!aa.is_valid(2));
2664        assert_eq!(1, aa.value(0));
2665        assert_eq!(1, aa.value(3));
2666        assert_eq!(5, aa.value(7));
2667    }
2668
2669    #[test]
2670    fn test_time_from_json_nanoseconds() {
2671        let schema = Schema::new(vec![Field::new(
2672            "a",
2673            DataType::Time64(TimeUnit::Nanosecond),
2674            true,
2675        )]);
2676
2677        let mut reader = read_file("test/data/basic_nulls.json", Some(schema));
2678        let batch = reader.next().unwrap().unwrap();
2679
2680        assert_eq!(1, batch.num_columns());
2681        assert_eq!(12, batch.num_rows());
2682
2683        let schema = reader.schema();
2684        let batch_schema = batch.schema();
2685        assert_eq!(schema, batch_schema);
2686
2687        let a = schema.column_with_name("a").unwrap();
2688        assert_eq!(&DataType::Time64(TimeUnit::Nanosecond), a.1.data_type());
2689
2690        let aa = batch.column(a.0).as_primitive::<Time64NanosecondType>();
2691        assert!(aa.is_valid(0));
2692        assert!(!aa.is_valid(1));
2693        assert!(!aa.is_valid(2));
2694        assert_eq!(1, aa.value(0));
2695        assert_eq!(1, aa.value(3));
2696        assert_eq!(5, aa.value(7));
2697    }
2698
2699    #[test]
2700    fn test_json_iterator() {
2701        let file = File::open("test/data/basic.json").unwrap();
2702        let mut reader = BufReader::new(file);
2703        let (schema, _) = infer_json_schema(&mut reader, None).unwrap();
2704        reader.rewind().unwrap();
2705
2706        let builder = ReaderBuilder::new(Arc::new(schema)).with_batch_size(5);
2707        let reader = builder.build(reader).unwrap();
2708        let schema = reader.schema();
2709        let (col_a_index, _) = schema.column_with_name("a").unwrap();
2710
2711        let mut sum_num_rows = 0;
2712        let mut num_batches = 0;
2713        let mut sum_a = 0;
2714        for batch in reader {
2715            let batch = batch.unwrap();
2716            assert_eq!(8, batch.num_columns());
2717            sum_num_rows += batch.num_rows();
2718            num_batches += 1;
2719            let batch_schema = batch.schema();
2720            assert_eq!(schema, batch_schema);
2721            let a_array = batch.column(col_a_index).as_primitive::<Int64Type>();
2722            sum_a += (0..a_array.len()).map(|i| a_array.value(i)).sum::<i64>();
2723        }
2724        assert_eq!(12, sum_num_rows);
2725        assert_eq!(3, num_batches);
2726        assert_eq!(100000000000011, sum_a);
2727    }
2728
2729    #[test]
2730    fn test_decoder_error() {
2731        let schema = Arc::new(Schema::new(vec![Field::new_struct(
2732            "a",
2733            vec![Field::new("child", DataType::Int32, false)],
2734            true,
2735        )]));
2736
2737        let mut decoder = ReaderBuilder::new(schema.clone()).build_decoder().unwrap();
2738        let _ = decoder.decode(r#"{"a": { "child":"#.as_bytes()).unwrap();
2739        assert!(decoder.tape_decoder.has_partial_row());
2740        assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2741        let _ = decoder.flush().unwrap_err();
2742        assert!(decoder.tape_decoder.has_partial_row());
2743        assert_eq!(decoder.tape_decoder.num_buffered_rows(), 1);
2744
2745        let parse_err = |s: &str| {
2746            ReaderBuilder::new(schema.clone())
2747                .build(Cursor::new(s.as_bytes()))
2748                .unwrap()
2749                .next()
2750                .unwrap()
2751                .unwrap_err()
2752                .to_string()
2753        };
2754
2755        let err = parse_err(r#"{"a": 123}"#);
2756        assert_eq!(
2757            err,
2758            "Json error: whilst decoding field 'a': expected { got 123"
2759        );
2760
2761        let err = parse_err(r#"{"a": ["bar"]}"#);
2762        assert_eq!(
2763            err,
2764            r#"Json error: whilst decoding field 'a': expected { got ["bar"]"#
2765        );
2766
2767        let err = parse_err(r#"{"a": []}"#);
2768        assert_eq!(
2769            err,
2770            "Json error: whilst decoding field 'a': expected { got []"
2771        );
2772
2773        let err = parse_err(r#"{"a": [{"child": 234}]}"#);
2774        assert_eq!(
2775            err,
2776            r#"Json error: whilst decoding field 'a': expected { got [{"child": 234}]"#
2777        );
2778
2779        let err = parse_err(r#"{"a": [{"child": {"foo": [{"foo": ["bar"]}]}}]}"#);
2780        assert_eq!(
2781            err,
2782            r#"Json error: whilst decoding field 'a': expected { got [{"child": {"foo": [{"foo": ["bar"]}]}}]"#
2783        );
2784
2785        let err = parse_err(r#"{"a": true}"#);
2786        assert_eq!(
2787            err,
2788            "Json error: whilst decoding field 'a': expected { got true"
2789        );
2790
2791        let err = parse_err(r#"{"a": false}"#);
2792        assert_eq!(
2793            err,
2794            "Json error: whilst decoding field 'a': expected { got false"
2795        );
2796
2797        let err = parse_err(r#"{"a": "foo"}"#);
2798        assert_eq!(
2799            err,
2800            "Json error: whilst decoding field 'a': expected { got \"foo\""
2801        );
2802
2803        let err = parse_err(r#"{"a": {"child": false}}"#);
2804        assert_eq!(
2805            err,
2806            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got false"
2807        );
2808
2809        let err = parse_err(r#"{"a": {"child": []}}"#);
2810        assert_eq!(
2811            err,
2812            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got []"
2813        );
2814
2815        let err = parse_err(r#"{"a": {"child": [123]}}"#);
2816        assert_eq!(
2817            err,
2818            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123]"
2819        );
2820
2821        let err = parse_err(r#"{"a": {"child": [123, 3465346]}}"#);
2822        assert_eq!(
2823            err,
2824            "Json error: whilst decoding field 'a': whilst decoding field 'child': expected primitive got [123, 3465346]"
2825        );
2826    }
2827
2828    #[test]
2829    fn test_serialize_timestamp() {
2830        let json = vec![
2831            json!({"timestamp": 1681319393}),
2832            json!({"timestamp": "1970-01-01T00:00:00+02:00"}),
2833        ];
2834        let schema = Schema::new(vec![Field::new(
2835            "timestamp",
2836            DataType::Timestamp(TimeUnit::Second, None),
2837            true,
2838        )]);
2839        let mut decoder = ReaderBuilder::new(Arc::new(schema))
2840            .build_decoder()
2841            .unwrap();
2842        decoder.serialize(&json).unwrap();
2843        let batch = decoder.flush().unwrap().unwrap();
2844        assert_eq!(batch.num_rows(), 2);
2845        assert_eq!(batch.num_columns(), 1);
2846        let values = batch.column(0).as_primitive::<TimestampSecondType>();
2847        assert_eq!(values.values(), &[1681319393, -7200]);
2848    }
2849
2850    #[test]
2851    fn test_serialize_decimal() {
2852        let json = vec![
2853            json!({"decimal": 1.234}),
2854            json!({"decimal": "1.234"}),
2855            json!({"decimal": 1234}),
2856            json!({"decimal": "1234"}),
2857        ];
2858        let schema = Schema::new(vec![Field::new(
2859            "decimal",
2860            DataType::Decimal128(10, 3),
2861            true,
2862        )]);
2863        let mut decoder = ReaderBuilder::new(Arc::new(schema))
2864            .build_decoder()
2865            .unwrap();
2866        decoder.serialize(&json).unwrap();
2867        let batch = decoder.flush().unwrap().unwrap();
2868        assert_eq!(batch.num_rows(), 4);
2869        assert_eq!(batch.num_columns(), 1);
2870        let values = batch.column(0).as_primitive::<Decimal128Type>();
2871        assert_eq!(values.values(), &[1234, 1234, 1234000, 1234000]);
2872    }
2873
2874    #[test]
2875    fn test_serde_field() {
2876        let field = Field::new("int", DataType::Int32, true);
2877        let mut decoder = ReaderBuilder::new_with_field(field)
2878            .build_decoder()
2879            .unwrap();
2880        decoder.serialize(&[1_i32, 2, 3, 4]).unwrap();
2881        let b = decoder.flush().unwrap().unwrap();
2882        let values = b.column(0).as_primitive::<Int32Type>().values();
2883        assert_eq!(values, &[1, 2, 3, 4]);
2884    }
2885
2886    #[test]
2887    fn test_serde_large_numbers() {
2888        let field = Field::new("int", DataType::Int64, true);
2889        let mut decoder = ReaderBuilder::new_with_field(field)
2890            .build_decoder()
2891            .unwrap();
2892
2893        decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2894        let b = decoder.flush().unwrap().unwrap();
2895        let values = b.column(0).as_primitive::<Int64Type>().values();
2896        assert_eq!(values, &[1699148028689, 2, 3, 4]);
2897
2898        let field = Field::new(
2899            "int",
2900            DataType::Timestamp(TimeUnit::Microsecond, None),
2901            true,
2902        );
2903        let mut decoder = ReaderBuilder::new_with_field(field)
2904            .build_decoder()
2905            .unwrap();
2906
2907        decoder.serialize(&[1699148028689_u64, 2, 3, 4]).unwrap();
2908        let b = decoder.flush().unwrap().unwrap();
2909        let values = b
2910            .column(0)
2911            .as_primitive::<TimestampMicrosecondType>()
2912            .values();
2913        assert_eq!(values, &[1699148028689, 2, 3, 4]);
2914    }
2915
2916    #[test]
2917    fn test_coercing_primitive_into_string_decoder() {
2918        let buf = &format!(
2919            r#"[{{"a": 1, "b": "A", "c": "T"}}, {{"a": 2, "b": "BB", "c": "F"}}, {{"a": {}, "b": 123, "c": false}}, {{"a": {}, "b": 789, "c": true}}]"#,
2920            (i32::MAX as i64 + 10),
2921            i64::MAX - 10
2922        );
2923        let schema = Schema::new(vec![
2924            Field::new("a", DataType::Float64, true),
2925            Field::new("b", DataType::Utf8, true),
2926            Field::new("c", DataType::Utf8, true),
2927        ]);
2928        let json_array: Vec<serde_json::Value> = serde_json::from_str(buf).unwrap();
2929        let schema_ref = Arc::new(schema);
2930
2931        // read record batches
2932        let reader = ReaderBuilder::new(schema_ref.clone()).with_coerce_primitive(true);
2933        let mut decoder = reader.build_decoder().unwrap();
2934        decoder.serialize(json_array.as_slice()).unwrap();
2935        let batch = decoder.flush().unwrap().unwrap();
2936        assert_eq!(
2937            batch,
2938            RecordBatch::try_new(
2939                schema_ref,
2940                vec![
2941                    Arc::new(Float64Array::from(vec![
2942                        1.0,
2943                        2.0,
2944                        (i32::MAX as i64 + 10) as f64,
2945                        (i64::MAX - 10) as f64
2946                    ])),
2947                    Arc::new(StringArray::from(vec!["A", "BB", "123", "789"])),
2948                    Arc::new(StringArray::from(vec!["T", "F", "false", "true"])),
2949                ]
2950            )
2951            .unwrap()
2952        );
2953    }
2954
2955    #[test]
2956    fn test_serialize_f32_into_string() {
2957        // Coercing an f32 into a string column must render the value, not its raw bit pattern.
2958        let field = Field::new("f", DataType::Utf8, true);
2959        let mut decoder = ReaderBuilder::new_with_field(field)
2960            .with_coerce_primitive(true)
2961            .build_decoder()
2962            .unwrap();
2963        decoder.serialize(&[1.5_f32, -2.25_f32]).unwrap();
2964        let batch = decoder.flush().unwrap().unwrap();
2965        let values = batch.column(0).as_string::<i32>();
2966        assert_eq!(values.value(0), "1.5");
2967        assert_eq!(values.value(1), "-2.25");
2968    }
2969
2970    // Parse the given `row` in `struct_mode` as a type given by fields.
2971    //
2972    // If as_struct == true, wrap the fields in a Struct field with name "r".
2973    // If as_struct == false, wrap the fields in a Schema.
2974    fn _parse_structs(
2975        row: &str,
2976        struct_mode: StructMode,
2977        fields: Fields,
2978        as_struct: bool,
2979    ) -> Result<RecordBatch, ArrowError> {
2980        let builder = if as_struct {
2981            ReaderBuilder::new_with_field(Field::new("r", DataType::Struct(fields), true))
2982        } else {
2983            ReaderBuilder::new(Arc::new(Schema::new(fields)))
2984        };
2985        builder
2986            .with_struct_mode(struct_mode)
2987            .build(Cursor::new(row.as_bytes()))
2988            .unwrap()
2989            .next()
2990            .unwrap()
2991    }
2992
2993    #[test]
2994    fn test_struct_decoding_list_length() {
2995        use arrow_array::array;
2996
2997        let row = "[1, 2]";
2998
2999        let mut fields = vec![Field::new("a", DataType::Int32, true)];
3000        let too_few_fields = Fields::from(fields.clone());
3001        fields.push(Field::new("b", DataType::Int32, true));
3002        let correct_fields = Fields::from(fields.clone());
3003        fields.push(Field::new("c", DataType::Int32, true));
3004        let too_many_fields = Fields::from(fields.clone());
3005
3006        let parse = |fields: Fields, as_struct: bool| {
3007            _parse_structs(row, StructMode::ListOnly, fields, as_struct)
3008        };
3009
3010        let expected_row = StructArray::new(
3011            correct_fields.clone(),
3012            vec![
3013                Arc::new(array::Int32Array::from(vec![1])),
3014                Arc::new(array::Int32Array::from(vec![2])),
3015            ],
3016            None,
3017        );
3018        let row_field = Field::new("r", DataType::Struct(correct_fields.clone()), true);
3019
3020        assert_eq!(
3021            parse(too_few_fields.clone(), true).unwrap_err().to_string(),
3022            "Json error: found extra columns for 1 fields".to_string()
3023        );
3024        assert_eq!(
3025            parse(too_few_fields, false).unwrap_err().to_string(),
3026            "Json error: found extra columns for 1 fields".to_string()
3027        );
3028        assert_eq!(
3029            parse(correct_fields.clone(), true).unwrap(),
3030            RecordBatch::try_new(
3031                Arc::new(Schema::new(vec![row_field])),
3032                vec![Arc::new(expected_row.clone())]
3033            )
3034            .unwrap()
3035        );
3036        assert_eq!(
3037            parse(correct_fields, false).unwrap(),
3038            RecordBatch::from(expected_row)
3039        );
3040        assert_eq!(
3041            parse(too_many_fields.clone(), true)
3042                .unwrap_err()
3043                .to_string(),
3044            "Json error: found 2 columns for 3 fields".to_string()
3045        );
3046        assert_eq!(
3047            parse(too_many_fields, false).unwrap_err().to_string(),
3048            "Json error: found 2 columns for 3 fields".to_string()
3049        );
3050    }
3051
3052    #[test]
3053    fn test_struct_decoding() {
3054        use arrow_array::builder;
3055
3056        let nested_object_json = r#"{"a": {"b": [1, 2], "c": {"d": 3}}}"#;
3057        let nested_list_json = r#"[[[1, 2], {"d": 3}]]"#;
3058        let nested_mixed_json = r#"{"a": [[1, 2], {"d": 3}]}"#;
3059
3060        let struct_fields = Fields::from(vec![
3061            Field::new("b", DataType::new_list(DataType::Int32, true), true),
3062            Field::new_map(
3063                "c",
3064                "entries",
3065                Field::new("keys", DataType::Utf8, false),
3066                Field::new("values", DataType::Int32, true),
3067                false,
3068                false,
3069            ),
3070        ]);
3071
3072        let list_array =
3073            ListArray::from_iter_primitive::<Int32Type, _, _>(vec![Some(vec![Some(1), Some(2)])]);
3074
3075        let map_array = {
3076            let mut map_builder = builder::MapBuilder::new(
3077                None,
3078                builder::StringBuilder::new(),
3079                builder::Int32Builder::new(),
3080            );
3081            map_builder.keys().append_value("d");
3082            map_builder.values().append_value(3);
3083            map_builder.append(true).unwrap();
3084            map_builder.finish()
3085        };
3086
3087        let struct_array = StructArray::new(
3088            struct_fields.clone(),
3089            vec![Arc::new(list_array), Arc::new(map_array)],
3090            None,
3091        );
3092
3093        let fields = Fields::from(vec![Field::new("a", DataType::Struct(struct_fields), true)]);
3094        let schema = Arc::new(Schema::new(fields.clone()));
3095        let expected = RecordBatch::try_new(schema.clone(), vec![Arc::new(struct_array)]).unwrap();
3096
3097        let parse = |row: &str, struct_mode: StructMode| {
3098            _parse_structs(row, struct_mode, fields.clone(), false)
3099        };
3100
3101        assert_eq!(
3102            parse(nested_object_json, StructMode::ObjectOnly).unwrap(),
3103            expected
3104        );
3105        assert_eq!(
3106            parse(nested_list_json, StructMode::ObjectOnly)
3107                .unwrap_err()
3108                .to_string(),
3109            "Json error: expected { got [[[1, 2], {\"d\": 3}]]".to_owned()
3110        );
3111        assert_eq!(
3112            parse(nested_mixed_json, StructMode::ObjectOnly)
3113                .unwrap_err()
3114                .to_string(),
3115            "Json error: whilst decoding field 'a': expected { got [[1, 2], {\"d\": 3}]".to_owned()
3116        );
3117
3118        assert_eq!(
3119            parse(nested_list_json, StructMode::ListOnly).unwrap(),
3120            expected
3121        );
3122        assert_eq!(
3123            parse(nested_object_json, StructMode::ListOnly)
3124                .unwrap_err()
3125                .to_string(),
3126            "Json error: expected [ got {\"a\": {\"b\": [1, 2]\"c\": {\"d\": 3}}}".to_owned()
3127        );
3128        assert_eq!(
3129            parse(nested_mixed_json, StructMode::ListOnly)
3130                .unwrap_err()
3131                .to_string(),
3132            "Json error: expected [ got {\"a\": [[1, 2], {\"d\": 3}]}".to_owned()
3133        );
3134    }
3135
3136    // Test cases:
3137    // [] -> RecordBatch row with no entries.  Schema = [('a', Int32)] -> Error
3138    // [] -> RecordBatch row with no entries. Schema = [('r', [('a', Int32)])] -> Error
3139    // [] -> StructArray row with no entries. Fields [('a', Int32')] -> Error
3140    // [[]] -> RecordBatch row with empty struct entry. Schema = [('r', [('a', Int32)])] -> Error
3141    #[test]
3142    fn test_struct_decoding_empty_list() {
3143        let int_field = Field::new("a", DataType::Int32, true);
3144        let struct_field = Field::new(
3145            "r",
3146            DataType::Struct(Fields::from(vec![int_field.clone()])),
3147            true,
3148        );
3149
3150        let parse = |row: &str, as_struct: bool, field: Field| {
3151            _parse_structs(
3152                row,
3153                StructMode::ListOnly,
3154                Fields::from(vec![field]),
3155                as_struct,
3156            )
3157        };
3158
3159        // Missing fields
3160        assert_eq!(
3161            parse("[]", true, struct_field.clone())
3162                .unwrap_err()
3163                .to_string(),
3164            "Json error: found 0 columns for 1 fields".to_owned()
3165        );
3166        assert_eq!(
3167            parse("[]", false, int_field.clone())
3168                .unwrap_err()
3169                .to_string(),
3170            "Json error: found 0 columns for 1 fields".to_owned()
3171        );
3172        assert_eq!(
3173            parse("[]", false, struct_field.clone())
3174                .unwrap_err()
3175                .to_string(),
3176            "Json error: found 0 columns for 1 fields".to_owned()
3177        );
3178        assert_eq!(
3179            parse("[[]]", false, struct_field.clone())
3180                .unwrap_err()
3181                .to_string(),
3182            "Json error: whilst decoding field 'r': found 0 columns for 1 fields".to_owned()
3183        );
3184    }
3185
3186    #[test]
3187    fn test_decode_list_struct_with_wrong_types() {
3188        let int_field = Field::new("a", DataType::Int32, true);
3189        let struct_field = Field::new(
3190            "r",
3191            DataType::Struct(Fields::from(vec![int_field.clone()])),
3192            true,
3193        );
3194
3195        let parse = |row: &str, as_struct: bool, field: Field| {
3196            _parse_structs(
3197                row,
3198                StructMode::ListOnly,
3199                Fields::from(vec![field]),
3200                as_struct,
3201            )
3202        };
3203
3204        // Wrong values
3205        assert_eq!(
3206            parse(r#"[["a"]]"#, false, struct_field.clone())
3207                .unwrap_err()
3208                .to_string(),
3209            "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3210        );
3211        assert_eq!(
3212            parse(r#"[["a"]]"#, true, struct_field.clone())
3213                .unwrap_err()
3214                .to_string(),
3215            "Json error: whilst decoding field 'r': whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3216        );
3217        assert_eq!(
3218            parse(r#"["a"]"#, true, int_field.clone())
3219                .unwrap_err()
3220                .to_string(),
3221            "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3222        );
3223        assert_eq!(
3224            parse(r#"["a"]"#, false, int_field.clone())
3225                .unwrap_err()
3226                .to_string(),
3227            "Json error: whilst decoding field 'a': failed to parse \"a\" as Int32".to_owned()
3228        );
3229    }
3230
3231    #[test]
3232    fn test_type_conflict_nulls() {
3233        let schema = Schema::new(vec![
3234            Field::new("null", DataType::Null, true),
3235            Field::new("bool", DataType::Boolean, true),
3236            Field::new("primitive", DataType::Int32, true),
3237            Field::new("numeric", DataType::Decimal128(10, 3), true),
3238            Field::new("string", DataType::Utf8, true),
3239            Field::new("string_view", DataType::Utf8View, true),
3240            Field::new(
3241                "timestamp",
3242                DataType::Timestamp(TimeUnit::Second, None),
3243                true,
3244            ),
3245            Field::new(
3246                "array",
3247                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3248                true,
3249            ),
3250            Field::new(
3251                "map",
3252                DataType::Map(
3253                    Arc::new(Field::new(
3254                        "entries",
3255                        DataType::Struct(Fields::from(vec![
3256                            Field::new("keys", DataType::Utf8, false),
3257                            Field::new("values", DataType::Utf8, true),
3258                        ])),
3259                        false, // not nullable
3260                    )),
3261                    false, // not sorted
3262                ),
3263                true, // nullable
3264            ),
3265            Field::new(
3266                "struct",
3267                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3268                true,
3269            ),
3270        ]);
3271
3272        // A compatible value for each schema field above, in schema order
3273        let json_values = vec![
3274            json!(null),
3275            json!(true),
3276            json!(42),
3277            json!(1.234),
3278            json!("hi"),
3279            json!("ho"),
3280            json!("1970-01-01T00:00:00+02:00"),
3281            json!([1, "ho", 3]),
3282            json!({"k": "value"}),
3283            json!({"a": 1}),
3284        ];
3285
3286        // Create a set of JSON rows that rotates each value past every field
3287        let json: Vec<_> = (0..json_values.len())
3288            .map(|i| {
3289                let pairs = json_values[i..]
3290                    .iter()
3291                    .chain(json_values[..i].iter())
3292                    .zip(&schema.fields)
3293                    .map(|(v, f)| (f.name().to_string(), v.clone()))
3294                    .collect();
3295                serde_json::Value::Object(pairs)
3296            })
3297            .collect();
3298        let mut decoder = ReaderBuilder::new(Arc::new(schema))
3299            .with_ignore_type_conflicts(true)
3300            .with_coerce_primitive(true)
3301            .build_decoder()
3302            .unwrap();
3303        decoder.serialize(&json).unwrap();
3304        let batch = decoder.flush().unwrap().unwrap();
3305        assert_eq!(batch.num_rows(), 10);
3306        assert_eq!(batch.num_columns(), 10);
3307
3308        // NOTE: NullArray doesn't materialize any values (they're all NULL by definition)
3309        let _ = batch
3310            .column(0)
3311            .as_any()
3312            .downcast_ref::<NullArray>()
3313            .unwrap();
3314
3315        assert!(
3316            batch
3317                .column(1)
3318                .as_any()
3319                .downcast_ref::<BooleanArray>()
3320                .unwrap()
3321                .iter()
3322                .eq([
3323                    Some(true),
3324                    None,
3325                    None,
3326                    None,
3327                    None,
3328                    None,
3329                    None,
3330                    None,
3331                    None,
3332                    None
3333                ])
3334        );
3335
3336        assert!(batch.column(2).as_primitive::<Int32Type>().iter().eq([
3337            Some(42),
3338            Some(1),
3339            None,
3340            None,
3341            None,
3342            None,
3343            None,
3344            None,
3345            None,
3346            None
3347        ]));
3348
3349        assert!(batch.column(3).as_primitive::<Decimal128Type>().iter().eq([
3350            Some(1234),
3351            None,
3352            None,
3353            None,
3354            None,
3355            None,
3356            None,
3357            None,
3358            None,
3359            Some(42000)
3360        ]));
3361
3362        assert!(
3363            batch
3364                .column(4)
3365                .as_any()
3366                .downcast_ref::<StringArray>()
3367                .unwrap()
3368                .iter()
3369                .eq([
3370                    Some("hi"),
3371                    Some("ho"),
3372                    Some("1970-01-01T00:00:00+02:00"),
3373                    None,
3374                    None,
3375                    None,
3376                    None,
3377                    Some("true"),
3378                    Some("42"),
3379                    Some("1.234"),
3380                ])
3381        );
3382
3383        assert!(
3384            batch
3385                .column(5)
3386                .as_any()
3387                .downcast_ref::<StringViewArray>()
3388                .unwrap()
3389                .iter()
3390                .eq([
3391                    Some("ho"),
3392                    Some("1970-01-01T00:00:00+02:00"),
3393                    None,
3394                    None,
3395                    None,
3396                    None,
3397                    Some("true"),
3398                    Some("42"),
3399                    Some("1.234"),
3400                    Some("hi"),
3401                ])
3402        );
3403
3404        assert!(
3405            batch
3406                .column(6)
3407                .as_primitive::<TimestampSecondType>()
3408                .iter()
3409                .eq([
3410                    Some(-7200),
3411                    None,
3412                    None,
3413                    None,
3414                    None,
3415                    None,
3416                    Some(42),
3417                    None,
3418                    None,
3419                    None,
3420                ])
3421        );
3422
3423        let arrays = batch
3424            .column(7)
3425            .as_any()
3426            .downcast_ref::<ListArray>()
3427            .unwrap();
3428        assert_eq!(
3429            arrays.nulls(),
3430            Some(&NullBuffer::from(
3431                &[
3432                    true, false, false, false, false, false, false, false, false, false
3433                ][..]
3434            ))
3435        );
3436        assert_eq!(arrays.offsets()[1], 3);
3437        let array_values = arrays
3438            .values()
3439            .as_any()
3440            .downcast_ref::<Int32Array>()
3441            .unwrap();
3442        assert!(array_values.iter().eq([Some(1), None, Some(3)]));
3443
3444        let maps = batch.column(8).as_any().downcast_ref::<MapArray>().unwrap();
3445        assert_eq!(
3446            maps.nulls(),
3447            Some(&NullBuffer::from(
3448                // Both map and struct can parse
3449                &[
3450                    true, true, false, false, false, false, false, false, false, false
3451                ][..]
3452            ))
3453        );
3454        let map_keys = maps.keys().as_any().downcast_ref::<StringArray>().unwrap();
3455        assert!(map_keys.iter().eq([Some("k"), Some("a")]));
3456        let map_values = maps
3457            .values()
3458            .as_any()
3459            .downcast_ref::<StringArray>()
3460            .unwrap();
3461        assert!(map_values.iter().eq([Some("value"), Some("1")]));
3462
3463        let structs = batch
3464            .column(9)
3465            .as_any()
3466            .downcast_ref::<StructArray>()
3467            .unwrap();
3468        assert_eq!(
3469            structs.nulls(),
3470            Some(&NullBuffer::from(
3471                // Both map and struct can parse
3472                &[
3473                    true, false, false, false, false, false, false, false, false, true
3474                ][..]
3475            ))
3476        );
3477        let struct_fields = structs
3478            .column(0)
3479            .as_any()
3480            .downcast_ref::<Int32Array>()
3481            .unwrap();
3482        assert!(struct_fields.slice(0, 2).iter().eq([Some(1), None]));
3483    }
3484
3485    #[test]
3486    fn test_type_conflict_non_nullable() {
3487        let fields = [
3488            Field::new("bool", DataType::Boolean, false),
3489            Field::new("primitive", DataType::Int32, false),
3490            Field::new("numeric", DataType::Decimal128(10, 3), false),
3491            Field::new("string", DataType::Utf8, false),
3492            Field::new("string_view", DataType::Utf8View, false),
3493            Field::new(
3494                "timestamp",
3495                DataType::Timestamp(TimeUnit::Second, None),
3496                false,
3497            ),
3498            Field::new(
3499                "array",
3500                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3501                false,
3502            ),
3503            Field::new(
3504                "fixed_size_list",
3505                DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3506                false,
3507            ),
3508            Field::new(
3509                "map",
3510                DataType::Map(
3511                    Arc::new(Field::new(
3512                        "entries",
3513                        DataType::Struct(Fields::from(vec![
3514                            Field::new("keys", DataType::Utf8, false),
3515                            Field::new("values", DataType::Utf8, true),
3516                        ])),
3517                        false, // not nullable
3518                    )),
3519                    false, // not sorted
3520                ),
3521                false, // not nullable
3522            ),
3523            Field::new(
3524                "struct",
3525                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3526                false,
3527            ),
3528        ];
3529
3530        // Every field above will have a type conflict with at least one of these values
3531        let json_values = vec![json!(true), json!({"a": 1})];
3532
3533        for field in fields {
3534            let mut decoder = ReaderBuilder::new_with_field(field)
3535                .with_ignore_type_conflicts(true)
3536                .build_decoder()
3537                .unwrap();
3538            decoder.serialize(&json_values).unwrap();
3539            decoder
3540                .flush()
3541                .expect_err("type conflict on non-nullable type");
3542        }
3543    }
3544
3545    #[test]
3546    fn test_ignore_type_conflicts_disabled() {
3547        let fields = [
3548            Field::new("null", DataType::Null, true),
3549            Field::new("bool", DataType::Boolean, true),
3550            Field::new("primitive", DataType::Int32, true),
3551            Field::new("numeric", DataType::Decimal128(10, 3), true),
3552            Field::new("string", DataType::Utf8, true),
3553            Field::new("string_view", DataType::Utf8View, true),
3554            Field::new(
3555                "timestamp",
3556                DataType::Timestamp(TimeUnit::Second, None),
3557                true,
3558            ),
3559            Field::new(
3560                "array",
3561                DataType::List(Arc::new(Field::new("item", DataType::Int32, true))),
3562                true,
3563            ),
3564            Field::new(
3565                "fixed_size_list",
3566                DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Int32, true)), 2),
3567                true,
3568            ),
3569            Field::new(
3570                "map",
3571                DataType::Map(
3572                    Arc::new(Field::new(
3573                        "entries",
3574                        DataType::Struct(Fields::from(vec![
3575                            Field::new("keys", DataType::Utf8, false),
3576                            Field::new("values", DataType::Utf8, true),
3577                        ])),
3578                        false, // not nullable
3579                    )),
3580                    false, // not sorted
3581                ),
3582                true, // not nullable
3583            ),
3584            Field::new(
3585                "struct",
3586                DataType::Struct(Fields::from(vec![Field::new("a", DataType::Int32, true)])),
3587                true,
3588            ),
3589        ];
3590
3591        // Every field above will have a type conflict with at least one of these values
3592        let json_values = vec![json!(true), json!({"a": 1})];
3593
3594        for field in fields {
3595            let mut decoder = ReaderBuilder::new_with_field(field)
3596                .build_decoder()
3597                .unwrap();
3598            decoder.serialize(&json_values).unwrap();
3599            decoder
3600                .flush()
3601                .expect_err("type conflict on non-nullable type");
3602        }
3603    }
3604
3605    #[test]
3606    fn test_read_run_end_encoded() {
3607        let buf = r#"
3608        {"a": "x"}
3609        {"a": "x"}
3610        {"a": "y"}
3611        {"a": "y"}
3612        {"a": "y"}
3613        "#;
3614
3615        let ree_type = DataType::RunEndEncoded(
3616            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3617            Arc::new(Field::new("values", DataType::Utf8, true)),
3618        );
3619        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3620        let batches = do_read(buf, 1024, false, false, schema);
3621        assert_eq!(batches.len(), 1);
3622
3623        let col = batches[0].column(0);
3624        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3625
3626        // 5 logical values compressed into 2 runs
3627        assert_eq!(run_array.len(), 5);
3628        assert_eq!(run_array.run_ends().values(), &[2, 5]);
3629
3630        let values = run_array.values().as_string::<i32>();
3631        assert_eq!(values.len(), 2);
3632        assert_eq!(values.value(0), "x");
3633        assert_eq!(values.value(1), "y");
3634    }
3635
3636    #[test]
3637    fn test_read_run_end_encoded_consecutive_nulls() {
3638        let buf = r#"
3639        {"a": "x"}
3640        {}
3641        {}
3642        {}
3643        {"a": "y"}
3644        "#;
3645
3646        let ree_type = DataType::RunEndEncoded(
3647            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3648            Arc::new(Field::new("values", DataType::Utf8, true)),
3649        );
3650        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3651        let batches = do_read(buf, 1024, false, false, schema);
3652        assert_eq!(batches.len(), 1);
3653
3654        let col = batches[0].column(0);
3655        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3656
3657        // 5 logical values: "x", null, null, null, "y" → 3 runs
3658        assert_eq!(run_array.len(), 5);
3659        assert_eq!(run_array.run_ends().values(), &[1, 4, 5]);
3660
3661        let values = run_array.values().as_string::<i32>();
3662        assert_eq!(values.len(), 3);
3663        assert_eq!(values.value(0), "x");
3664        assert!(values.is_null(1));
3665        assert_eq!(values.value(2), "y");
3666    }
3667
3668    #[test]
3669    fn test_read_run_end_encoded_all_unique() {
3670        let buf = r#"
3671        {"a": 1}
3672        {"a": 2}
3673        {"a": 3}
3674        "#;
3675
3676        let ree_type = DataType::RunEndEncoded(
3677            Arc::new(Field::new("run_ends", DataType::Int32, false)),
3678            Arc::new(Field::new("values", DataType::Int32, true)),
3679        );
3680        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3681        let batches = do_read(buf, 1024, false, false, schema);
3682        assert_eq!(batches.len(), 1);
3683
3684        let col = batches[0].column(0);
3685        let run_array = col.as_run::<arrow_array::types::Int32Type>();
3686
3687        // No compression: 3 unique values → 3 runs
3688        assert_eq!(run_array.len(), 3);
3689        assert_eq!(run_array.run_ends().values(), &[1, 2, 3]);
3690    }
3691
3692    #[test]
3693    fn test_read_run_end_encoded_int16_run_ends() {
3694        let buf = r#"
3695        {"a": "x"}
3696        {"a": "x"}
3697        {"a": "y"}
3698        "#;
3699
3700        let ree_type = DataType::RunEndEncoded(
3701            Arc::new(Field::new("run_ends", DataType::Int16, false)),
3702            Arc::new(Field::new("values", DataType::Utf8, true)),
3703        );
3704        let schema = Arc::new(Schema::new(vec![Field::new("a", ree_type, true)]));
3705        let batches = do_read(buf, 1024, false, false, schema);
3706        assert_eq!(batches.len(), 1);
3707
3708        let col = batches[0].column(0);
3709        let run_array = col.as_run::<arrow_array::types::Int16Type>();
3710
3711        assert_eq!(run_array.len(), 3);
3712        assert_eq!(run_array.run_ends().values(), &[2i16, 3]);
3713    }
3714
3715    #[test]
3716    fn test_read_nested_run_end_encoded() {
3717        let buf = r#"
3718        {"a": "x"}
3719        {"a": "x"}
3720        {"a": "y"}
3721        "#;
3722
3723        // The outer REE compresses whole rows, while the inner REE compresses the
3724        // repeated string values produced by decoding those rows.
3725        let inner_type = DataType::RunEndEncoded(
3726            Arc::new(Field::new("run_ends", DataType::Int64, false)),
3727            Arc::new(Field::new("values", DataType::Utf8, true)),
3728        );
3729        let outer_type = DataType::RunEndEncoded(
3730            Arc::new(Field::new("run_ends", DataType::Int64, false)),
3731            Arc::new(Field::new("values", inner_type, true)),
3732        );
3733        let schema = Arc::new(Schema::new(vec![Field::new("a", outer_type, true)]));
3734        let batches = do_read(buf, 1024, false, false, schema);
3735        assert_eq!(batches.len(), 1);
3736
3737        let col = batches[0].column(0);
3738        let outer = col.as_run::<arrow_array::types::Int64Type>();
3739        // Three logical rows compress to two outer runs: ["x", "x"] and ["y"].
3740        assert_eq!(outer.len(), 3);
3741        assert_eq!(outer.run_ends().values(), &[2, 3]);
3742
3743        let nested = outer.values().as_run::<arrow_array::types::Int64Type>();
3744        // The physical values of the outer REE are themselves a two-element REE.
3745        assert_eq!(nested.len(), 2);
3746        assert_eq!(nested.run_ends().values(), &[1, 2]);
3747
3748        let nested_values = nested.values().as_string::<i32>();
3749        assert_eq!(nested_values.len(), 2);
3750        assert_eq!(nested_values.value(0), "x");
3751        assert_eq!(nested_values.value(1), "y");
3752    }
3753}