Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 7 additions & 7 deletions ffi/examples/read-table/read_table.c
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ bool scan_row_callback(
int64_t mod_time,
const Stats* stats,
HandleSharedDvInfo dv_info,
OptionalValue_HandleSharedExpression transform,
OptionalValueHandleSharedExpression transform,
const CStringMap* partition_values)
{
(void)mod_time; // not using this at the moment
Expand All @@ -74,8 +74,8 @@ bool scan_row_callback(
if (selection_vector_res.tag != OkKernelBoolSlice) {
printf("Could not get selection vector from kernel\n");
free_kernel_dv_info(dv_info);
if (transform.tag == OptionalValue_HandleSharedExpression_Tag_Some) {
free_kernel_expression(transform.some._0);
if (transform.tag == SomeHandleSharedExpression) {
free_kernel_expression(transform.some);
}
exit(-1);
}
Expand All @@ -94,17 +94,17 @@ bool scan_row_callback(
print_partition_info(context, partition_values);
#ifdef PRINT_ARROW_DATA
const Expression* transform_expr = NULL;
if (transform.tag == OptionalValue_HandleSharedExpression_Tag_Some) {
transform_expr = (const Expression*)transform.some._0;
if (transform.tag == SomeHandleSharedExpression) {
transform_expr = (const Expression*)transform.some;
}
c_read_parquet_file(context, path, selection_vector, transform_expr);
#endif
free_bool_slice(selection_vector);
context->partition_values = NULL;

free_kernel_dv_info(dv_info);
if (transform.tag == OptionalValue_HandleSharedExpression_Tag_Some) {
free_kernel_expression(transform.some._0);
if (transform.tag == SomeHandleSharedExpression) {
free_kernel_expression(transform.some);
}

return true; // Continue iteration
Expand Down
13 changes: 10 additions & 3 deletions ffi/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -522,7 +522,11 @@ pub unsafe extern "C" fn builder_build(
let allocate_fn = builder_box.allocate_fn;
unsafe {
catch_unwind_into_extern_result(&allocate_fn, move || {
get_default_engine_impl(builder_box.url, builder_box.options, builder_box.allocate_fn)
get_default_engine_impl(
builder_box.url,
builder_box.options,
builder_box.allocate_fn,
)
})
}
}
Expand Down Expand Up @@ -555,8 +559,11 @@ unsafe fn catch_unwind_into_extern_result<T>(
alloc: &dyn AllocateError,
f: impl FnOnce() -> DeltaResult<T>,
) -> ExternResult<T> {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(f))
.unwrap_or_else(|_| Err(Error::generic("delta-kernel-rs panicked across the FFI boundary")));
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(f)).unwrap_or_else(|_| {
Err(Error::generic(
"delta-kernel-rs panicked across the FFI boundary",
))
});
unsafe { result.into_extern_result(alloc) }
}

Expand Down
1 change: 0 additions & 1 deletion ffi/src/scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -624,5 +624,4 @@ mod tests {
let final_map: HashMap<String, String> = *unsafe { Box::from_raw(map_ptr) };
assert_eq!(test_map, final_map);
}

}
11 changes: 7 additions & 4 deletions ffi/src/table_changes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,8 @@ use crate::expressions::kernel_visitor::{unwrap_kernel_predicate, KernelExpressi
use crate::scan::EnginePredicate;
use crate::{
catch_unwind_into_extern_result, kernel_string_slice, unwrap_and_parse_path_as_url,
AllocateStringFn, ExternEngine, ExternResult, IntoExternResult, KernelStringSlice,
NullableCvoid, SharedExternEngine, SharedSchema,
AllocateStringFn, ExternEngine, ExternResult, KernelStringSlice, NullableCvoid,
SharedExternEngine, SharedSchema,
};

#[handle_descriptor(target=TableChanges, mutable=true, sized=true)]
Expand Down Expand Up @@ -471,8 +471,11 @@ mod tests {

pub fn generate_batch_with_id(start_i: i32) -> Result<RecordBatch, ArrowError> {
generate_batch(vec![
("id", vec![start_i, start_i + 1, start_i + 2].into_array()),
("val", vec!["a", "b", "c"].into_array()),
(
"id",
vec![start_i, start_i + 1, start_i + 2].into_arrow_array(),
),
("val", vec!["a", "b", "c"].into_arrow_array()),
])
}

Expand Down
7 changes: 6 additions & 1 deletion kernel/src/actions/visitors.rs
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,7 @@ pub(crate) struct AddVisitor {
}

impl AddVisitor {
#[allow(dead_code)]
#[internal_api]
fn visit_add<'a>(
row_index: usize,
Expand Down Expand Up @@ -141,6 +142,7 @@ impl AddVisitor {
clustering_provider,
})
}
#[allow(dead_code)]
pub(crate) fn names_and_types() -> (&'static [ColumnName], &'static [DataType]) {
static NAMES_AND_TYPES: LazyLock<ColumnNamesAndTypes> =
LazyLock::new(|| Add::to_schema().leaves(ADD_NAME));
Expand Down Expand Up @@ -171,6 +173,7 @@ pub(crate) struct RemoveVisitor {
}

impl RemoveVisitor {
#[allow(dead_code)]
#[internal_api]
pub(crate) fn visit_remove<'a>(
row_index: usize,
Expand Down Expand Up @@ -217,6 +220,7 @@ impl RemoveVisitor {
default_row_commit_version,
})
}
#[allow(dead_code)]
pub(crate) fn names_and_types() -> (&'static [ColumnName], &'static [DataType]) {
static NAMES_AND_TYPES: LazyLock<ColumnNamesAndTypes> =
LazyLock::new(|| Remove::to_schema().leaves(REMOVE_NAME));
Expand Down Expand Up @@ -247,6 +251,7 @@ pub(crate) struct CdcVisitor {
}

impl CdcVisitor {
#[allow(dead_code)]
#[internal_api]
pub(crate) fn visit_cdc<'a>(
row_index: usize,
Expand Down Expand Up @@ -826,7 +831,7 @@ mod tests {
};
let expected = vec![add1, add2, add3];
assert_eq!(add_visitor.adds.len(), expected.len());
for (add, expected) in add_visitor.adds.into_iter().zip(expected.into_iter()) {
for (add, expected) in add_visitor.adds.into_iter().zip(expected) {
assert_eq!(add, expected);
}
}
Expand Down
1 change: 1 addition & 0 deletions kernel/src/engine/parquet_row_group_skipping.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,6 +208,7 @@ impl ParquetStatsProvider for RowGroupFilter<'_> {
// physical name mapping has been performed. Because we currently lack both the
// validation and the name mapping support, we must disable this optimization for the
// time being. See https://github.com/delta-io/delta-kernel-rs/issues/434.
#[allow(unknown_lints, clippy::some_filter)]
return Some(self.get_parquet_rowcount_stat()).filter(|_| false);
};

Expand Down
5 changes: 4 additions & 1 deletion kernel/src/scan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -589,7 +589,10 @@ impl Scan {
&self,
engine: Arc<dyn Engine>,
) -> DeltaResult<impl Iterator<Item = DeltaResult<Box<dyn EngineData>>>> {
fn scan_metadata_callback(batches: &mut Vec<state::ScanFile>, file: state::ScanFile) -> bool {
fn scan_metadata_callback(
batches: &mut Vec<state::ScanFile>,
file: state::ScanFile,
) -> bool {
batches.push(file);
true
}
Expand Down
1 change: 1 addition & 0 deletions kernel/src/schema/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1110,6 +1110,7 @@ impl ColumnNamesAndTypes {
(&self.0, &self.1)
}

#[allow(dead_code)]
pub(crate) fn extend(&mut self, other: ColumnNamesAndTypes) {
self.0.extend(other.0);
self.1.extend(other.1);
Expand Down
30 changes: 15 additions & 15 deletions kernel/tests/read.rs
Original file line number Diff line number Diff line change
Expand Up @@ -198,8 +198,8 @@ async fn stats() -> Result<(), Box<dyn std::error::Error>> {

let batch1 = generate_simple_batch()?;
let batch2 = generate_batch(vec![
("id", vec![5, 7].into_array()),
("val", vec!["e", "g"].into_array()),
("id", vec![5, 7].into_arrow_array()),
("val", vec!["e", "g"].into_arrow_array()),
])?;
let storage = Arc::new(InMemory::new());
// valid commit with min/max (0, 2)
Expand Down Expand Up @@ -996,7 +996,7 @@ fn with_predicate_and_removes() -> Result<(), Box<dyn std::error::Error>> {
#[tokio::test]
async fn predicate_on_non_nullable_partition_column() -> Result<(), Box<dyn std::error::Error>> {
// Test for https://github.com/delta-io/delta-kernel-rs/issues/698
let batch = generate_batch(vec![("val", vec!["a", "b", "c"].into_array())])?;
let batch = generate_batch(vec![("val", vec!["a", "b", "c"].into_arrow_array())])?;

let storage = Arc::new(InMemory::new());
let actions = [
Expand Down Expand Up @@ -1048,8 +1048,8 @@ async fn predicate_on_non_nullable_partition_column() -> Result<(), Box<dyn std:
#[tokio::test]
async fn predicate_on_non_nullable_column_missing_stats() -> Result<(), Box<dyn std::error::Error>>
{
let batch_1 = generate_batch(vec![("val", vec!["a", "b", "c"].into_array())])?;
let batch_2 = generate_batch(vec![("val", vec!["d", "e", "f"].into_array())])?;
let batch_1 = generate_batch(vec![("val", vec!["a", "b", "c"].into_arrow_array())])?;
let batch_2 = generate_batch(vec![("val", vec!["d", "e", "f"].into_arrow_array())])?;

let storage = Arc::new(InMemory::new());
let actions = [
Expand Down Expand Up @@ -1319,16 +1319,16 @@ fn unshredded_variant_table() -> Result<(), Box<dyn std::error::Error>> {
async fn test_row_index_metadata_column() -> Result<(), Box<dyn std::error::Error>> {
// Setup up an in-memory table with different numbers of rows in each file
let batch1 = generate_batch(vec![
("id", vec![1i32, 2, 3, 4, 5].into_array()),
("value", vec!["a", "b", "c", "d", "e"].into_array()),
("id", vec![1i32, 2, 3, 4, 5].into_arrow_array()),
("value", vec!["a", "b", "c", "d", "e"].into_arrow_array()),
])?;
let batch2 = generate_batch(vec![
("id", vec![10i32, 20, 30].into_array()),
("value", vec!["x", "y", "z"].into_array()),
("id", vec![10i32, 20, 30].into_arrow_array()),
("value", vec!["x", "y", "z"].into_arrow_array()),
])?;
let batch3 = generate_batch(vec![
("id", vec![100i32, 200, 300, 400].into_array()),
("value", vec!["p", "q", "r", "s"].into_array()),
("id", vec![100i32, 200, 300, 400].into_arrow_array()),
("value", vec!["p", "q", "r", "s"].into_arrow_array()),
])?;

let storage = Arc::new(InMemory::new());
Expand Down Expand Up @@ -1418,12 +1418,12 @@ async fn test_file_path_metadata_column() -> Result<(), Box<dyn std::error::Erro

// Set up an in-memory table with multiple data files
let batch1 = generate_batch(vec![
("id", vec![1i32, 2, 3].into_array()),
("value", vec!["a", "b", "c"].into_array()),
("id", vec![1i32, 2, 3].into_arrow_array()),
("value", vec!["a", "b", "c"].into_arrow_array()),
])?;
let batch2 = generate_batch(vec![
("id", vec![10i32, 20].into_array()),
("value", vec!["x", "y"].into_array()),
("id", vec![10i32, 20].into_arrow_array()),
("value", vec!["x", "y"].into_arrow_array()),
])?;

let storage = Arc::new(InMemory::new());
Expand Down
14 changes: 7 additions & 7 deletions test-utils/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,29 +138,29 @@ pub fn record_batch_to_bytes_with_props(

/// Anything that implements `IntoArray` can turn itself into a reference to an arrow array
pub trait IntoArray {
fn into_array(self) -> ArrayRef;
fn into_arrow_array(self) -> ArrayRef;
}

impl IntoArray for Vec<i32> {
fn into_array(self) -> ArrayRef {
fn into_arrow_array(self) -> ArrayRef {
Arc::new(Int32Array::from(self))
}
}

impl IntoArray for Vec<i64> {
fn into_array(self) -> ArrayRef {
fn into_arrow_array(self) -> ArrayRef {
Arc::new(Int64Array::from(self))
}
}

impl IntoArray for Vec<bool> {
fn into_array(self) -> ArrayRef {
fn into_arrow_array(self) -> ArrayRef {
Arc::new(BooleanArray::from(self))
}
}

impl IntoArray for Vec<&'static str> {
fn into_array(self) -> ArrayRef {
fn into_arrow_array(self) -> ArrayRef {
Arc::new(StringArray::from(self))
}
}
Expand All @@ -179,8 +179,8 @@ where
/// respectively
pub fn generate_simple_batch() -> Result<RecordBatch, ArrowError> {
generate_batch(vec![
("id", vec![1, 2, 3].into_array()),
("val", vec!["a", "b", "c"].into_array()),
("id", vec![1, 2, 3].into_arrow_array()),
("val", vec!["a", "b", "c"].into_arrow_array()),
])
}

Expand Down
Loading