Skip to content

Commit 4c1f548

Browse files
committed
fix(datafusion): reject non-append insert operations
1 parent 93879b2 commit 4c1f548

1 file changed

Lines changed: 26 additions & 44 deletions

File tree

  • crates/integrations/datafusion/src/table

crates/integrations/datafusion/src/table/mod.rs

Lines changed: 26 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -159,11 +159,11 @@ impl TableProvider for IcebergTableProvider {
159159
&self,
160160
state: &dyn Session,
161161
input: Arc<dyn ExecutionPlan>,
162-
insert_op: InsertOp,
162+
_insert_op: InsertOp,
163163
) -> DFResult<Arc<dyn ExecutionPlan>> {
164-
if insert_op != InsertOp::Append {
164+
if _insert_op != InsertOp::Append {
165165
return Err(DataFusionError::NotImplemented(format!(
166-
"IcebergTableProvider supports only append inserts, got {insert_op}"
166+
"IcebergTableProvider supports only append inserts, got {_insert_op}"
167167
)));
168168
}
169169

@@ -717,57 +717,39 @@ mod tests {
717717
}
718718

719719
#[tokio::test]
720-
async fn test_catalog_backed_provider_rejects_overwrite_insert() {
720+
async fn test_catalog_backed_provider_rejects_non_append_op() {
721721
use datafusion::physical_plan::empty::EmptyExec;
722722

723723
let (catalog, namespace, table_name, _temp_dir) = get_test_catalog_and_table().await;
724724
let provider = IcebergTableProvider::try_new(catalog, namespace, table_name)
725725
.await
726726
.unwrap();
727727
let ctx = SessionContext::new();
728-
let input = Arc::new(EmptyExec::new(provider.schema())) as Arc<dyn ExecutionPlan>;
729-
730-
let error = provider
731-
.insert_into(&ctx.state(), input, InsertOp::Overwrite)
732-
.await
733-
.expect_err("overwrite inserts should be rejected");
734728

735-
assert!(
736-
matches!(
737-
error,
738-
DataFusionError::NotImplemented(ref message)
739-
if message
740-
== "IcebergTableProvider supports only append inserts, got Insert Overwrite"
729+
for (insert_op, expected_message) in [
730+
(
731+
InsertOp::Overwrite,
732+
"IcebergTableProvider supports only append inserts, got Insert Overwrite",
741733
),
742-
"unexpected error: {error}"
743-
);
744-
}
745-
746-
#[tokio::test]
747-
async fn test_catalog_backed_provider_rejects_replace_insert() {
748-
use datafusion::physical_plan::empty::EmptyExec;
749-
750-
let (catalog, namespace, table_name, _temp_dir) = get_test_catalog_and_table().await;
751-
let provider = IcebergTableProvider::try_new(catalog, namespace, table_name)
752-
.await
753-
.unwrap();
754-
let ctx = SessionContext::new();
755-
let input = Arc::new(EmptyExec::new(provider.schema())) as Arc<dyn ExecutionPlan>;
756-
757-
let error = provider
758-
.insert_into(&ctx.state(), input, InsertOp::Replace)
759-
.await
760-
.expect_err("replace inserts should be rejected");
761-
762-
assert!(
763-
matches!(
764-
error,
765-
DataFusionError::NotImplemented(ref message)
766-
if message
767-
== "IcebergTableProvider supports only append inserts, got Replace Into"
734+
(
735+
InsertOp::Replace,
736+
"IcebergTableProvider supports only append inserts, got Replace Into",
768737
),
769-
"unexpected error: {error}"
770-
);
738+
] {
739+
let input = Arc::new(EmptyExec::new(provider.schema())) as Arc<dyn ExecutionPlan>;
740+
let error = provider
741+
.insert_into(&ctx.state(), input, insert_op)
742+
.await
743+
.expect_err("non-append inserts should be rejected");
744+
745+
assert!(
746+
matches!(
747+
error,
748+
DataFusionError::NotImplemented(ref message) if message == expected_message
749+
),
750+
"unexpected error: {error}"
751+
);
752+
}
771753
}
772754

773755
#[tokio::test]

0 commit comments

Comments
 (0)