Skip to main content

akar_storage/
parquet_reader.rs

1//! Parquet reader for the COPY FROM command.
2//!
3//! Reads Parquet files and converts Arrow columnar data to Akar `Value`
4//! types based on a provided schema (column names + types from the catalog).
5
6use akar_catalog::CatalogColumn;
7use akar_common::types::{Date, Interval, LogicalTypeID, Timestamp, Value};
8use arrow::array::*;
9use arrow::datatypes::{DataType as ArrowDataType, TimeUnit};
10use arrow::record_batch::RecordBatch;
11
12/// Error type for Parquet reader operations.
13#[derive(Debug)]
14pub enum ParquetReaderError {
15    /// I/O or format error from the Parquet/Arrow layer.
16    ParquetError(String),
17    /// Schema mismatch: a column name from the catalog was not found in the file.
18    ColumnNotFound {
19        column_name: String,
20        available: Vec<String>,
21    },
22    /// Type mismatch: column exists but has incompatible type.
23    TypeMismatch {
24        column_name: String,
25        arrow_type: String,
26        expected_type: String,
27    },
28    /// Row-level conversion error.
29    ConversionError {
30        column_name: String,
31        row: usize,
32        message: String,
33    },
34}
35
36impl std::fmt::Display for ParquetReaderError {
37    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
38        match self {
39            ParquetReaderError::ParquetError(e) => write!(f, "Parquet error: {e}"),
40            ParquetReaderError::ColumnNotFound { column_name, available } => write!(
41                f,
42                "Column '{}' not found in Parquet file. Available columns: [{}]",
43                column_name,
44                available.join(", ")
45            ),
46            ParquetReaderError::TypeMismatch {
47                column_name,
48                arrow_type,
49                expected_type,
50            } => write!(
51                f,
52                "Type mismatch for column '{}': Parquet has {arrow_type}, expected {expected_type}",
53                column_name
54            ),
55            ParquetReaderError::ConversionError {
56                column_name,
57                row,
58                message,
59            } => write!(
60                f,
61                "Conversion error for column '{}' at row {row}: {message}",
62                column_name
63            ),
64        }
65    }
66}
67
68impl std::error::Error for ParquetReaderError {}
69
70impl From<parquet::errors::ParquetError> for ParquetReaderError {
71    fn from(e: parquet::errors::ParquetError) -> Self {
72        ParquetReaderError::ParquetError(e.to_string())
73    }
74}
75
76impl From<arrow::error::ArrowError> for ParquetReaderError {
77    fn from(e: arrow::error::ArrowError) -> Self {
78        ParquetReaderError::ParquetError(e.to_string())
79    }
80}
81
82/// Result alias for Parquet reader operations.
83pub type ParquetResult<T> = Result<T, ParquetReaderError>;
84
85/// Read a Parquet file and convert columns to Akar `Value`s matching the schema.
86///
87/// The `columns` parameter defines the target schema. Columns are matched by
88/// name against the Parquet file's schema; order in the result follows the
89/// `columns` slice order.
90///
91/// # Arguments
92///
93/// * `path` - Path to the `.parquet` file.
94/// * `columns` - Target column schema (name + type). Only these columns are
95///   read from the file.
96///
97/// # Returns
98///
99/// A vector of rows, where each row is a `Vec<Value>` with length equal to
100/// `columns.len()`.
101pub fn read_parquet(
102    path: &str,
103    vfs: &akar_common::file_system::VirtualFileSystemRegistry,
104    columns: &[CatalogColumn],
105) -> ParquetResult<Vec<Vec<Value>>> {
106    let mut file = vfs
107        .open_read(path)
108        .map_err(|e| ParquetReaderError::ParquetError(format!("Cannot open file: {e}")))?;
109
110    // For now, read the entire file into memory to pass it to ParquetRecordBatchReaderBuilder
111    // as bytes::Bytes, since Box<dyn FileRead> doesn't natively implement ChunkReader.
112    let mut buffer = Vec::new();
113    file.read_to_end(&mut buffer)
114        .map_err(|e| ParquetReaderError::ParquetError(format!("Read error: {e}")))?;
115    let bytes = bytes::Bytes::from(buffer);
116
117    let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(bytes)?;
118    let schema = builder.schema().clone();
119    let reader = builder.build()?;
120
121    // Build a map of column name → Arrow field for quick lookup
122    let arrow_fields: Vec<(String, &ArrowDataType)> = schema
123        .fields()
124        .iter()
125        .map(|f| (f.name().clone(), f.data_type()))
126        .collect();
127
128    let available_names: Vec<String> = arrow_fields.iter().map(|(n, _)| n.clone()).collect();
129
130    // For each catalog column, find its index in the Arrow schema and validate types
131    let mut col_indices: Vec<usize> = Vec::with_capacity(columns.len());
132    for col in columns {
133        let pos = arrow_fields
134            .iter()
135            .position(|(name, _)| name.eq_ignore_ascii_case(&col.name));
136        match pos {
137            Some(idx) => {
138                let (_, arrow_type) = &arrow_fields[idx];
139                validate_type_compatibility(arrow_type, col.logical_type).map_err(|_| {
140                    ParquetReaderError::TypeMismatch {
141                        column_name: col.name.clone(),
142                        arrow_type: format!("{arrow_type:?}"),
143                        expected_type: format!("{:?}", col.logical_type),
144                    }
145                })?;
146                col_indices.push(idx);
147            }
148            None => {
149                return Err(ParquetReaderError::ColumnNotFound {
150                    column_name: col.name.clone(),
151                    available: available_names.clone(),
152                });
153            }
154        }
155    }
156
157    // Read all row groups and convert
158    let mut results: Vec<Vec<Value>> = Vec::new();
159
160    for batch_result in reader {
161        let batch: RecordBatch = batch_result?;
162        let num_rows = batch.num_rows();
163
164        // Pre-allocate result rows
165        if results.is_empty() {
166            results.reserve(num_rows * 4); // rough initial estimate
167        }
168
169        for row_idx in 0..num_rows {
170            let mut row = Vec::with_capacity(columns.len());
171            for (catalog_idx, &arrow_col_idx) in col_indices.iter().enumerate() {
172                let col = &columns[catalog_idx];
173                let array = batch.column(arrow_col_idx);
174                let value = arrow_array_to_value(array, row_idx, &col.name, col.logical_type, results.len())?;
175                row.push(value);
176            }
177            results.push(row);
178        }
179    }
180
181    Ok(results)
182}
183
184/// A streaming parquet reader that yields batches of rows on demand.
185///
186/// Unlike `read_parquet`, this avoids materializing the entire file into
187/// `Vec<Vec<Value>>` at once. Each call to `next()` reads and converts one
188/// Arrow `RecordBatch`.
189pub struct ParquetStreamReader {
190    reader: parquet::arrow::arrow_reader::ParquetRecordBatchReader,
191    col_indices: Vec<usize>,
192    columns: Vec<CatalogColumn>,
193}
194
195impl Iterator for ParquetStreamReader {
196    type Item = ParquetResult<Vec<Vec<Value>>>;
197
198    fn next(&mut self) -> Option<Self::Item> {
199        match self.reader.next() {
200            Some(Ok(batch)) => {
201                let num_rows = batch.num_rows();
202                let mut rows = Vec::with_capacity(num_rows);
203                for row_idx in 0..num_rows {
204                    let mut row = Vec::with_capacity(self.columns.len());
205                    for (catalog_idx, &arrow_col_idx) in self.col_indices.iter().enumerate() {
206                        let col = &self.columns[catalog_idx];
207                        let array = batch.column(arrow_col_idx);
208                        let value = match arrow_array_to_value(array, row_idx, &col.name, col.logical_type, rows.len())
209                        {
210                            Ok(v) => v,
211                            Err(e) => return Some(Err(e)),
212                        };
213                        row.push(value);
214                    }
215                    rows.push(row);
216                }
217                Some(Ok(rows))
218            }
219            Some(Err(e)) => Some(Err(ParquetReaderError::ParquetError(e.to_string()))),
220            None => None,
221        }
222    }
223}
224
225/// Open a Parquet file and return a streaming reader that yields row batches.
226///
227/// The file is read into memory (required by the Parquet format for footer
228/// access), but rows are converted to `Vec<Value>` per batch rather than
229/// materializing the entire dataset at once.
230pub fn stream_parquet(
231    path: &str,
232    vfs: &akar_common::file_system::VirtualFileSystemRegistry,
233    columns: &[CatalogColumn],
234) -> ParquetResult<ParquetStreamReader> {
235    let mut file = vfs
236        .open_read(path)
237        .map_err(|e| ParquetReaderError::ParquetError(format!("Cannot open file: {e}")))?;
238
239    let mut buffer = Vec::new();
240    file.read_to_end(&mut buffer)
241        .map_err(|e| ParquetReaderError::ParquetError(format!("Read error: {e}")))?;
242    let bytes = bytes::Bytes::from(buffer);
243
244    let builder = parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder::try_new(bytes)?;
245    let schema = builder.schema().clone();
246    let reader = builder.build()?;
247
248    let arrow_fields: Vec<(String, &arrow::datatypes::DataType)> = schema
249        .fields()
250        .iter()
251        .map(|f| (f.name().clone(), f.data_type()))
252        .collect();
253
254    let available_names: Vec<String> = arrow_fields.iter().map(|(n, _)| n.clone()).collect();
255
256    let mut col_indices: Vec<usize> = Vec::with_capacity(columns.len());
257    for col in columns {
258        let pos = arrow_fields
259            .iter()
260            .position(|(name, _)| name.eq_ignore_ascii_case(&col.name));
261        match pos {
262            Some(idx) => {
263                let (_, arrow_type) = &arrow_fields[idx];
264                validate_type_compatibility(arrow_type, col.logical_type).map_err(|_| {
265                    ParquetReaderError::TypeMismatch {
266                        column_name: col.name.clone(),
267                        arrow_type: format!("{arrow_type:?}"),
268                        expected_type: format!("{:?}", col.logical_type),
269                    }
270                })?;
271                col_indices.push(idx);
272            }
273            None => {
274                return Err(ParquetReaderError::ColumnNotFound {
275                    column_name: col.name.clone(),
276                    available: available_names.clone(),
277                });
278            }
279        }
280    }
281
282    Ok(ParquetStreamReader {
283        reader,
284        col_indices,
285        columns: columns.to_vec(),
286    })
287}
288
289// ─── Type validation ────────────────────────────────────────────────────────────
290
291/// Check that an Arrow `DataType` is compatible with the expected Akar `LogicalTypeID`.
292fn validate_type_compatibility(arrow_type: &ArrowDataType, expected: LogicalTypeID) -> Result<(), ()> {
293    match (arrow_type, expected) {
294        (ArrowDataType::Boolean, LogicalTypeID::Bool) => Ok(()),
295        (ArrowDataType::Int8, LogicalTypeID::Int8) => Ok(()),
296        (ArrowDataType::Int16, LogicalTypeID::Int16) => Ok(()),
297        (ArrowDataType::Int32, LogicalTypeID::Int32) => Ok(()),
298        (ArrowDataType::Int64, LogicalTypeID::Int64 | LogicalTypeID::Serial) => Ok(()),
299        (ArrowDataType::UInt8, LogicalTypeID::UInt8) => Ok(()),
300        (ArrowDataType::UInt16, LogicalTypeID::UInt16) => Ok(()),
301        (ArrowDataType::UInt32, LogicalTypeID::UInt32) => Ok(()),
302        (ArrowDataType::UInt64, LogicalTypeID::UInt64) => Ok(()),
303        (ArrowDataType::Float32, LogicalTypeID::Float) => Ok(()),
304        (ArrowDataType::Float64, LogicalTypeID::Double) => Ok(()),
305        (ArrowDataType::Utf8 | ArrowDataType::LargeUtf8, LogicalTypeID::String) => Ok(()),
306        (ArrowDataType::Binary | ArrowDataType::LargeBinary, LogicalTypeID::Blob) => Ok(()),
307        (ArrowDataType::Date32 | ArrowDataType::Date64, LogicalTypeID::Date) => Ok(()),
308        (ArrowDataType::Timestamp(TimeUnit::Second, _), LogicalTypeID::TimestampSec) => Ok(()),
309        (ArrowDataType::Timestamp(TimeUnit::Millisecond, _), LogicalTypeID::TimestampMs) => Ok(()),
310        (ArrowDataType::Timestamp(TimeUnit::Microsecond, _), LogicalTypeID::Timestamp) => Ok(()),
311        (ArrowDataType::Timestamp(TimeUnit::Nanosecond, _), LogicalTypeID::TimestampNs) => Ok(()),
312        (ArrowDataType::Duration(_), LogicalTypeID::Interval) => Ok(()),
313        (ArrowDataType::List(_), LogicalTypeID::List) => Ok(()),
314        (ArrowDataType::Struct(_), LogicalTypeID::Struct) => Ok(()),
315        (ArrowDataType::Map(_, _), LogicalTypeID::Map) => Ok(()),
316        // Allow numeric widening: smaller ints can be read as larger targets
317        (ArrowDataType::Int8, LogicalTypeID::Int64 | LogicalTypeID::Int32 | LogicalTypeID::Int16) => Ok(()),
318        (ArrowDataType::Int16, LogicalTypeID::Int64 | LogicalTypeID::Int32) => Ok(()),
319        (ArrowDataType::Int32, LogicalTypeID::Int64) => Ok(()),
320        (ArrowDataType::UInt8, LogicalTypeID::UInt64 | LogicalTypeID::UInt32 | LogicalTypeID::UInt16) => Ok(()),
321        (ArrowDataType::UInt16, LogicalTypeID::UInt64 | LogicalTypeID::UInt32) => Ok(()),
322        (ArrowDataType::UInt32, LogicalTypeID::UInt64) => Ok(()),
323        (ArrowDataType::Float32, LogicalTypeID::Double) => Ok(()),
324        // Fallback: accept String target for any Arrow type (coercion done during conversion)
325        (_, LogicalTypeID::String) => Ok(()),
326        _ => Err(()),
327    }
328}
329
330// ─── Array → Value conversion ───────────────────────────────────────────────────
331
332/// Convert a value from an Arrow array at the given row index to a Akar `Value`.
333fn arrow_array_to_value(
334    array: &dyn Array,
335    row: usize,
336    column_name: &str,
337    target_type: LogicalTypeID,
338    _global_row: usize,
339) -> ParquetResult<Value> {
340    // Handle nulls
341    if array.is_null(row) {
342        return Ok(Value::Null);
343    }
344
345    match target_type {
346        LogicalTypeID::Bool => {
347            let arr = downcast::<BooleanArray>(array, column_name)?;
348            Ok(Value::Bool(arr.value(row)))
349        }
350        LogicalTypeID::Int64 | LogicalTypeID::Serial => {
351            let val = cast_int_to_i64(array, row, column_name)?;
352            Ok(Value::Int64(val))
353        }
354        LogicalTypeID::Int32 => {
355            let val = cast_int_to_i64(array, row, column_name)?;
356            Ok(Value::Int32(val as i32))
357        }
358        LogicalTypeID::Int16 => {
359            let val = cast_int_to_i64(array, row, column_name)?;
360            Ok(Value::Int16(val as i16))
361        }
362        LogicalTypeID::Int8 => {
363            let val = cast_int_to_i64(array, row, column_name)?;
364            Ok(Value::Int8(val as i8))
365        }
366        LogicalTypeID::UInt64 => {
367            let val = cast_uint_to_u64(array, row, column_name)?;
368            Ok(Value::UInt64(val))
369        }
370        LogicalTypeID::UInt32 => {
371            let val = cast_uint_to_u64(array, row, column_name)?;
372            Ok(Value::UInt32(val as u32))
373        }
374        LogicalTypeID::UInt16 => {
375            let val = cast_uint_to_u64(array, row, column_name)?;
376            Ok(Value::UInt16(val as u16))
377        }
378        LogicalTypeID::UInt8 => {
379            let val = cast_uint_to_u64(array, row, column_name)?;
380            Ok(Value::UInt8(val as u8))
381        }
382        LogicalTypeID::Double => {
383            let val = cast_to_f64(array, row, column_name)?;
384            Ok(Value::Double(val))
385        }
386        LogicalTypeID::Float => {
387            let val = cast_to_f64(array, row, column_name)?;
388            Ok(Value::Float(val as f32))
389        }
390        LogicalTypeID::String => {
391            let s = array_to_string(array, row, column_name)?;
392            Ok(Value::String(s))
393        }
394        LogicalTypeID::Blob => {
395            let arr = downcast::<BinaryArray>(array, column_name)?;
396            Ok(Value::Blob(arr.value(row).to_vec()))
397        }
398        LogicalTypeID::Date => {
399            let val = cast_date_to_days(array, row, column_name)?;
400            Ok(Value::Date(Date::from_days_since_epoch(val)))
401        }
402        LogicalTypeID::Timestamp => {
403            let micros = cast_timestamp_to_micros(array, row, column_name)?;
404            Ok(Value::Timestamp(Timestamp::from_micros_since_epoch(micros)))
405        }
406        LogicalTypeID::TimestampMs => {
407            let micros = cast_timestamp_to_micros(array, row, column_name)?;
408            Ok(Value::TimestampMs(Timestamp::from_micros_since_epoch(micros)))
409        }
410        LogicalTypeID::TimestampSec => {
411            let micros = cast_timestamp_to_micros(array, row, column_name)?;
412            Ok(Value::TimestampSec(Timestamp(micros / 1_000_000)))
413        }
414        LogicalTypeID::TimestampNs => {
415            let micros = cast_timestamp_to_micros(array, row, column_name)?;
416            Ok(Value::TimestampNs(Timestamp(micros * 1000)))
417        }
418        LogicalTypeID::TimestampTz => {
419            let micros = cast_timestamp_to_micros(array, row, column_name)?;
420            Ok(Value::TimestampTz(akar_common::types::TimestampTZ(micros)))
421        }
422        LogicalTypeID::Interval => {
423            let arr = downcast::<DurationMicrosecondArray>(array, column_name)?;
424            Ok(Value::Interval(Interval::new(0, 0, arr.value(row))))
425        }
426        LogicalTypeID::List => {
427            let vals = array_list_to_values(array, row, column_name)?;
428            Ok(Value::List(vals))
429        }
430        LogicalTypeID::Struct => {
431            let vals = array_struct_to_values(array, row, column_name)?;
432            Ok(Value::Struct(vals))
433        }
434        LogicalTypeID::Map => {
435            let vals = array_map_to_values(array, row, column_name)?;
436            Ok(Value::Map(vals))
437        }
438        // Fallback: string representation
439        _ => {
440            let s = array_to_string(array, row, column_name)?;
441            Ok(Value::String(s))
442        }
443    }
444}
445
446// ─── Downcast helper ────────────────────────────────────────────────────────────
447
448fn downcast<'a, T: Array + 'static>(array: &'a dyn Array, column_name: &str) -> ParquetResult<&'a T> {
449    array
450        .as_any()
451        .downcast_ref::<T>()
452        .ok_or_else(|| ParquetReaderError::ConversionError {
453            column_name: column_name.to_string(),
454            row: 0,
455            message: format!(
456                "expected array type {} but got {:?}",
457                std::any::type_name::<T>(),
458                array.data_type()
459            ),
460        })
461}
462
463// ─── Numeric casting ────────────────────────────────────────────────────────────
464
465/// Extract an i64 from any integer Arrow array (with widening).
466fn cast_int_to_i64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i64> {
467    if let Some(arr) = array.as_any().downcast_ref::<Int8Array>() {
468        return Ok(arr.value(row) as i64);
469    }
470    if let Some(arr) = array.as_any().downcast_ref::<Int16Array>() {
471        return Ok(arr.value(row) as i64);
472    }
473    if let Some(arr) = array.as_any().downcast_ref::<Int32Array>() {
474        return Ok(arr.value(row) as i64);
475    }
476    if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
477        return Ok(arr.value(row));
478    }
479    Err(ParquetReaderError::ConversionError {
480        column_name: column_name.to_string(),
481        row,
482        message: format!("cannot cast {:?} to Int64", array.data_type()),
483    })
484}
485
486/// Extract a u64 from any unsigned integer Arrow array.
487fn cast_uint_to_u64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<u64> {
488    if let Some(arr) = array.as_any().downcast_ref::<UInt8Array>() {
489        return Ok(arr.value(row) as u64);
490    }
491    if let Some(arr) = array.as_any().downcast_ref::<UInt16Array>() {
492        return Ok(arr.value(row) as u64);
493    }
494    if let Some(arr) = array.as_any().downcast_ref::<UInt32Array>() {
495        return Ok(arr.value(row) as u64);
496    }
497    if let Some(arr) = array.as_any().downcast_ref::<UInt64Array>() {
498        return Ok(arr.value(row));
499    }
500    Err(ParquetReaderError::ConversionError {
501        column_name: column_name.to_string(),
502        row,
503        message: format!("cannot cast {:?} to UInt64", array.data_type()),
504    })
505}
506
507/// Extract an f64 from float or integer Arrow arrays.
508fn cast_to_f64(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<f64> {
509    if let Some(arr) = array.as_any().downcast_ref::<Float32Array>() {
510        return Ok(arr.value(row) as f64);
511    }
512    if let Some(arr) = array.as_any().downcast_ref::<Float64Array>() {
513        return Ok(arr.value(row));
514    }
515    if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
516        return Ok(arr.value(row) as f64);
517    }
518    if let Some(arr) = array.as_any().downcast_ref::<Int32Array>() {
519        return Ok(arr.value(row) as f64);
520    }
521    Err(ParquetReaderError::ConversionError {
522        column_name: column_name.to_string(),
523        row,
524        message: format!("cannot cast {:?} to Float64", array.data_type()),
525    })
526}
527
528/// Extract a string from various Arrow array types.
529fn array_to_string(array: &dyn Array, row: usize, _column_name: &str) -> ParquetResult<String> {
530    if let Some(arr) = array.as_any().downcast_ref::<StringArray>() {
531        return Ok(arr.value(row).to_string());
532    }
533    if let Some(arr) = array.as_any().downcast_ref::<LargeStringArray>() {
534        return Ok(arr.value(row).to_string());
535    }
536    if let Some(arr) = array.as_any().downcast_ref::<BinaryArray>() {
537        return Ok(String::from_utf8_lossy(arr.value(row)).to_string());
538    }
539    if let Some(arr) = array.as_any().downcast_ref::<Int64Array>() {
540        return Ok(arr.value(row).to_string());
541    }
542    if let Some(arr) = array.as_any().downcast_ref::<Float64Array>() {
543        return Ok(arr.value(row).to_string());
544    }
545    if let Some(arr) = array.as_any().downcast_ref::<BooleanArray>() {
546        return Ok(arr.value(row).to_string());
547    }
548    // Fallback: use Debug formatting
549    Ok(format!("{:?}", array))
550}
551
552// ─── Date/timestamp casting ─────────────────────────────────────────────────────
553
554/// Extract days since epoch from Date32/Date64 Arrow arrays.
555fn cast_date_to_days(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i32> {
556    if let Some(arr) = array.as_any().downcast_ref::<Date32Array>() {
557        return Ok(arr.value(row));
558    }
559    if let Some(arr) = array.as_any().downcast_ref::<Date64Array>() {
560        // Date64 is milliseconds since epoch
561        return Ok((arr.value(row) / 86_400_000) as i32);
562    }
563    Err(ParquetReaderError::ConversionError {
564        column_name: column_name.to_string(),
565        row,
566        message: format!("cannot cast {:?} to Date", array.data_type()),
567    })
568}
569
570/// Extract microseconds since epoch from Timestamp Arrow arrays.
571fn cast_timestamp_to_micros(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<i64> {
572    if let Some(arr) = array.as_any().downcast_ref::<TimestampSecondArray>() {
573        return Ok(arr.value(row) * 1_000_000);
574    }
575    if let Some(arr) = array.as_any().downcast_ref::<TimestampMillisecondArray>() {
576        return Ok(arr.value(row) * 1_000);
577    }
578    if let Some(arr) = array.as_any().downcast_ref::<TimestampMicrosecondArray>() {
579        return Ok(arr.value(row));
580    }
581    if let Some(arr) = array.as_any().downcast_ref::<TimestampNanosecondArray>() {
582        return Ok(arr.value(row) / 1_000);
583    }
584    Err(ParquetReaderError::ConversionError {
585        column_name: column_name.to_string(),
586        row,
587        message: format!("cannot cast {:?} to Timestamp", array.data_type()),
588    })
589}
590
591// ─── Complex type helpers ───────────────────────────────────────────────────────
592
593/// Convert a List array entry to `Vec<Value>`, preserving the element type
594/// (Float64 → `Value::Double`, Int64 → `Value::Int64`, ...) so FLOAT[]
595/// embeddings round-trip through parquet (P53.37). The previous String
596/// fallback turned every element into `Value::String`, which the FLOAT[]
597/// column could not coerce.
598fn array_list_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<Value>> {
599    let list_arr = downcast::<ListArray>(array, column_name)?;
600    let values = list_arr.value(row);
601    let mut result = Vec::with_capacity(values.len());
602    for i in 0..values.len() {
603        if values.is_null(i) {
604            result.push(Value::Null);
605        } else if let Some(f) = values.as_any().downcast_ref::<Float64Array>() {
606            result.push(Value::Double(f.value(i)));
607        } else if let Some(f) = values.as_any().downcast_ref::<Float32Array>() {
608            result.push(Value::Float(f.value(i)));
609        } else if let Some(a) = values.as_any().downcast_ref::<Int64Array>() {
610            result.push(Value::Int64(a.value(i)));
611        } else if let Some(a) = values.as_any().downcast_ref::<Int32Array>() {
612            result.push(Value::Int32(a.value(i)));
613        } else if let Some(a) = values.as_any().downcast_ref::<Int16Array>() {
614            result.push(Value::Int16(a.value(i)));
615        } else if let Some(a) = values.as_any().downcast_ref::<Int8Array>() {
616            result.push(Value::Int8(a.value(i)));
617        } else if let Some(a) = values.as_any().downcast_ref::<BooleanArray>() {
618            result.push(Value::Bool(a.value(i)));
619        } else if let Some(s) = values.as_any().downcast_ref::<StringArray>() {
620            result.push(Value::String(s.value(i).to_string()));
621        } else {
622            // Unknown element type — generic string fallback.
623            let s = array_to_string(&values, i, column_name)?;
624            result.push(Value::String(s));
625        }
626    }
627    Ok(result)
628}
629
630/// Convert a Struct array entry to `Vec<(String, Value)>`.
631fn array_struct_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<(String, Value)>> {
632    let struct_arr = downcast::<StructArray>(array, column_name)?;
633    let mut result = Vec::with_capacity(struct_arr.num_columns());
634    for col_idx in 0..struct_arr.num_columns() {
635        let field = struct_arr.column(col_idx);
636        let field_name = struct_arr
637            .fields()
638            .get(col_idx)
639            .map(|f| f.name().clone())
640            .unwrap_or_else(|| format!("_{col_idx}"));
641        let val = if field.is_null(row) {
642            Value::Null
643        } else {
644            Value::String(array_to_string(field.as_ref(), row, column_name)?)
645        };
646        result.push((field_name, val));
647    }
648    Ok(result)
649}
650
651/// Convert a Map array entry to `Vec<(Value, Value)>`.
652fn array_map_to_values(array: &dyn Array, row: usize, column_name: &str) -> ParquetResult<Vec<(Value, Value)>> {
653    let map_arr = downcast::<MapArray>(array, column_name)?;
654    let entries = map_arr.value(row);
655    let keys = map_arr.keys();
656    let values = map_arr.values();
657
658    let mut result = Vec::new();
659    if let Some(_entries_struct) = entries.as_any().downcast_ref::<StructArray>() {
660        // Map entries are stored as a list of structs {key, value}
661        for i in 0..entries.len() {
662            let key_val = if keys.is_null(i) {
663                Value::Null
664            } else {
665                Value::String(array_to_string(keys, i, column_name)?)
666            };
667            let val_val = if values.is_null(i) {
668                Value::Null
669            } else {
670                Value::String(array_to_string(values, i, column_name)?)
671            };
672            result.push((key_val, val_val));
673        }
674    }
675    Ok(result)
676}
677
678// ─── Tests ──────────────────────────────────────────────────────────────────────
679
680#[cfg(test)]
681mod tests {
682    use super::*;
683    use arrow::datatypes::{DataType, Field, Schema};
684    use std::sync::Arc;
685
686    /// Write a RecordBatch to a temporary parquet file and return the path.
687    fn write_parquet_batch(dir: &tempfile::TempDir, filename: &str, batch: &RecordBatch) -> std::path::PathBuf {
688        let path = dir.path().join(filename);
689        let file = std::fs::File::create(&path).unwrap();
690        let schema = batch.schema();
691        let mut writer = parquet::arrow::ArrowWriter::try_new(file, schema, None).unwrap();
692        writer.write(batch).unwrap();
693        writer.close().unwrap();
694        path
695    }
696
697    fn test_schema() -> Vec<CatalogColumn> {
698        vec![
699            CatalogColumn {
700                compression: akar_common::enums::CompressionType::Uncompressed,
701                name: "name".into(),
702                logical_type: LogicalTypeID::String,
703                is_primary_key: true,
704                default_value: None,
705            },
706            CatalogColumn {
707                compression: akar_common::enums::CompressionType::Uncompressed,
708                name: "age".into(),
709                logical_type: LogicalTypeID::Int64,
710                is_primary_key: false,
711                default_value: None,
712            },
713            CatalogColumn {
714                compression: akar_common::enums::CompressionType::Uncompressed,
715                name: "score".into(),
716                logical_type: LogicalTypeID::Double,
717                is_primary_key: false,
718                default_value: None,
719            },
720            CatalogColumn {
721                compression: akar_common::enums::CompressionType::Uncompressed,
722                name: "active".into(),
723                logical_type: LogicalTypeID::Bool,
724                is_primary_key: false,
725                default_value: None,
726            },
727        ]
728    }
729
730    #[test]
731    fn test_validate_type_compatibility() {
732        assert!(validate_type_compatibility(&ArrowDataType::Int64, LogicalTypeID::Int64).is_ok());
733        assert!(validate_type_compatibility(&ArrowDataType::Utf8, LogicalTypeID::String).is_ok());
734        assert!(validate_type_compatibility(&ArrowDataType::Boolean, LogicalTypeID::Bool).is_ok());
735        assert!(validate_type_compatibility(&ArrowDataType::Float64, LogicalTypeID::Double).is_ok());
736        // Widening
737        assert!(validate_type_compatibility(&ArrowDataType::Int32, LogicalTypeID::Int64).is_ok());
738        assert!(validate_type_compatibility(&ArrowDataType::Int8, LogicalTypeID::Int64).is_ok());
739        // Mismatch
740        assert!(validate_type_compatibility(&ArrowDataType::Int64, LogicalTypeID::String).is_ok()); // String accepts anything
741        assert!(validate_type_compatibility(&ArrowDataType::Boolean, LogicalTypeID::Int64).is_err());
742    }
743
744    #[test]
745    fn test_cast_int_to_i64() {
746        let i32_arr = Int32Array::from(vec![42, -1]);
747        assert_eq!(cast_int_to_i64(&i32_arr, 0, "col").unwrap(), 42i64);
748        assert_eq!(cast_int_to_i64(&i32_arr, 1, "col").unwrap(), -1i64);
749
750        let i64_arr = Int64Array::from(vec![999_999_999_999i64]);
751        assert_eq!(cast_int_to_i64(&i64_arr, 0, "col").unwrap(), 999_999_999_999i64);
752    }
753
754    #[test]
755    fn test_cast_to_f64() {
756        let f64_arr = Float64Array::from(vec![3.15]);
757        assert!((cast_to_f64(&f64_arr, 0, "col").unwrap() - 3.15).abs() < 1e-10);
758
759        let f32_arr = Float32Array::from(vec![2.5f32]);
760        assert!((cast_to_f64(&f32_arr, 0, "col").unwrap() - 2.5).abs() < 1e-10);
761    }
762
763    #[test]
764    fn test_array_to_string() {
765        let str_arr = StringArray::from(vec!["hello"]);
766        assert_eq!(array_to_string(&str_arr, 0, "col").unwrap(), "hello");
767
768        let int_arr = Int64Array::from(vec![42]);
769        assert_eq!(array_to_string(&int_arr, 0, "col").unwrap(), "42");
770    }
771
772    #[test]
773    fn test_cast_date_to_days() {
774        let date_arr = Date32Array::from(vec![0i32, 19723i32]); // 1970-01-01, 2024-01-01
775        assert_eq!(cast_date_to_days(&date_arr, 0, "col").unwrap(), 0);
776        assert_eq!(cast_date_to_days(&date_arr, 1, "col").unwrap(), 19723);
777    }
778
779    #[test]
780    fn test_cast_timestamp_to_micros() {
781        use arrow::array::TimestampMicrosecondArray;
782        let ts_arr = TimestampMicrosecondArray::from(vec![1_700_000_000_000_000i64]);
783        let micros = cast_timestamp_to_micros(&ts_arr, 0, "col").unwrap();
784        assert_eq!(micros, 1_700_000_000_000_000i64);
785    }
786
787    #[test]
788    fn test_column_not_found() {
789        let dir = tempfile::tempdir().unwrap();
790        // Create a parquet file with one valid column
791        let schema = Arc::new(Schema::new(vec![
792            Field::new("x", DataType::Int64, false),
793            Field::new("y", DataType::Utf8, false),
794        ]));
795        let batch = RecordBatch::try_new(
796            schema.clone(),
797            vec![
798                Arc::new(Int64Array::from(vec![1])),
799                Arc::new(StringArray::from(vec!["a"])),
800            ],
801        )
802        .unwrap();
803        let path = write_parquet_batch(&dir, "test.parquet", &batch);
804
805        let columns = vec![CatalogColumn {
806            compression: akar_common::enums::CompressionType::Uncompressed,
807            name: "missing_col".into(),
808            logical_type: LogicalTypeID::Int64,
809            is_primary_key: false,
810            default_value: None,
811        }];
812
813        let result = read_parquet(
814            path.to_str().unwrap(),
815            &akar_common::file_system::VirtualFileSystemRegistry::new(),
816            &columns,
817        );
818        assert!(result.is_err());
819        match result.unwrap_err() {
820            ParquetReaderError::ColumnNotFound { column_name, .. } => {
821                assert_eq!(column_name, "missing_col");
822            }
823            e => panic!("Expected ColumnNotFound, got: {e}"),
824        }
825    }
826
827    #[test]
828    fn test_type_mismatch() {
829        let dir = tempfile::tempdir().unwrap();
830        let schema = Arc::new(Schema::new(vec![Field::new("val", DataType::Boolean, false)]));
831        let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(BooleanArray::from(vec![true]))]).unwrap();
832        let path = write_parquet_batch(&dir, "mismatch.parquet", &batch);
833
834        let columns = vec![CatalogColumn {
835            compression: akar_common::enums::CompressionType::Uncompressed,
836            name: "val".into(),
837            logical_type: LogicalTypeID::Int64,
838            is_primary_key: false,
839            default_value: None,
840        }];
841
842        let result = read_parquet(
843            path.to_str().unwrap(),
844            &akar_common::file_system::VirtualFileSystemRegistry::new(),
845            &columns,
846        );
847        assert!(result.is_err());
848        match result.unwrap_err() {
849            ParquetReaderError::TypeMismatch { .. } => {} // expected
850            e => panic!("Expected TypeMismatch, got: {e}"),
851        }
852    }
853
854    #[test]
855    fn test_file_not_found() {
856        let result = read_parquet(
857            "nonexistent.parquet",
858            &akar_common::file_system::VirtualFileSystemRegistry::new(),
859            &test_schema(),
860        );
861        assert!(result.is_err());
862        match result.unwrap_err() {
863            ParquetReaderError::ParquetError(_) => {} // expected
864            _ => panic!("Expected ParquetError"),
865        }
866    }
867
868    #[test]
869    fn test_arrow_array_to_value_basic() {
870        // Bool
871        let bool_arr = BooleanArray::from(vec![Some(true), None, Some(false)]);
872        assert_eq!(
873            arrow_array_to_value(&bool_arr, 0, "b", LogicalTypeID::Bool, 0).unwrap(),
874            Value::Bool(true)
875        );
876        assert_eq!(
877            arrow_array_to_value(&bool_arr, 1, "b", LogicalTypeID::Bool, 0).unwrap(),
878            Value::Null
879        );
880
881        // Int64
882        let int_arr = Int64Array::from(vec![Some(42), None]);
883        assert_eq!(
884            arrow_array_to_value(&int_arr, 0, "i", LogicalTypeID::Int64, 0).unwrap(),
885            Value::Int64(42)
886        );
887        assert_eq!(
888            arrow_array_to_value(&int_arr, 1, "i", LogicalTypeID::Int64, 0).unwrap(),
889            Value::Null
890        );
891
892        // String
893        let str_arr = StringArray::from(vec![Some("hello"), None]);
894        assert_eq!(
895            arrow_array_to_value(&str_arr, 0, "s", LogicalTypeID::String, 0).unwrap(),
896            Value::String("hello".into())
897        );
898        assert_eq!(
899            arrow_array_to_value(&str_arr, 1, "s", LogicalTypeID::String, 0).unwrap(),
900            Value::Null
901        );
902
903        // Double
904        let f64_arr = Float64Array::from(vec![Some(3.15), None]);
905        assert_eq!(
906            arrow_array_to_value(&f64_arr, 0, "d", LogicalTypeID::Double, 0).unwrap(),
907            Value::Double(3.15)
908        );
909    }
910
911    #[test]
912    fn test_arrow_array_widening() {
913        // Int32 → Int64 widening
914        let i32_arr = Int32Array::from(vec![100]);
915        assert_eq!(
916            arrow_array_to_value(&i32_arr, 0, "i", LogicalTypeID::Int64, 0).unwrap(),
917            Value::Int64(100)
918        );
919
920        // Float32 → Double widening
921        let f32_arr = Float32Array::from(vec![2.5f32]);
922        let val = arrow_array_to_value(&f32_arr, 0, "f", LogicalTypeID::Double, 0).unwrap();
923        if let Value::Double(d) = val {
924            assert!((d - 2.5).abs() < 1e-6);
925        } else {
926            panic!("Expected Double");
927        }
928    }
929
930    #[test]
931    fn test_unsigned_int_types() {
932        let u8_arr = UInt8Array::from(vec![200u8]);
933        assert_eq!(
934            arrow_array_to_value(&u8_arr, 0, "u", LogicalTypeID::UInt8, 0).unwrap(),
935            Value::UInt8(200)
936        );
937
938        let u32_arr = UInt32Array::from(vec![100000u32]);
939        assert_eq!(
940            arrow_array_to_value(&u32_arr, 0, "u", LogicalTypeID::UInt32, 0).unwrap(),
941            Value::UInt32(100000)
942        );
943
944        let u64_arr = UInt64Array::from(vec![u64::MAX]);
945        assert_eq!(
946            arrow_array_to_value(&u64_arr, 0, "u", LogicalTypeID::UInt64, 0).unwrap(),
947            Value::UInt64(u64::MAX)
948        );
949    }
950
951    #[test]
952    fn test_date_conversion() {
953        let date_arr = Date32Array::from(vec![7439i32]); // 1990-05-15
954        let val = arrow_array_to_value(&date_arr, 0, "d", LogicalTypeID::Date, 0).unwrap();
955        if let Value::Date(d) = val {
956            assert_eq!(d.days_since_epoch(), 7439);
957        } else {
958            panic!("Expected Date");
959        }
960    }
961
962    #[test]
963    fn test_timestamp_conversion() {
964        use arrow::array::TimestampMicrosecondArray;
965        let ts_arr = TimestampMicrosecondArray::from(vec![1_704_198_600_000_000i64]); // 2024-01-02 12:30:00.000000
966        let val = arrow_array_to_value(&ts_arr, 0, "t", LogicalTypeID::Timestamp, 0).unwrap();
967        if let Value::Timestamp(ts) = val {
968            assert_eq!(ts.micros_since_epoch(), 1_704_198_600_000_000);
969        } else {
970            panic!("Expected Timestamp");
971        }
972    }
973
974    #[test]
975    fn test_blob_conversion() {
976        let blob_arr = BinaryArray::from(vec![&b"hello"[..]]);
977        let val = arrow_array_to_value(&blob_arr, 0, "b", LogicalTypeID::Blob, 0).unwrap();
978        assert_eq!(val, Value::Blob(b"hello".to_vec()));
979    }
980
981    #[test]
982    fn test_round_trip_parquet() {
983        let dir = tempfile::tempdir().unwrap();
984
985        // Write a parquet file with known data
986        let schema = Arc::new(Schema::new(vec![
987            Field::new("name", DataType::Utf8, false),
988            Field::new("age", DataType::Int64, false),
989            Field::new("score", DataType::Float64, false),
990            Field::new("active", DataType::Boolean, false),
991        ]));
992
993        let names = StringArray::from(vec!["Alice", "Bob", "Charlie"]);
994        let ages = Int64Array::from(vec![30, 25, 35]);
995        let scores = Float64Array::from(vec![95.5, 87.3, 91.2]);
996        let actives = BooleanArray::from(vec![true, false, true]);
997
998        let batch = RecordBatch::try_new(
999            schema,
1000            vec![Arc::new(names), Arc::new(ages), Arc::new(scores), Arc::new(actives)],
1001        )
1002        .unwrap();
1003
1004        let parquet_path = write_parquet_batch(&dir, "roundtrip.parquet", &batch);
1005
1006        // Read it back using our reader
1007        let columns = test_schema();
1008        let rows = read_parquet(
1009            parquet_path.to_str().unwrap(),
1010            &akar_common::file_system::VirtualFileSystemRegistry::new(),
1011            &columns,
1012        )
1013        .unwrap();
1014
1015        assert_eq!(rows.len(), 3);
1016        assert_eq!(rows[0][0], Value::String("Alice".into()));
1017        assert_eq!(rows[0][1], Value::Int64(30));
1018        assert_eq!(rows[0][2], Value::Double(95.5));
1019        assert_eq!(rows[0][3], Value::Bool(true));
1020
1021        assert_eq!(rows[1][0], Value::String("Bob".into()));
1022        assert_eq!(rows[1][1], Value::Int64(25));
1023        assert_eq!(rows[1][2], Value::Double(87.3));
1024        assert_eq!(rows[1][3], Value::Bool(false));
1025
1026        assert_eq!(rows[2][0], Value::String("Charlie".into()));
1027        assert_eq!(rows[2][1], Value::Int64(35));
1028        assert_eq!(rows[2][2], Value::Double(91.2));
1029        assert_eq!(rows[2][3], Value::Bool(true));
1030    }
1031
1032    #[test]
1033    fn test_round_trip_with_nulls() {
1034        let dir = tempfile::tempdir().unwrap();
1035
1036        let schema = Arc::new(Schema::new(vec![
1037            Field::new("name", DataType::Utf8, true),
1038            Field::new("age", DataType::Int64, true),
1039        ]));
1040
1041        let names = StringArray::from(vec![Some("Alice"), None, Some("Charlie")]);
1042        let ages = Int64Array::from(vec![Some(30), Some(25), None]);
1043
1044        let batch = RecordBatch::try_new(schema, vec![Arc::new(names), Arc::new(ages)]).unwrap();
1045
1046        let parquet_path = write_parquet_batch(&dir, "nulls.parquet", &batch);
1047
1048        let columns = vec![
1049            CatalogColumn {
1050                compression: akar_common::enums::CompressionType::Uncompressed,
1051                name: "name".into(),
1052                logical_type: LogicalTypeID::String,
1053                is_primary_key: false,
1054                default_value: None,
1055            },
1056            CatalogColumn {
1057                compression: akar_common::enums::CompressionType::Uncompressed,
1058                name: "age".into(),
1059                logical_type: LogicalTypeID::Int64,
1060                is_primary_key: false,
1061                default_value: None,
1062            },
1063        ];
1064
1065        let rows = read_parquet(
1066            parquet_path.to_str().unwrap(),
1067            &akar_common::file_system::VirtualFileSystemRegistry::new(),
1068            &columns,
1069        )
1070        .unwrap();
1071        assert_eq!(rows.len(), 3);
1072        assert_eq!(rows[0][0], Value::String("Alice".into()));
1073        assert_eq!(rows[0][1], Value::Int64(30));
1074        assert_eq!(rows[1][0], Value::Null);
1075        assert_eq!(rows[1][1], Value::Int64(25));
1076        assert_eq!(rows[2][0], Value::String("Charlie".into()));
1077        assert_eq!(rows[2][1], Value::Null);
1078    }
1079
1080    #[test]
1081    fn test_round_trip_unsigned_and_floats() {
1082        let dir = tempfile::tempdir().unwrap();
1083
1084        let schema = Arc::new(Schema::new(vec![
1085            Field::new("small", DataType::UInt8, false),
1086            Field::new("medium", DataType::UInt32, false),
1087            Field::new("large", DataType::UInt64, false),
1088            Field::new("temp", DataType::Float32, false),
1089        ]));
1090
1091        let small = UInt8Array::from(vec![100u8, 200u8]);
1092        let medium = UInt32Array::from(vec![1000u32, 50000u32]);
1093        let large = UInt64Array::from(vec![100000u64, u64::MAX]);
1094        let temp = Float32Array::from(vec![36.5f32, 98.6f32]);
1095
1096        let batch = RecordBatch::try_new(
1097            schema,
1098            vec![Arc::new(small), Arc::new(medium), Arc::new(large), Arc::new(temp)],
1099        )
1100        .unwrap();
1101
1102        let parquet_path = write_parquet_batch(&dir, "uints.parquet", &batch);
1103
1104        let columns = vec![
1105            CatalogColumn {
1106                compression: akar_common::enums::CompressionType::Uncompressed,
1107                name: "small".into(),
1108                logical_type: LogicalTypeID::UInt8,
1109                is_primary_key: false,
1110                default_value: None,
1111            },
1112            CatalogColumn {
1113                compression: akar_common::enums::CompressionType::Uncompressed,
1114                name: "medium".into(),
1115                logical_type: LogicalTypeID::UInt32,
1116                is_primary_key: false,
1117                default_value: None,
1118            },
1119            CatalogColumn {
1120                compression: akar_common::enums::CompressionType::Uncompressed,
1121                name: "large".into(),
1122                logical_type: LogicalTypeID::UInt64,
1123                is_primary_key: false,
1124                default_value: None,
1125            },
1126            CatalogColumn {
1127                compression: akar_common::enums::CompressionType::Uncompressed,
1128                name: "temp".into(),
1129                logical_type: LogicalTypeID::Float,
1130                is_primary_key: false,
1131                default_value: None,
1132            },
1133        ];
1134
1135        let rows = read_parquet(
1136            parquet_path.to_str().unwrap(),
1137            &akar_common::file_system::VirtualFileSystemRegistry::new(),
1138            &columns,
1139        )
1140        .unwrap();
1141        assert_eq!(rows.len(), 2);
1142        assert_eq!(rows[0][0], Value::UInt8(100));
1143        assert_eq!(rows[0][1], Value::UInt32(1000));
1144        assert_eq!(rows[0][2], Value::UInt64(100000));
1145        if let Value::Float(f) = rows[0][3] {
1146            assert!((f - 36.5).abs() < 1e-5);
1147        } else {
1148            panic!("Expected Float");
1149        }
1150        assert_eq!(rows[1][0], Value::UInt8(200));
1151        assert_eq!(rows[1][2], Value::UInt64(u64::MAX));
1152    }
1153}