diff --git a/crates/integrations/datafusion/src/relation_planner.rs b/crates/integrations/datafusion/src/relation_planner.rs index 60978e3d..d9cf0492 100644 --- a/crates/integrations/datafusion/src/relation_planner.rs +++ b/crates/integrations/datafusion/src/relation_planner.rs @@ -33,7 +33,7 @@ use datafusion::sql::sqlparser::ast::{self, TableFactor, TableVersion}; use paimon::spec::{SCAN_TIMESTAMP_MILLIS_OPTION, SCAN_VERSION_OPTION}; use crate::catalog::ReadOnlyTableProvider; -use crate::table::PaimonTableProvider; +use crate::table::{ObjectTableProvider, PaimonTableProvider}; /// A [`RelationPlanner`] that intercepts `VERSION AS OF` and `TIMESTAMP AS OF` /// clauses on Paimon tables and resolves them to time travel options. @@ -93,17 +93,36 @@ impl RelationPlanner for PaimonRelationPlanner { )); } - let extra_options = match version { - TableVersion::VersionAsOf(expr) => resolve_version_as_of(expr)?, - TableVersion::TimestampAsOf(expr) => resolve_timestamp_as_of(expr)?, - _ => return Ok(RelationPlanning::Original(Box::new(relation))), - }; - - // Check if this is a Paimon table. let Some(paimon_provider) = provider.downcast_ref::() else { + // Ours, but with no snapshot to rewrite onto: refuse, because + // handing them on reaches the default planner, which drops the + // clause. Any other provider is another engine's, clause and all. + if provider.is::() + || crate::system_tables::is_system_table_provider(provider.as_ref()) + { + return Err(plan_datafusion_err!( + "time travel is not supported for '{table_ref}'" + )); + } return Ok(RelationPlanning::Original(Box::new(relation))); }; + let extra_options = match version { + TableVersion::VersionAsOf(expr) => resolve_version_as_of(expr)?, + // Same timestamp expression as `TIMESTAMP AS OF`. + TableVersion::TimestampAsOf(expr) | TableVersion::ForSystemTimeAsOf(expr) => { + resolve_timestamp_as_of(expr)? + } + // `AT(...)` and `CHANGES(...)` name neither a snapshot nor a + // timestamp; dropping the clause would answer with current rows. + _ => { + return Err(plan_datafusion_err!( + "this time-travel syntax is not supported for Paimon tables; \ + use VERSION AS OF, TIMESTAMP AS OF or FOR SYSTEM_TIME AS OF" + )) + } + }; + // Resolving time travel may switch the table to the snapshot's schema, // which requires async IO; this planner hook is synchronous, so bridge // through the shared runtime like other sync DataFusion callbacks. diff --git a/crates/integrations/datafusion/src/system_tables/branches.rs b/crates/integrations/datafusion/src/system_tables/branches.rs index 01d942c5..c925def7 100644 --- a/crates/integrations/datafusion/src/system_tables/branches.rs +++ b/crates/integrations/datafusion/src/system_tables/branches.rs @@ -53,7 +53,7 @@ fn branches_schema() -> SchemaRef { } #[derive(Debug)] -struct BranchesTable { +pub(super) struct BranchesTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/consumers.rs b/crates/integrations/datafusion/src/system_tables/consumers.rs index b8a33d4d..40c922ce 100644 --- a/crates/integrations/datafusion/src/system_tables/consumers.rs +++ b/crates/integrations/datafusion/src/system_tables/consumers.rs @@ -51,7 +51,7 @@ fn consumers_schema() -> SchemaRef { } #[derive(Debug)] -struct ConsumersTable { +pub(super) struct ConsumersTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/files.rs b/crates/integrations/datafusion/src/system_tables/files.rs index 07264a3d..e9749007 100644 --- a/crates/integrations/datafusion/src/system_tables/files.rs +++ b/crates/integrations/datafusion/src/system_tables/files.rs @@ -84,7 +84,7 @@ fn files_schema() -> SchemaRef { } #[derive(Debug)] -struct FilesTable { +pub(super) struct FilesTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/manifests.rs b/crates/integrations/datafusion/src/system_tables/manifests.rs index 6f8748f7..9380b316 100644 --- a/crates/integrations/datafusion/src/system_tables/manifests.rs +++ b/crates/integrations/datafusion/src/system_tables/manifests.rs @@ -61,7 +61,7 @@ fn manifests_schema() -> SchemaRef { } #[derive(Debug)] -struct ManifestsTable { +pub(super) struct ManifestsTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/mod.rs b/crates/integrations/datafusion/src/system_tables/mod.rs index 3b3a0b09..c81b8403 100644 --- a/crates/integrations/datafusion/src/system_tables/mod.rs +++ b/crates/integrations/datafusion/src/system_tables/mod.rs @@ -97,6 +97,26 @@ pub(crate) fn is_registered(name: &str) -> bool { } /// Wraps an already-loaded base table as the system table `name`. +/// Does `provider` serve one of the [`TABLES`] above? +/// +/// By type, not by the `base$name` spelling — another engine may name a table +/// with a `$`. Keep in step with `TABLES`: one missing here goes back to +/// having its time-travel clause silently dropped. +pub(crate) fn is_system_table_provider(provider: &dyn TableProvider) -> bool { + provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() + || provider.is::() +} + fn wrap_to_system_table(name: &str, base_table: Table) -> Option>> { TABLES .iter() diff --git a/crates/integrations/datafusion/src/system_tables/options.rs b/crates/integrations/datafusion/src/system_tables/options.rs index b611ba37..04d85f87 100644 --- a/crates/integrations/datafusion/src/system_tables/options.rs +++ b/crates/integrations/datafusion/src/system_tables/options.rs @@ -47,7 +47,7 @@ fn options_schema() -> SchemaRef { } #[derive(Debug)] -struct OptionsTable { +pub(super) struct OptionsTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/partitions.rs b/crates/integrations/datafusion/src/system_tables/partitions.rs index a0b86c2b..749bb282 100644 --- a/crates/integrations/datafusion/src/system_tables/partitions.rs +++ b/crates/integrations/datafusion/src/system_tables/partitions.rs @@ -88,7 +88,7 @@ fn partitions_schema() -> SchemaRef { .clone() } -struct PartitionsTable { +pub(super) struct PartitionsTable { catalog: Arc, identifier: Identifier, table: Table, diff --git a/crates/integrations/datafusion/src/system_tables/physical_files_size.rs b/crates/integrations/datafusion/src/system_tables/physical_files_size.rs index a3e242dd..01ef43c1 100644 --- a/crates/integrations/datafusion/src/system_tables/physical_files_size.rs +++ b/crates/integrations/datafusion/src/system_tables/physical_files_size.rs @@ -54,7 +54,7 @@ fn output_schema() -> SchemaRef { } #[derive(Debug)] -struct PhysicalFilesSizeTable { +pub(super) struct PhysicalFilesSizeTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs b/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs index c8f2b40a..568663ca 100644 --- a/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs +++ b/crates/integrations/datafusion/src/system_tables/referenced_files_size.rs @@ -55,7 +55,7 @@ fn output_schema() -> SchemaRef { } #[derive(Debug)] -struct ReferencedFilesSizeTable { +pub(super) struct ReferencedFilesSizeTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/schemas.rs b/crates/integrations/datafusion/src/system_tables/schemas.rs index 800188d5..7575b3b0 100644 --- a/crates/integrations/datafusion/src/system_tables/schemas.rs +++ b/crates/integrations/datafusion/src/system_tables/schemas.rs @@ -59,7 +59,7 @@ fn schemas_schema() -> SchemaRef { } #[derive(Debug)] -struct SchemasTable { +pub(super) struct SchemasTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/snapshots.rs b/crates/integrations/datafusion/src/system_tables/snapshots.rs index da3428f5..040c51c3 100644 --- a/crates/integrations/datafusion/src/system_tables/snapshots.rs +++ b/crates/integrations/datafusion/src/system_tables/snapshots.rs @@ -65,7 +65,7 @@ fn snapshots_schema() -> SchemaRef { } #[derive(Debug)] -struct SnapshotsTable { +pub(super) struct SnapshotsTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/table_indexes.rs b/crates/integrations/datafusion/src/system_tables/table_indexes.rs index 2533b828..cbd1c2c1 100644 --- a/crates/integrations/datafusion/src/system_tables/table_indexes.rs +++ b/crates/integrations/datafusion/src/system_tables/table_indexes.rs @@ -83,7 +83,7 @@ fn dv_meta_fields() -> Fields { } #[derive(Debug)] -struct TableIndexesTable { +pub(super) struct TableIndexesTable { table: Table, } diff --git a/crates/integrations/datafusion/src/system_tables/tags.rs b/crates/integrations/datafusion/src/system_tables/tags.rs index 09af4231..433d59f1 100644 --- a/crates/integrations/datafusion/src/system_tables/tags.rs +++ b/crates/integrations/datafusion/src/system_tables/tags.rs @@ -62,7 +62,7 @@ fn tags_schema() -> SchemaRef { } #[derive(Debug)] -struct TagsTable { +pub(super) struct TagsTable { table: Table, } diff --git a/crates/integrations/datafusion/src/table_loader.rs b/crates/integrations/datafusion/src/table_loader.rs index c72d0385..664a58c0 100644 --- a/crates/integrations/datafusion/src/table_loader.rs +++ b/crates/integrations/datafusion/src/table_loader.rs @@ -66,10 +66,33 @@ pub(crate) async fn load_table_for_read( identifier.database().to_string(), parsed.table().to_string(), ); - let mut table = catalog - .get_table(&base_identifier) - .await - .map_err(to_datafusion_error)?; + let mut table = match catalog.load_table(&base_identifier).await { + Ok(paimon::catalog::LoadedTable::Paimon(table)) => *table, + Ok(paimon::catalog::LoadedTable::Object(_)) => { + // The search UDTFs land here too, so the message names the table + // rather than what the caller asked of it. + return Err(DataFusionError::Plan(format!( + "table '{}' is declared 'object-table' and cannot be read as a Paimon table", + base_identifier.full_name() + ))); + } + Ok(paimon::catalog::LoadedTable::External(external)) => { + return Err(DataFusionError::Plan(format!( + "table '{}' is declared '{}' and cannot be read as a Paimon table", + base_identifier.full_name(), + external.declared() + ))); + } + // `LoadedTable` is non_exhaustive: a variant added upstream is not a + // Paimon table until this path says how to read one. + Ok(_) => { + return Err(DataFusionError::Plan(format!( + "table '{}' cannot be read as a Paimon table", + base_identifier.full_name() + ))); + } + Err(err) => return Err(to_datafusion_error(err)), + }; let system_table = parsed.system_table().map(str::to_string); if let Some(branch) = parsed.branch() { let is_branches_table = system_table diff --git a/crates/integrations/datafusion/tests/table_type_routing.rs b/crates/integrations/datafusion/tests/table_type_routing.rs index 5620d6bc..3e979637 100644 --- a/crates/integrations/datafusion/tests/table_type_routing.rs +++ b/crates/integrations/datafusion/tests/table_type_routing.rs @@ -647,6 +647,113 @@ async fn object_and_lance_tables_route_to_engines() { assert_eq!(column_i32(&batches), vec![1, 3]); } +#[tokio::test] +async fn an_unsupported_time_travel_clause_on_a_paimon_table_is_rejected() { + use datafusion::prelude::SessionContext; + use paimon_datafusion::PaimonCatalogProvider; + + let paimon_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", paimon_dir.path().display()); + let mut options = Options::new(); + options.set(CatalogOptions::WAREHOUSE, warehouse); + let fs_catalog = Arc::new(FileSystemCatalog::new(options).unwrap()); + fs_catalog + .create_database(DB, false, HashMap::new()) + .await + .unwrap(); + let plain = PaimonSchema::builder() + .column( + "id", + paimon::spec::DataType::Int(paimon::spec::IntType::new()), + ) + .build() + .unwrap(); + fs_catalog + .create_table(&Identifier::new(DB, "pt"), plain, false) + .await + .unwrap(); + + let ctx = SessionContext::new(); + ctx.register_catalog( + CATALOG, + Arc::new(PaimonCatalogProvider::new( + Some(CATALOG.to_string()), + fs_catalog, + Default::default(), + Default::default(), + None, + )), + ); + paimon_datafusion::register_catalog_table_engine( + &ctx, + CATALOG, + TableType::IcebergTable, + Arc::new(FakeEngineResolver), + ) + .unwrap(); + for dialect in ["databricks", "mssql", "bigquery", "snowflake"] { + ctx.sql(&format!("SET datafusion.sql_parser.dialect = '{dialect}'")) + .await + .unwrap() + .collect() + .await + .unwrap(); + + let run = |sql: String| { + let ctx = &ctx; + async move { + match ctx.sql(&sql).await { + Err(err) => Err(err.to_string()), + Ok(df) => df.collect().await.map(|_| ()).map_err(|e| e.to_string()), + } + } + }; + let ts = "2020-01-01 00:00:00"; + let system_time = run(format!( + "SELECT * FROM {CATALOG}.{DB}.pt FOR SYSTEM_TIME AS OF '{ts}'" + )) + .await; + let timestamp = run(format!( + "SELECT * FROM {CATALOG}.{DB}.pt TIMESTAMP AS OF '{ts}'" + )) + .await; + assert_eq!( + system_time, timestamp, + "[{dialect}] the two spellings diverged" + ); + + let Err(msg) = run(format!("SELECT * FROM {CATALOG}.{DB}.pt AT('{ts}')")).await else { + panic!("[{dialect}] a historical clause must not answer with current rows"); + }; + assert!( + msg.contains("this time-travel syntax is not supported"), + "[{dialect}] {msg}" + ); + + let err = ctx + .sql(&format!("SELECT * FROM {CATALOG}.{DB}.pt VERSION AS OF 1")) + .await + .unwrap() + .collect() + .await + .err() + .unwrap_or_else(|| panic!("[{dialect}] snapshot 1 does not exist in this fixture")); + assert!(err.to_string().contains("Snapshot 1"), "[{dialect}] {err}"); + + let Err(msg) = run(format!( + "SELECT * FROM {CATALOG}.{DB}.\"pt$schemas\" VERSION AS OF 1" + )) + .await + else { + panic!("[{dialect}] a system table answered a historical clause with current rows"); + }; + assert!( + msg.contains("time travel is not supported for"), + "[{dialect}] {msg}" + ); + } +} + #[tokio::test] async fn time_travel_on_routed_tables_is_rejected() { let env = setup().await; @@ -1252,3 +1359,84 @@ async fn a_rejected_external_table_does_not_pollute_the_blob_registry() { "a rejected table must leave no registration behind" ); } + +#[derive(Debug)] +struct ForeignRelationPlanner; + +impl datafusion::logical_expr::planner::RelationPlanner for ForeignRelationPlanner { + fn plan_relation( + &self, + relation: datafusion::sql::sqlparser::ast::TableFactor, + _context: &mut dyn datafusion::logical_expr::planner::RelationPlannerContext, + ) -> DFResult { + use datafusion::sql::sqlparser::ast::TableFactor; + if matches!( + relation, + TableFactor::Table { + version: Some(_), + .. + } + ) { + return Err(DataFusionError::Plan("foreign planner reached".into())); + } + Ok(datafusion::logical_expr::planner::RelationPlanning::Original(Box::new(relation))) + } +} + +#[tokio::test] +async fn a_foreign_provider_keeps_its_own_version_clause() { + use datafusion::prelude::SessionContext; + + let ctx = SessionContext::new(); + for name in ["foreign_table", "foreign$history", "foreign$schemas"] { + let schema = Arc::new(ArrowSchema::new(vec![Field::new( + "id", + DataType::Int32, + false, + )])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vec![1, 2])) as _], + ) + .unwrap(); + ctx.register_table( + datafusion::common::TableReference::bare(name), + Arc::new(MemTable::try_new(schema, vec![vec![batch]]).unwrap()), + ) + .unwrap(); + } + ctx.register_relation_planner(Arc::new(ForeignRelationPlanner)) + .unwrap(); + ctx.register_relation_planner(Arc::new(paimon_datafusion::PaimonRelationPlanner::new())) + .unwrap(); + ctx.sql("SET datafusion.sql_parser.dialect = 'databricks'") + .await + .unwrap() + .collect() + .await + .unwrap(); + + for table in [ + "foreign_table", + "\"foreign$history\"", + "\"foreign$schemas\"", + ] { + for clause in [ + "FOR SYSTEM_TIME AS OF CURRENT_TIMESTAMP() - INTERVAL '1' DAY", + "TIMESTAMP AS OF '2020-01-01T00:00:00Z'", + "VERSION AS OF 1", + ] { + let err = ctx + .sql(&format!("SELECT * FROM {table} {clause}")) + .await + .err() + .unwrap_or_else(|| { + panic!("[{table} {clause}] expected the foreign planner to claim it") + }); + assert!( + err.to_string().contains("foreign planner reached"), + "[{table} {clause}] Paimon took a clause that is not its own: {err}" + ); + } + } +}