use crate::types::{TypedValue, ValueKind};
use crate::ast::ArgValue;
use std::error::Error;
use std::fmt;
#[derive(Debug)]
pub enum InterchangeError {
ConversionError(String),
UnsupportedType(String),
}
impl fmt::Display for InterchangeError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
InterchangeError::ConversionError(msg) => write!(f, "Conversion error: {}", msg),
InterchangeError::UnsupportedType(msg) => write!(f, "Unsupported type: {}", msg),
}
}
}
impl Error for InterchangeError {}
#[cfg(feature = "duckdb")]
pub fn record_batch_to_table(
batch: duckdb::arrow::record_batch::RecordBatch,
) -> Result<TypedValue, Box<dyn Error>> {
record_batch_to_typed_value(batch)
}
#[cfg(feature = "duckdb")]
pub fn table_to_record_batch(
table: &TypedValue,
) -> Result<duckdb::arrow::record_batch::RecordBatch, Box<dyn Error>> {
typed_value_to_record_batch(table)
}
use duckdb::arrow::record_batch::RecordBatch as DuckDBRecordBatch;
#[cfg(feature = "duckdb")]
pub fn record_batch_to_typed_value(_batch: DuckDBRecordBatch) -> Result<TypedValue, Box<dyn Error>> {
let schema = _batch.schema();
let _num_rows = _batch.num_rows();
let num_cols = _batch.num_columns();
if _num_rows == 0 || num_cols == 0 {
return Ok(TypedValue::new(ValueKind::Table, ArgValue::Table(vec![])));
}
let column_names: Vec<String> = schema.fields().iter().map(|f| f.name().clone()).collect();
let _ = column_names.len();
let mut table_data: Vec<Vec<TypedValue>> = Vec::with_capacity(_num_rows);
for row_idx in 0.._num_rows {
let mut row_data: Vec<TypedValue> = Vec::with_capacity(num_cols);
for col_idx in 0..num_cols {
let column = _batch.column(col_idx);
let typed_val = match duckdb_array_to_typed_value(column.as_ref(), row_idx) {
Ok(val) => val,
Err(e) => {
eprintln!("Error converting value at row {}, col {}: {}", row_idx, col_idx, e);
TypedValue::new(ValueKind::Null, ArgValue::Null)
}
};
row_data.push(typed_val);
}
table_data.push(row_data);
}
Ok(TypedValue::new(ValueKind::Table, ArgValue::Table(table_data)))
}
#[cfg(feature = "duckdb")]
pub fn typed_value_to_record_batch(_value: &TypedValue) -> Result<DuckDBRecordBatch, Box<dyn Error>> {
if !_value.is_table() {
return Err("Only Table type can be converted to RecordBatch".into());
}
let table_data = match &_value.value {
ArgValue::Table(data) => data,
_ => return Err("Expected Table data".into()),
};
if table_data.is_empty() || table_data[0].is_empty() {
use duckdb::arrow::datatypes::{Field, Schema};
use std::sync::Arc;
let schema = Arc::new(Schema::new(Vec::<Field>::new()));
let columns = Vec::new();
return Ok(DuckDBRecordBatch::try_new(schema, columns)?);
}
let num_cols = table_data[0].len();
let mut fields: Vec<duckdb::arrow::datatypes::Field> = Vec::with_capacity(num_cols);
for col_idx in 0..num_cols {
let col_type = infer_column_type(table_data, col_idx);
let field = duckdb::arrow::datatypes::Field::new(
&format!("col_{}", col_idx),
col_type,
true );
fields.push(field);
}
let schema = std::sync::Arc::new(duckdb::arrow::datatypes::Schema::new(fields));
let mut columns: Vec<std::sync::Arc<dyn duckdb::arrow::array::Array>> = Vec::with_capacity(num_cols);
for col_idx in 0..num_cols {
let array = build_column_array(table_data, col_idx)?;
columns.push(array);
}
Ok(DuckDBRecordBatch::try_new(schema, columns)?)
}
#[cfg(feature = "duckdb")]
fn duckdb_array_to_typed_value(array: &dyn duckdb::arrow::array::Array, index: usize) -> Result<TypedValue, Box<dyn Error>> {
use duckdb::arrow::array::*;
use duckdb::arrow::datatypes::DataType;
if array.is_null(index) {
return Ok(TypedValue::new(ValueKind::Null, ArgValue::Null));
}
match array.data_type() {
DataType::Boolean => {
let bool_array = array.as_any().downcast_ref::<BooleanArray>().unwrap();
let value = bool_array.value(index);
Ok(TypedValue::new(ValueKind::Bool, ArgValue::Bool(value)))
},
DataType::Int8 | DataType::Int16 | DataType::Int32 | DataType::Int64 |
DataType::UInt8 | DataType::UInt16 | DataType::UInt32 | DataType::UInt64 |
DataType::Float32 | DataType::Float64 => {
let value = match array.data_type() {
DataType::Int8 => {
let arr = array.as_any().downcast_ref::<Int8Array>().unwrap();
arr.value(index) as f64
},
DataType::Int16 => {
let arr = array.as_any().downcast_ref::<Int16Array>().unwrap();
arr.value(index) as f64
},
DataType::Int32 => {
let arr = array.as_any().downcast_ref::<Int32Array>().unwrap();
arr.value(index) as f64
},
DataType::Int64 => {
let arr = array.as_any().downcast_ref::<Int64Array>().unwrap();
arr.value(index) as f64
},
DataType::UInt8 => {
let arr = array.as_any().downcast_ref::<UInt8Array>().unwrap();
arr.value(index) as f64
},
DataType::UInt16 => {
let arr = array.as_any().downcast_ref::<UInt16Array>().unwrap();
arr.value(index) as f64
},
DataType::UInt32 => {
let arr = array.as_any().downcast_ref::<UInt32Array>().unwrap();
arr.value(index) as f64
},
DataType::UInt64 => {
let arr = array.as_any().downcast_ref::<UInt64Array>().unwrap();
arr.value(index) as f64
},
DataType::Float32 => {
let arr = array.as_any().downcast_ref::<Float32Array>().unwrap();
arr.value(index) as f64
},
DataType::Float64 => {
let arr = array.as_any().downcast_ref::<Float64Array>().unwrap();
arr.value(index)
},
_ => unreachable!(),
};
Ok(TypedValue::new(ValueKind::Number, ArgValue::Number(value)))
},
DataType::Utf8 => {
let string_array = array.as_any().downcast_ref::<StringArray>().unwrap();
let value = string_array.value(index).to_string();
Ok(TypedValue::new(ValueKind::String, ArgValue::String(value)))
},
_ => Err(format!("Unsupported DuckDB array type: {:?}", array.data_type()).into()),
}
}
#[cfg(feature = "duckdb")]
fn infer_column_type(table_data: &[Vec<TypedValue>], col_idx: usize) -> duckdb::arrow::datatypes::DataType {
use duckdb::arrow::datatypes::DataType;
for row in table_data {
if col_idx < row.len() {
let val = &row[col_idx];
if !val.is_null() {
return match val.kind {
ValueKind::Number => DataType::Float64,
ValueKind::String => DataType::Utf8,
ValueKind::Bool => DataType::Boolean,
_ => DataType::Utf8, };
}
}
}
DataType::Utf8 }
#[cfg(feature = "duckdb")]
fn build_column_array(table_data: &[Vec<TypedValue>], col_idx: usize) -> Result<std::sync::Arc<dyn duckdb::arrow::array::Array>, Box<dyn Error>> {
use duckdb::arrow::array::*;
use duckdb::arrow::datatypes::DataType;
let data_type = infer_column_type(table_data, col_idx);
let _num_rows = table_data.len();
match data_type {
DataType::Boolean => {
let mut builder = BooleanBuilder::new();
for row in table_data {
if col_idx < row.len() && row[col_idx].is_bool() {
if let ArgValue::Bool(b) = row[col_idx].value {
builder.append_value(b);
} else {
builder.append_null();
}
} else {
builder.append_null();
}
}
Ok(std::sync::Arc::new(builder.finish()) as std::sync::Arc<dyn duckdb::arrow::array::Array>)
},
DataType::Float64 => {
let mut builder = Float64Builder::new();
for row in table_data {
if col_idx < row.len() && row[col_idx].is_number() {
if let ArgValue::Number(n) = row[col_idx].value {
builder.append_value(n);
} else {
builder.append_null();
}
} else {
builder.append_null();
}
}
Ok(std::sync::Arc::new(builder.finish()) as std::sync::Arc<dyn duckdb::arrow::array::Array>)
},
DataType::Utf8 => {
let mut builder = StringBuilder::new();
for row in table_data {
if col_idx < row.len() && row[col_idx].is_string() {
if let ArgValue::String(s) = &row[col_idx].value {
builder.append_value(s);
} else {
builder.append_null();
}
} else {
builder.append_null();
}
}
Ok(std::sync::Arc::new(builder.finish()) as std::sync::Arc<dyn duckdb::arrow::array::Array>)
},
_ => Err("Unsupported data type for column array".into()),
}
}
#[cfg(feature = "duckdb")]
pub fn duckdb_value_to_typed_value(value: &duckdb::types::Value) -> Result<TypedValue, Box<dyn Error>> {
match value {
duckdb::types::Value::Null => Ok(TypedValue::new(ValueKind::Null, ArgValue::Null)),
duckdb::types::Value::Boolean(b) => Ok(TypedValue::new(ValueKind::Bool, ArgValue::Bool(*b))),
duckdb::types::Value::Int(i) => Ok(TypedValue::new(ValueKind::Number, ArgValue::Number(*i as f64))),
duckdb::types::Value::BigInt(i) => Ok(TypedValue::new(ValueKind::Number, ArgValue::Number(*i as f64))),
duckdb::types::Value::Float(f) => Ok(TypedValue::new(ValueKind::Number, ArgValue::Number(*f as f64))),
duckdb::types::Value::Double(d) => Ok(TypedValue::new(ValueKind::Number, ArgValue::Number(*d))),
duckdb::types::Value::Text(s) => Ok(TypedValue::new(ValueKind::String, ArgValue::String(s.clone()))),
_ => Err(format!("Unsupported DuckDB value type: {:?}", value).into()),
}
}
pub fn extract_schema(_batch: &DuckDBRecordBatch) -> Result<(), Box<dyn Error>> {
Ok(())
}