NoahKusaba commented on code in PR #20:
URL: https://github.com/apache/datafusion-iceberg/pull/20#discussion_r4135752460


##########
crates/datafusion/tests/integration_datafusion_test.rs:
##########
@@ -977,3 +987,150 @@ async fn test_insert_into_partitioned() -> Result<(), 
Box<dyn Error>> {
 
     Ok(())
 }
+
+/// Executes the single partition of `plan` and renders its rows as a table.
+async fn run(
+    plan: &dyn ExecutionPlan,
+    ctx: &SessionContext,
+) -> Result<String, Box<dyn Error>> {
+    let stream = plan.execute(0, ctx.task_ctx())?;
+    let batches = datafusion::physical_plan::common::collect(stream).await?;
+    Ok(pretty_format_batches(&batches)?.to_string())
+}
+
+/// Returns the first node of type `T` in `plan`, depth first.
+fn find_node<T: ExecutionPlan + 'static>(plan: &Arc<dyn ExecutionPlan>) -> 
Option<&T> {
+    plan.downcast_ref::<T>()
+        .or_else(|| plan.children().into_iter().find_map(find_node::<T>))
+}
+
+/// The plan nodes and providers can be named and inspected from outside this
+/// crate, and rebuilt from their parts, as a codec that serializes them does.
+#[tokio::test]
+async fn test_plan_nodes_are_inspectable() -> Result<(), Box<dyn Error>> {
+    let iceberg_catalog = get_iceberg_catalog().await;
+    let namespace = NamespaceIdent::new("test_plan_nodes".to_string());
+    set_test_namespace(&iceberg_catalog, &namespace).await?;
+    let creation = get_table_creation(temp_path(), "my_table", None)?;
+    iceberg_catalog.create_table(&namespace, creation).await?;
+    let ident = TableIdent::new(namespace.clone(), "my_table".to_string());
+    let client: Arc<dyn Catalog> = Arc::new(iceberg_catalog);
+
+    let ctx = SessionContext::new();
+    let catalog = IcebergCatalogProvider::try_new(client.clone()).await?;
+    ctx.register_catalog("catalog", Arc::new(catalog));
+    let provider = ctx
+        .table_provider("catalog.test_plan_nodes.my_table")
+        .await?;
+    let provider = provider
+        .downcast_ref::<IcebergTableProvider>()
+        .expect("a catalog-backed provider");
+    assert_eq!(provider.table_ident(), &ident);
+    assert!(Arc::ptr_eq(provider.catalog(), &client));
+    let rebuilt = IcebergTableProvider::try_new(
+        provider.catalog().clone(),
+        provider.table_ident().namespace().clone(),
+        provider.table_ident().name(),
+    )
+    .await?;
+    assert_eq!(rebuilt.table_ident(), &ident);
+    assert_eq!(rebuilt.schema(), provider.schema());
+
+    // Write path: a commit above a write, both holding the table, and the
+    // commit going through the provider's catalog. The plan that runs is
+    // rebuilt from their accessors and children alone. The optimizer drops
+    // the coalesce above a single-partition write, so the rebuilt commit
+    // always gets one, as a codec would.
+    let insert = ctx
+        .sql("INSERT INTO catalog.test_plan_nodes.my_table VALUES (1, 'alan'), 
(2, 'turing')")
+        .await?
+        .create_physical_plan()
+        .await?;
+    let commit = insert
+        .downcast_ref::<IcebergCommitExec>()
+        .expect("the insert plan is rooted at a commit");
+    assert_eq!(commit.table().identifier(), &ident);
+    assert!(Arc::ptr_eq(commit.catalog(), &client));
+    let write = find_node::<IcebergWriteExec>(&insert).expect("a write below 
the commit");
+    assert_eq!(write.table().identifier(), &ident);
+    let rebuilt_write: Arc<dyn ExecutionPlan> = Arc::new(IcebergWriteExec::new(
+        write.table().clone(),
+        write.children()[0].clone(),
+    ));
+    let rebuilt_commit = IcebergCommitExec::new(
+        commit.table().clone(),
+        commit.catalog().clone(),
+        Arc::new(CoalescePartitionsExec::new(rebuilt_write)),
+    );
+    let inserted = run(&rebuilt_commit, &ctx).await?;
+    assert!(inserted.contains("| 2     |"), "{inserted}");
+
+    // Read path: a scan pinned to a snapshot, rebuilt from its accessors,
+    // returns the same rows.
+    let table = client.load_table(&ident).await?;
+    let snapshot_id = table.metadata().current_snapshot_id().unwrap();
+    let pinned =
+        IcebergStaticTableProvider::try_new_from_table_snapshot(table, 
snapshot_id)
+            .await?;
+    assert_eq!(pinned.snapshot_id(), Some(snapshot_id));
+    ctx.register_table("pinned", Arc::new(pinned.clone()))?;
+    let plan = ctx
+        .sql("SELECT foo2 FROM pinned WHERE foo1 = 1")
+        .await?
+        .create_physical_plan()
+        .await?;
+    let scan = find_node::<IcebergTableScan>(&plan).expect("a scan");
+    assert!(scan.predicates().is_some(), "the filter is pushed down");
+    let rebuilt = IcebergTableScan::new_with_predicate(
+        scan.table().clone(),
+        scan.snapshot_id(),
+        scan.schema(),
+        scan.projection().map(<[String]>::to_vec),
+        scan.predicates().cloned(),
+        scan.limit(),
+    )?;
+    assert_eq!(rebuilt.schema(), scan.schema());
+    let expected = run(scan, &ctx).await?;
+    assert!(
+        expected.contains("alan") && !expected.contains("turing"),
+        "{expected}"
+    );
+    assert_eq!(run(&rebuilt, &ctx).await?, expected);
+
+    // Without a projection the scan reads every column, and its limit is kept.
+    let plan = pinned.scan(&ctx.state(), None, &[], Some(1)).await?;
+    let scan = plan.downcast_ref::<IcebergTableScan>().expect("a scan");
+    assert_eq!(scan.projection(), None);
+    let rebuilt = IcebergTableScan::new_with_predicate(
+        scan.table().clone(),
+        scan.snapshot_id(),
+        scan.schema(),
+        scan.projection().map(<[String]>::to_vec),
+        scan.predicates().cloned(),
+        scan.limit(),
+    )?;
+    assert_eq!(rebuilt.limit(), Some(1));
+    let expected = run(scan, &ctx).await?;
+    assert_eq!(expected.lines().count(), 5, "one row:\n{expected}");

Review Comment:
   You're absolutely right, I'll fix that in my next commit. 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to