From 4b5a6dba2cc6b30ad2376025b5b304a97650bb14 Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Sun, 20 Sep 2026 20:33:59 +0700 Subject: [PATCH 1/4] feat: publish the native Iceberg storage scheme list over JNI The native storage factory now exposes the schemes it can open per access mode, and a JNI entry returns them, so the JVM read and write gates can load the list instead of mirroring it by hand. Tests keep the list and the factory's match arms in step in both directions. --- .../src/execution/operators/iceberg_common.rs | 54 +++++++++++++++++-- native/core/src/execution/operators/mod.rs | 2 +- native/core/src/lib.rs | 24 +++++++++ 3 files changed, 75 insertions(+), 5 deletions(-) diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 509e3b1f986..3bafca8b31b 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -44,10 +44,10 @@ const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client."]; /// class is configured; `access_mode` is forwarded to the JVM SPI so the read and write paths can /// be granted different (e.g. read-only vs read-write) credentials. /// -/// The JVM planner mirrors the schemes admitted here by hand, so that it can decline a scan -/// cleanly instead of failing at execution. Changing the arms below means updating -/// `CometScanRule.icebergReadableSchemes` (reads) and -/// `CometIcebergNativeWrite.SupportedStorageSchemes` (writes). +/// The JVM planner loads `builtin_storage_schemes` over JNI so it can decline a scan or a write +/// cleanly instead of failing at execution. Opt-in S3-compliant alias schemes come from catalog +/// properties and are added on the JVM side. The tests below keep that list and the match arms +/// here in step. pub(crate) fn storage_factory_for( path: &str, catalog_properties: &HashMap, @@ -96,6 +96,15 @@ pub(crate) fn storage_factory_for( } } +/// The built-in storage schemes `storage_factory_for` admits for `access_mode`, without any +/// opted-in S3-compliant alias. This is the list the JVM read and write gates load over JNI. +pub(crate) fn builtin_storage_schemes(access_mode: AccessMode) -> &'static [&'static str] { + match access_mode { + AccessMode::Read => &["file", "memory", "gs", "oss", "s3", "s3a"], + AccessMode::Write => &["file", "memory", "gs", "s3", "s3a"], + } +} + /// Build a `FileIO` whose storage scheme is inferred from `reference_path` and whose properties /// come from the catalog. The reference path is the metadata location for reads or the data /// location for writes — anything that carries the right URI scheme. `catalog_name` is the @@ -276,6 +285,43 @@ mod tests { ); } + #[test] + fn exposed_scheme_list_matches_storage_factory() { + for mode in [AccessMode::Read, AccessMode::Write] { + for scheme in builtin_storage_schemes(mode) { + let url = format!("{scheme}://bucket/path"); + assert!( + factory_result(&url, mode).is_ok(), + "{scheme} is listed for {mode:?} but the factory rejects it" + ); + } + } + assert!(builtin_storage_schemes(AccessMode::Read).contains(&"oss")); + assert!(!builtin_storage_schemes(AccessMode::Write).contains(&"oss")); + assert!(builtin_storage_schemes(AccessMode::Read).contains(&"memory")); + assert!(builtin_storage_schemes(AccessMode::Write).contains(&"memory")); + } + + #[test] + fn unlisted_schemes_are_rejected_for_both_modes() { + let unlisted = [ + "hdfs", "abfs", "abfss", "wasb", "wasbs", "gcs", "http", "https", "azure", + ]; + for mode in [AccessMode::Read, AccessMode::Write] { + for scheme in unlisted { + let err = factory_result(&format!("{scheme}://bucket/path"), mode).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {scheme} in {mode:?}: {err}" + ); + assert!( + !builtin_storage_schemes(mode).contains(&scheme), + "{scheme} must not be listed for {mode:?}" + ); + } + } + } + #[test] fn scheme_of_extracts_scheme_from_all_uri_forms() { // Host-bearing and hostless/opaque vendor forms must resolve to the same scheme, so an diff --git a/native/core/src/execution/operators/mod.rs b/native/core/src/execution/operators/mod.rs index d09b0b4fb37..37f56fe1c5c 100644 --- a/native/core/src/execution/operators/mod.rs +++ b/native/core/src/execution/operators/mod.rs @@ -34,7 +34,7 @@ mod expand; pub use expand::ExpandExec; mod explode; pub use explode::ExplodeExec; -mod iceberg_common; +pub(crate) mod iceberg_common; mod iceberg_partition_path; mod iceberg_scan; mod iceberg_write; diff --git a/native/core/src/lib.rs b/native/core/src/lib.rs index 0872a478f93..72f3bd7491f 100644 --- a/native/core/src/lib.rs +++ b/native/core/src/lib.rs @@ -52,6 +52,9 @@ pub mod jvm_bridge { use errors::{try_unwrap_or_throw, CometError, CometResult}; +use crate::cloud::s3::credential_bridge::AccessMode; +use crate::execution::operators::iceberg_common::builtin_storage_schemes; + pub mod alloc_accounting; pub mod cloud; pub mod execution; @@ -227,6 +230,27 @@ pub extern "system" fn Java_org_apache_comet_NativeBase_isObjectStoreSchemeSuppo }) } +/// JNI: the comma-joined built-in storage schemes the native Iceberg scan (`for_write` false) or +/// write (`for_write` true) path admits. This is the source of truth the JVM read and write gates +/// load. Opt-in S3-compliant alias schemes are not answered here; the JVM adds them from catalog +/// properties. +#[no_mangle] +pub extern "system" fn Java_org_apache_comet_NativeBase_icebergStorageSchemes( + env: EnvUnowned, + _: JClass, + for_write: jni::sys::jboolean, +) -> jni::sys::jstring { + try_unwrap_or_throw(&env, |env| { + let access_mode = if for_write == jni::sys::JNI_TRUE { + AccessMode::Write + } else { + AccessMode::Read + }; + let joined = builtin_storage_schemes(access_mode).join(","); + Ok(env.new_string(joined)?.into_raw()) + }) +} + // Creates a default log4rs config, which logs to console with log level. fn default_logger_config(log_level: &str) -> CometResult { let console_append = ConsoleAppender::builder() From d0a6bded7a2c60091c037b6b0187425db8987039 Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Sun, 20 Sep 2026 21:04:33 +0700 Subject: [PATCH 2/4] fix: load the Iceberg gate scheme lists from the native storage factory The JVM Iceberg scan and write gates hardcoded the schemes the native storage factory can open and had already drifted from it. Load both lists over JNI through a new IcebergStorageSchemes object, keeping fallback constants only for when the library cannot load, and pin those constants to what native publishes in CometScanSchemeFallbackSuite. A wasb-backed Iceberg table now falls back to Spark with a reason naming the scheme instead of failing natively at execution. --- .../src/execution/operators/iceberg_common.rs | 2 + native/core/src/lib.rs | 7 +- .../java/org/apache/comet/NativeBase.java | 12 +++ .../comet/iceberg/IcebergStorageSchemes.scala | 75 ++++++++++++++++ .../apache/comet/rules/CometScanRule.scala | 20 ++--- .../operator/CometIcebergNativeWrite.scala | 5 +- .../hadoop/fs/FakeWasbSchemeFileSystem.java | 51 +++++++++++ .../rules/CometScanSchemeFallbackSuite.scala | 86 +++++++++++++++++-- 8 files changed, 231 insertions(+), 27 deletions(-) create mode 100644 spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala create mode 100644 spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 3bafca8b31b..7f9cd1aae47 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -98,6 +98,8 @@ pub(crate) fn storage_factory_for( /// The built-in storage schemes `storage_factory_for` admits for `access_mode`, without any /// opted-in S3-compliant alias. This is the list the JVM read and write gates load over JNI. +/// The rejection test below covers a fixed set of unlisted schemes, so a new factory arm must +/// be added to this list as well or the JVM gate keeps declining it. pub(crate) fn builtin_storage_schemes(access_mode: AccessMode) -> &'static [&'static str] { match access_mode { AccessMode::Read => &["file", "memory", "gs", "oss", "s3", "s3a"], diff --git a/native/core/src/lib.rs b/native/core/src/lib.rs index 72f3bd7491f..a549d0d36f5 100644 --- a/native/core/src/lib.rs +++ b/native/core/src/lib.rs @@ -31,6 +31,7 @@ extern crate datafusion_comet_jni_bridge; use jni::{ objects::{JClass, JString}, + sys::{jboolean, jstring, JNI_FALSE}, EnvUnowned, }; use log::info; @@ -238,10 +239,10 @@ pub extern "system" fn Java_org_apache_comet_NativeBase_isObjectStoreSchemeSuppo pub extern "system" fn Java_org_apache_comet_NativeBase_icebergStorageSchemes( env: EnvUnowned, _: JClass, - for_write: jni::sys::jboolean, -) -> jni::sys::jstring { + for_write: jboolean, +) -> jstring { try_unwrap_or_throw(&env, |env| { - let access_mode = if for_write == jni::sys::JNI_TRUE { + let access_mode = if for_write != JNI_FALSE { AccessMode::Write } else { AccessMode::Read diff --git a/spark/src/main/java/org/apache/comet/NativeBase.java b/spark/src/main/java/org/apache/comet/NativeBase.java index 15ddf8bd8b3..c8c2b7e84a8 100644 --- a/spark/src/main/java/org/apache/comet/NativeBase.java +++ b/spark/src/main/java/org/apache/comet/NativeBase.java @@ -331,4 +331,16 @@ public static void releaseNative() throws Throwable { * @return true if object_store can construct both a store and an object key for this URL */ public static native boolean isObjectStoreSchemeSupported(String url); + + /** + * The comma-joined URL schemes the native Iceberg storage factory can open, for reads ({@code + * forWrite} false) or writes ({@code forWrite} true). This is the authoritative list the JVM + * Iceberg scan and write gates load, so the planner never hardcodes (and drifts from) the set the + * native factory actually builds. Opt-in S3-compliant alias schemes are not included; the JVM + * adds them from {@code fs.comet.s3Compliant.schemes}. + * + * @param forWrite true for the write path's schemes, false for the scan path's + * @return the supported schemes joined by commas, e.g. "file,memory,gs,s3,s3a" + */ + public static native String icebergStorageSchemes(boolean forWrite); } diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala new file mode 100644 index 00000000000..e96d7792745 --- /dev/null +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala @@ -0,0 +1,75 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import java.util.Locale + +import scala.util.control.NonFatal + +import org.apache.spark.internal.Logging + +import org.apache.comet.NativeBase + +/** + * The storage schemes the native Iceberg storage factory publishes over JNI, so the JVM scan and + * write gates decline what native cannot open instead of failing at execution. The fallback + * constants are used only when the library cannot be loaded: assuming a scheme is supported would + * recreate the execution-time failure, and without the library nothing runs natively anyway. + */ +private[comet] object IcebergStorageSchemes extends Logging { + + private[comet] val FallbackRead: Set[String] = Set("file", "memory", "s3", "s3a", "gs", "oss") + private[comet] val FallbackWrite: Set[String] = Set("file", "memory", "s3", "s3a", "gs") + + lazy val read: Set[String] = load(forWrite = false, FallbackRead) + lazy val write: Set[String] = load(forWrite = true, FallbackWrite) + + private def load(forWrite: Boolean, fallback: Set[String]): Set[String] = { + val path = if (forWrite) "write" else "read" + val joined = + try NativeBase.icebergStorageSchemes(forWrite) + catch { + case e: UnsatisfiedLinkError => + logWarning( + s"Comet native library is not loaded; using the fallback Iceberg $path scheme list " + + s"${fallback.toSeq.sorted.mkString(", ")}: ${e.getMessage}") + return fallback + case NonFatal(e) => + logWarning( + s"Failed to load the Iceberg $path scheme list from the Comet native library; using " + + s"the fallback list ${fallback.toSeq.sorted.mkString(", ")}", + e) + return fallback + } + val schemes = Option(joined).toSeq + .flatMap(_.split(",")) + .map(_.trim.toLowerCase(Locale.ROOT)) + .filter(_.nonEmpty) + .toSet + if (schemes.isEmpty) { + logWarning( + s"Comet native library published an empty Iceberg $path scheme list; using the fallback " + + s"list ${fallback.toSeq.sorted.mkString(", ")}") + fallback + } else { + schemes + } + } +} diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index 5f37761a051..95ec5390caa 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -46,7 +46,7 @@ import org.apache.spark.sql.types._ import org.apache.comet.{CometConf, DataTypeSupport, NativeBase} import org.apache.comet.CometConf._ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plus, withFallbackReason, withFallbackReasons} -import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} +import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection, IcebergStorageSchemes} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, readFieldId} import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} @@ -1192,20 +1192,12 @@ object CometScanRule extends Logging { org.apache.spark.sql.catalyst.trees.TreeNodeTag[Unit]("comet.skipCometScan") /** - * Schemes Comet's native Iceberg scan can actually open, mirroring the match arms in - * `native/core/src/execution/operators/iceberg_common.rs::storage_factory_for`. Deliberately - * NOT delegated to `isNativelyReadableScheme`: object_store recognizes schemes (http/https, - * azure, memory) that iceberg-rust's OpenDAL storage factory cannot build, and admitting them - * here turns a clean JVM fallback into a native runtime "Unsupported storage scheme" error. Add - * here what you add to `storage_factory_for` (currently Aliyun `oss` and GCS `gs`). - * S3-compliant aliases like `blob` are opt-in via `fs.comet.s3Compliant.schemes` (see - * `isIcebergReadableScheme`), not hardcoded, since the native planner opens them via S3. The - * write path keeps its own list (`CometIcebergNativeWrite.SupportedStorageSchemes`), which - * differs deliberately: it excludes `oss` (fails closed, see `storage_factory_for`) and - * includes `memory`. + * Schemes Comet's native Iceberg scan can open, loaded from the native storage factory over JNI + * so this gate cannot drift from `storage_factory_for`. Opt-in aliases from + * `fs.comet.s3Compliant.schemes` are additive (see `isIcebergReadableScheme`); the write path + * loads its own set (`CometIcebergNativeWrite.SupportedStorageSchemes`). */ - private val icebergReadableSchemes: Set[String] = - Set("file", "s3", "s3a", "gs", "oss") + private val icebergReadableSchemes: Set[String] = IcebergStorageSchemes.read /** * "Supported schemes: ..." suffix shared by the Iceberg scheme-fallback messages. Lists the diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala index 94f6b139273..1209f5af640 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala @@ -29,7 +29,7 @@ import org.apache.spark.sql.comet.{CometIcebergWriteExec, CometNativeExec, Icebe import org.apache.comet.{CometConf, ConfigEntry} import org.apache.comet.CometSparkSessionExtensions.withFallbackReason -import org.apache.comet.iceberg.IcebergReflection +import org.apache.comet.iceberg.{IcebergReflection, IcebergStorageSchemes} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.serde.{CometOperatorSerde, Compatible, OperatorOuterClass, SupportLevel, Unsupported} import org.apache.comet.serde.OperatorOuterClass.Operator @@ -85,8 +85,7 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { // `oss.*` catalog properties to it and no functional test covers the path, so an OSS write // could silently drop endpoint/credential configuration. Fail closed until it is covered. // `gs` is additionally gated on the resolved FileIO (`requireGcsFileIOForGcsDataLocation`). - private val SupportedStorageSchemes: Set[String] = - Set("file", "memory", "s3", "s3a", "gs") + private val SupportedStorageSchemes: Set[String] = IcebergStorageSchemes.write private val MinUnsupportedFormatVersion = 3 private val ParquetWritePropertyPrefix = "write.parquet." private val ParquetMrPropertyPrefix = "parquet." diff --git a/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java new file mode 100644 index 00000000000..02830e5408b --- /dev/null +++ b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java @@ -0,0 +1,51 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.hadoop.fs; + +import java.net.URI; + +import org.apache.hadoop.fs.RawLocalFileSystem; + +/** + * A local-disk-backed FileSystem that reports the {@code wasb} scheme, so a test can create and + * read an Iceberg table under a {@code wasb://} warehouse without Azure. Used to assert that {@code + * CometScanRule} declines a native Iceberg scan whose storage scheme object_store recognizes but + * the native Iceberg storage factory cannot build, instead of failing at execution. + */ +public class FakeWasbSchemeFileSystem extends RawLocalFileSystem { + + public static final String PREFIX = "wasb://fake-container"; + + public FakeWasbSchemeFileSystem() { + // Avoid `URI scheme is not "file"` error on + // RawLocalFileSystem$DeprecatedRawLocalFileStatus.getOwner + RawLocalFileSystem.useStatIfAvailable(); + } + + @Override + public String getScheme() { + return "wasb"; + } + + @Override + public URI getUri() { + return URI.create(PREFIX); + } +} diff --git a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala index 5a14f9f592e..9de3c129184 100644 --- a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala +++ b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala @@ -27,11 +27,12 @@ import java.util.UUID import org.apache.commons.io.FileUtils import org.apache.spark.SparkConf import org.apache.spark.sql.{CometTestBase, SaveMode} -import org.apache.spark.sql.comet.CometScanExec +import org.apache.spark.sql.comet.{CometIcebergNativeScanExec, CometScanExec} import org.apache.spark.sql.execution.{FileSourceScanExec, SparkPlan} -import org.apache.comet.CometConf -import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem} +import org.apache.comet.{CometConf, CometIcebergTestBase, NativeBase} +import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem, FakeWasbSchemeFileSystem} +import org.apache.comet.iceberg.IcebergStorageSchemes /** * Comet's native readers go through object_store, which understands a fixed set of URL schemes. A @@ -43,7 +44,7 @@ import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem} * native probe is cached per probe URL so an authorityless URL can't poison the authority-bearing * form of the same scheme. */ -class CometScanSchemeFallbackSuite extends CometTestBase { +class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBase { private var fakeRootDir: File = _ @@ -54,6 +55,9 @@ class CometScanSchemeFallbackSuite extends CometTestBase { // cluster. `hdfs` is natively readable by default, so this scan must be CLAIMED, not declined. conf.set("spark.hadoop.fs.hdfs.impl", "org.apache.comet.hadoop.fs.FakeHdfsSchemeFileSystem") conf.set("spark.hadoop.fs.defaultFS", FakeHDFSFileSystem.PREFIX) + // Back `wasb` with a local FS so an Iceberg table can live under a `wasb://` warehouse. The + // native Iceberg storage factory has no arm for it, so that scan must be DECLINED, not claimed. + conf.set("spark.hadoop.fs.wasb.impl", classOf[FakeWasbSchemeFileSystem].getName) // Intentionally NOT setting CometConf.COMET_LIBHDFS_SCHEMES -- `fake` is not natively readable, // and `hdfs` must still be claimed by default (mirrors the native `is_hdfs_scheme` default). conf @@ -120,12 +124,46 @@ class CometScanSchemeFallbackSuite extends CometTestBase { "gs://bucket/... must be admitted even after an authorityless gs:// URI was probed first") } + test("iceberg gate: native scheme lists match the JVM fallback constants") { + // The gates load their scheme sets from the native storage factory over JNI. The JVM + // fallback constants only stand in when the library cannot be loaded, so this pins them to + // what native publishes: extending `builtin_storage_schemes` without updating the fallback + // fails here rather than drifting silently. + assume(NativeBase.isLoaded, "Comet native library not loaded") + assert( + IcebergStorageSchemes.read == IcebergStorageSchemes.FallbackRead, + s"native read schemes ${IcebergStorageSchemes.read.toSeq.sorted} differ from the JVM " + + s"fallback ${IcebergStorageSchemes.FallbackRead.toSeq.sorted}") + assert( + IcebergStorageSchemes.write == IcebergStorageSchemes.FallbackWrite, + s"native write schemes ${IcebergStorageSchemes.write.toSeq.sorted} differ from the JVM " + + s"fallback ${IcebergStorageSchemes.FallbackWrite.toSeq.sorted}") + Seq("memory", "oss").foreach { scheme => + assert( + IcebergStorageSchemes.read.contains(scheme), + s"$scheme must be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") + } + assert( + !IcebergStorageSchemes.write.contains("oss"), + s"oss must not be in the native write schemes ${IcebergStorageSchemes.write.toSeq.sorted}") + Seq("abfs", "abfss", "wasb", "wasbs", "gcs").foreach { scheme => + assert( + !IcebergStorageSchemes.read.contains(scheme), + s"$scheme must not be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") + assert( + !IcebergStorageSchemes.write.contains(scheme), + s"$scheme must not be in the native write schemes " + + s"${IcebergStorageSchemes.write.toSeq.sorted}") + } + } + test("iceberg gate: builtin allowlist admitted, unbuildable schemes rejected") { - // The Iceberg gate is an explicit allowlist mirroring `storage_factory_for`'s arms (keep in - // lockstep). Narrower than the Parquet gate: object_store recognizes http/abfs/wasb but - // iceberg-rust can't build them, so reject up-front rather than fail during native setup. + // The Iceberg gate is the allowlist the native storage factory publishes over JNI. Narrower + // than the Parquet gate: object_store recognizes http/abfs/wasb but iceberg-rust can't build + // them, so reject up-front rather than fail during native setup. Seq( "file:///tmp/key.parquet", + "memory:///key.parquet", "s3://bucket/key.parquet", "s3a://bucket/key.parquet", "gs://bucket/key.parquet", @@ -218,6 +256,40 @@ class CometScanSchemeFallbackSuite extends CometTestBase { "without opt-in, blob gets no bucket promotion, so a hostless blob location is unopenable") } + test("iceberg scan declines a wasb table instead of failing at execution") { + // End-to-end guard for the scheme gate: object_store recognizes `wasb`, but the native Iceberg + // storage factory has no arm for it. The planner must decline the scan with a reason naming + // the scheme so Spark reads the table, rather than claim it and fail at execution with + // "Unsupported storage scheme: wasb". + assume(icebergAvailable, "Iceberg not available in classpath") + withTempIcebergDir { dir => + val warehouse = s"${FakeWasbSchemeFileSystem.PREFIX}${dir.getAbsolutePath}/warehouse" + withSQLConf( + "spark.sql.catalog.wasb_catalog" -> "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.wasb_catalog.type" -> "hadoop", + "spark.sql.catalog.wasb_catalog.warehouse" -> warehouse, + CometConf.COMET_ENABLED.key -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") { + spark.sql("CREATE TABLE wasb_catalog.db.wasb_table (id INT, name STRING) USING iceberg") + spark.sql( + "INSERT INTO wasb_catalog.db.wasb_table VALUES (1, 'Alice'), (2, 'Bob'), (3, 'Carol')") + try { + val (_, cometPlan) = checkSparkAnswerAndFallbackReason( + "SELECT * FROM wasb_catalog.db.wasb_table ORDER BY id", + "wasb") + val nativeScans = cometPlan.collect { case s: CometIcebergNativeScanExec => s } + assert( + nativeScans.isEmpty, + "`wasb://` has no native Iceberg storage factory arm; the scan must fall back to " + + s"Spark, but Comet claimed it:\n$cometPlan") + } finally { + spark.sql("DROP TABLE wasb_catalog.db.wasb_table") + } + } + } + } + test("native scan claims hdfs:// when libhdfs.schemes is unset (native-default lockstep)") { // Native `is_hdfs_scheme` treats `hdfs` as readable when `fs.comet.libhdfs.schemes` is unset, // so the JVM gate must agree and CLAIM `hdfs://`. Guards the `case None => Set("hdfs")` default From ab9288bbdaf5e95916e88a3425d508da241bdf44 Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Sun, 20 Sep 2026 21:48:12 +0700 Subject: [PATCH 3/4] fix: gate the Iceberg scheme lists on the native list and match schemes verbatim Make builtin_storage_schemes the single point of change by rejecting any scheme it does not list before the factory match, drop memory from the read list since an OpenDAL memory backend is a fresh empty store, and load the JVM sets lazily behind the library-loaded check so a disabled Comet never touches the native library. The JVM gates now match the built-in set verbatim, as OpenDAL strips the scheme prefix case-sensitively at open time. Tests round-trip the JNI lists directly and cover the wasb write gate. --- .../src/execution/operators/iceberg_common.rs | 110 +++++++++++++----- .../comet/iceberg/IcebergStorageSchemes.scala | 48 ++++---- .../apache/comet/rules/CometScanRule.scala | 25 ++-- .../operator/CometIcebergNativeWrite.scala | 14 ++- .../hadoop/fs/FakeWasbSchemeFileSystem.java | 6 +- .../CometIcebergWriteDetectionSuite.scala | 14 +++ .../rules/CometScanSchemeFallbackSuite.scala | 99 +++++++++++++--- 7 files changed, 221 insertions(+), 95 deletions(-) diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 7f9cd1aae47..5860f21216f 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -38,40 +38,31 @@ const ICEBERG_PROVIDER_CLASS_PROPERTY: &str = "s3.comet.credential.provider.clas /// `CometS3CredentialBridge` can read whatever the vendor needs. const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client."]; -/// Pick an OpenDAL storage backend from a URI's scheme. `file` (or no scheme) falls through to -/// the local file system. `memory` is used by the write path to assemble manifest bytes that -/// stay entirely in-process. For S3, the Comet credential bridge is wired in when a provider -/// class is configured; `access_mode` is forwarded to the JVM SPI so the read and write paths can -/// be granted different (e.g. read-only vs read-write) credentials. -/// -/// The JVM planner loads `builtin_storage_schemes` over JNI so it can decline a scan or a write -/// cleanly instead of failing at execution. Opt-in S3-compliant alias schemes come from catalog -/// properties and are added on the JVM side. The tests below keep that list and the match arms -/// here in step. +/// Pick an OpenDAL storage backend for a URI whose scheme `builtin_storage_schemes` lists for +/// `access_mode`, or that is an opted-in S3-compliant alias; anything else is rejected before the +/// match, so the list alone decides what each mode admits. For S3, the Comet credential bridge is +/// wired in when a provider class is configured and `access_mode` is forwarded to the JVM SPI. pub(crate) fn storage_factory_for( path: &str, catalog_properties: &HashMap, catalog_name: &str, access_mode: AccessMode, ) -> Result, DataFusionError> { + // Verbatim match: OpenDAL strips the scheme prefix from every path case-sensitively at open + // time, so admitting `S3://` here would only defer the failure. The JVM gates match verbatim. let scheme = scheme_of(path); + if !builtin_storage_schemes(access_mode).contains(&scheme) + && !is_s3_family_scheme(scheme, catalog_properties) + { + return Err(DataFusionError::Execution(format!( + "Unsupported storage scheme: {scheme}" + ))); + } match scheme { "file" => Ok(Arc::new(OpenDalStorageFactory::Fs)), "memory" => Ok(Arc::new(OpenDalStorageFactory::Memory)), "gs" => Ok(Arc::new(OpenDalStorageFactory::Gcs)), - // Reads keep the OSS backend they have always had (CometScanRule admits `oss` scan - // locations through HadoopFileIO). Writes fail closed: Comet does not forward `oss.*` - // properties into the FileIO and no test covers the write path, so OSS-specific - // endpoint/credential configuration could silently be dropped. The JVM write gate - // already declines `oss` locations; this is the native-side backstop. - "oss" => match access_mode { - AccessMode::Read => Ok(Arc::new(OpenDalStorageFactory::Oss)), - AccessMode::Write => Err(DataFusionError::Execution( - "OSS is not supported for native Iceberg writes (oss.* properties are not \ - forwarded to the native FileIO)" - .to_string(), - )), - }, + "oss" => Ok(Arc::new(OpenDalStorageFactory::Oss)), // s3, s3a, and any opted-in s3-compliant alias (e.g. blob) route to the S3 backend. Listed // last so the built-in backends above stay authoritative even if one of their schemes is // also named in `fs.comet.s3Compliant.schemes`. An alias additionally gets a wrapper that @@ -90,19 +81,20 @@ pub(crate) fn storage_factory_for( })) } } + // Only a listed scheme without an arm reaches here; the tests below fail on that. _ => Err(DataFusionError::Execution(format!( "Unsupported storage scheme: {scheme}" ))), } } -/// The built-in storage schemes `storage_factory_for` admits for `access_mode`, without any -/// opted-in S3-compliant alias. This is the list the JVM read and write gates load over JNI. -/// The rejection test below covers a fixed set of unlisted schemes, so a new factory arm must -/// be added to this list as well or the JVM gate keeps declining it. +/// The single point of change for the built-in schemes `storage_factory_for` admits per access +/// mode; the JVM read and write gates load these lists over JNI. `memory` is write-only: an OpenDAL +/// memory backend is a fresh empty in-process store the write path assembles manifests in, so a +/// read finds nothing. `oss` is read-only: no `oss.*` property is forwarded and no test covers it. pub(crate) fn builtin_storage_schemes(access_mode: AccessMode) -> &'static [&'static str] { match access_mode { - AccessMode::Read => &["file", "memory", "gs", "oss", "s3", "s3a"], + AccessMode::Read => &["file", "gs", "oss", "s3", "s3a"], AccessMode::Write => &["file", "memory", "gs", "s3", "s3a"], } } @@ -263,7 +255,25 @@ mod tests { // oss.* property forwarding exists and is tested. assert!(factory_result("oss://bucket/db/table", AccessMode::Read).is_ok()); let err = factory_result("oss://bucket/db/table", AccessMode::Write).unwrap_err(); - assert!(err.contains("OSS"), "unexpected error: {err}"); + assert!( + err.contains("Unsupported storage scheme: oss"), + "unexpected error: {err}" + ); + } + + #[test] + fn memory_scheme_is_writable_but_not_readable() { + // The write path assembles manifests in a fresh in-process memory store. A read against + // a new memory store can never find data, so the factory must decline it up front. + assert!(factory_result("memory:manifest.avro", AccessMode::Write).is_ok()); + assert!(factory_result("memory:///manifest.avro", AccessMode::Write).is_ok()); + for path in ["memory:manifest.avro", "memory:///key.parquet"] { + let err = factory_result(path, AccessMode::Read).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: memory"), + "unexpected error for {path}: {err}" + ); + } } #[test] @@ -271,13 +281,49 @@ mod tests { for mode in [AccessMode::Read, AccessMode::Write] { assert!(factory_result("file:///tmp/x", mode).is_ok()); assert!(factory_result("/tmp/no-scheme", mode).is_ok()); - assert!(factory_result("memory:manifest.avro", mode).is_ok()); // No credential provider configured: the default chain applies in both modes. assert!(factory_result("s3://bucket/db/table", mode).is_ok()); assert!(factory_result("gs://bucket/db/table", mode).is_ok()); } } + #[test] + fn scheme_list_pre_check_is_load_bearing() { + // The oss and memory arms build a backend for either mode; only their absence from the + // list for the other mode rejects them. Both facts are asserted so that adding an + // access-mode branch back into an arm, or listing the scheme, breaks this test. + assert!(factory_result("oss://bucket/path", AccessMode::Read).is_ok()); + assert!(!builtin_storage_schemes(AccessMode::Write).contains(&"oss")); + let err = factory_result("oss://bucket/path", AccessMode::Write).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: oss"), + "unexpected error: {err}" + ); + assert!(factory_result("memory:///path", AccessMode::Write).is_ok()); + assert!(!builtin_storage_schemes(AccessMode::Read).contains(&"memory")); + let err = factory_result("memory:///path", AccessMode::Read).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: memory"), + "unexpected error: {err}" + ); + } + + #[test] + fn mixed_case_scheme_is_rejected_for_both_modes() { + // OpenDAL strips the scheme prefix from every path case-sensitively at open time + // (`S3://bucket/key` fails its `s3://bucket/` prefix check), so the factory must not + // admit what the open cannot serve. The JVM gates match the built-in set verbatim too. + for mode in [AccessMode::Read, AccessMode::Write] { + for path in ["S3://bucket/key", "File:///tmp/x", "GS://bucket/key"] { + let err = factory_result(path, mode).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {path} in {mode:?}: {err}" + ); + } + } + } + #[test] fn unknown_scheme_is_rejected() { let err = factory_result("hdfs://nn/db/table", AccessMode::Read).unwrap_err(); @@ -288,7 +334,7 @@ mod tests { } #[test] - fn exposed_scheme_list_matches_storage_factory() { + fn listed_schemes_are_accepted_by_storage_factory() { for mode in [AccessMode::Read, AccessMode::Write] { for scheme in builtin_storage_schemes(mode) { let url = format!("{scheme}://bucket/path"); @@ -300,7 +346,7 @@ mod tests { } assert!(builtin_storage_schemes(AccessMode::Read).contains(&"oss")); assert!(!builtin_storage_schemes(AccessMode::Write).contains(&"oss")); - assert!(builtin_storage_schemes(AccessMode::Read).contains(&"memory")); + assert!(!builtin_storage_schemes(AccessMode::Read).contains(&"memory")); assert!(builtin_storage_schemes(AccessMode::Write).contains(&"memory")); } diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala index e96d7792745..8d0ec85b1f0 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala @@ -21,43 +21,27 @@ package org.apache.comet.iceberg import java.util.Locale -import scala.util.control.NonFatal - import org.apache.spark.internal.Logging import org.apache.comet.NativeBase /** * The storage schemes the native Iceberg storage factory publishes over JNI, so the JVM scan and - * write gates decline what native cannot open instead of failing at execution. The fallback - * constants are used only when the library cannot be loaded: assuming a scheme is supported would - * recreate the execution-time failure, and without the library nothing runs natively anyway. + * write gates decline what native cannot open instead of failing at execution. Every caller sits + * behind `isCometLoaded`, so the fallback constants are only consulted in a JVM where nothing + * runs natively; the pinning test in `CometScanSchemeFallbackSuite` keeps them equal to the + * native lists. */ private[comet] object IcebergStorageSchemes extends Logging { - private[comet] val FallbackRead: Set[String] = Set("file", "memory", "s3", "s3a", "gs", "oss") + private[comet] val FallbackRead: Set[String] = Set("file", "s3", "s3a", "gs", "oss") private[comet] val FallbackWrite: Set[String] = Set("file", "memory", "s3", "s3a", "gs") lazy val read: Set[String] = load(forWrite = false, FallbackRead) lazy val write: Set[String] = load(forWrite = true, FallbackWrite) - private def load(forWrite: Boolean, fallback: Set[String]): Set[String] = { - val path = if (forWrite) "write" else "read" - val joined = - try NativeBase.icebergStorageSchemes(forWrite) - catch { - case e: UnsatisfiedLinkError => - logWarning( - s"Comet native library is not loaded; using the fallback Iceberg $path scheme list " + - s"${fallback.toSeq.sorted.mkString(", ")}: ${e.getMessage}") - return fallback - case NonFatal(e) => - logWarning( - s"Failed to load the Iceberg $path scheme list from the Comet native library; using " + - s"the fallback list ${fallback.toSeq.sorted.mkString(", ")}", - e) - return fallback - } + /** Splits the comma-joined native list; a null or blank list yields `fallback`. */ + private[comet] def parse(joined: String, fallback: Set[String]): Set[String] = { val schemes = Option(joined).toSeq .flatMap(_.split(",")) .map(_.trim.toLowerCase(Locale.ROOT)) @@ -65,11 +49,25 @@ private[comet] object IcebergStorageSchemes extends Logging { .toSet if (schemes.isEmpty) { logWarning( - s"Comet native library published an empty Iceberg $path scheme list; using the fallback " + - s"list ${fallback.toSeq.sorted.mkString(", ")}") + "Comet native library published an empty Iceberg scheme list; using the fallback list " + + fallback.toSeq.sorted.mkString(", ")) fallback } else { schemes } } + + // A native fault while answering the probe propagates: `isLoaded` true means the symbol + // resolves, and anything else is a build bug that must fail loudly rather than fall back. + private def load(forWrite: Boolean, fallback: Set[String]): Set[String] = { + if (NativeBase.isLoaded) { + parse(NativeBase.icebergStorageSchemes(forWrite), fallback) + } else { + val path = if (forWrite) "write" else "read" + logWarning( + s"Comet native library is not loaded; using the fallback Iceberg $path scheme list " + + fallback.toSeq.sorted.mkString(", ")) + fallback + } + } } diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index 95ec5390caa..b3a6d8ae510 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -1193,11 +1193,21 @@ object CometScanRule extends Logging { /** * Schemes Comet's native Iceberg scan can open, loaded from the native storage factory over JNI - * so this gate cannot drift from `storage_factory_for`. Opt-in aliases from + * so this gate cannot drift from `storage_factory_for`. Lazy so that constructing the rule does + * not touch the native library before `isCometLoaded` has been consulted. Opt-in aliases from * `fs.comet.s3Compliant.schemes` are additive (see `isIcebergReadableScheme`); the write path * loads its own set (`CometIcebergNativeWrite.SupportedStorageSchemes`). */ - private val icebergReadableSchemes: Set[String] = IcebergStorageSchemes.read + private lazy val icebergReadableSchemes: Set[String] = IcebergStorageSchemes.read + + /** + * True when the Iceberg scan gate admits `scheme`. The built-in set matches verbatim because + * native opens a location by its raw scheme and OpenDAL strips that prefix case-sensitively, so + * `S3://` is not `s3://`; the opt-in alias list is matched case-insensitively on both sides. + */ + private def isAdmittedIcebergScheme(scheme: String, s3CompliantSchemes: Set[String]): Boolean = + icebergReadableSchemes.contains(scheme) || + s3CompliantSchemes.contains(scheme.toLowerCase(Locale.ROOT)) /** * "Supported schemes: ..." suffix shared by the Iceberg scheme-fallback messages. Lists the @@ -1218,8 +1228,7 @@ object CometScanRule extends Logging { s3CompliantSchemes: Set[String]): Boolean = { val scheme = uri.getScheme if (scheme == null) return true - val lower = scheme.toLowerCase(Locale.ROOT) - icebergReadableSchemes.contains(lower) || s3CompliantSchemes.contains(lower) + isAdmittedIcebergScheme(scheme, s3CompliantSchemes) } /** @@ -1281,9 +1290,6 @@ object CometScanRule extends Logging { // hasOpenableAuthority); non-empty => decline. One example suffices for the message. var hostlessLocation: Option[String] = None - // Union of the two admitted scheme sets, built once so the per-file loop does one lookup. - val openableSchemes = icebergReadableSchemes ++ s3CompliantSchemes - // Classify one data/delete file location; see `icebergReadableSchemes` for why that allowlist // is narrower than the Parquet native gate. Runs per data and delete file, so the scheme and // the bucket are each derived once and threaded down. @@ -1294,9 +1300,8 @@ object CometScanRule extends Logging { // A schemeless local path routes to iceberg-rust's LocalFs and needs no host. val scheme = uri.getScheme if (scheme == null) return - val lower = scheme.toLowerCase(Locale.ROOT) - if (!openableSchemes.contains(lower)) { - unsupportedSchemes += lower + if (!isAdmittedIcebergScheme(scheme, s3CompliantSchemes)) { + unsupportedSchemes += scheme } else if (!hasOpenableAuthority(uri, s3CompliantSchemes)) { if (hostlessLocation.isEmpty) hostlessLocation = Some(rawPath) } else { diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala index 1209f5af640..f2276dc11e3 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala @@ -81,11 +81,11 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { private val EncryptionPropertyPrefix = "encryption." private val UnsupportedWriteTypeIds: Set[String] = Set("UUID") - // `oss` is deliberately absent: iceberg-rust has an OSS backend, but Comet does not forward - // `oss.*` catalog properties to it and no functional test covers the path, so an OSS write - // could silently drop endpoint/credential configuration. Fail closed until it is covered. - // `gs` is additionally gated on the resolved FileIO (`requireGcsFileIOForGcsDataLocation`). - private val SupportedStorageSchemes: Set[String] = IcebergStorageSchemes.write + // Loaded from the native storage factory: `builtin_storage_schemes` in + // `native/core/src/execution/operators/iceberg_common.rs` is the single point of change and + // explains why `oss` is read-only and `memory` write-only. Lazy so constructing the serde does + // not touch the native library. `gs` is additionally gated on the resolved FileIO below. + private lazy val SupportedStorageSchemes: Set[String] = IcebergStorageSchemes.write private val MinUnsupportedFormatVersion = 3 private val ParquetWritePropertyPrefix = "write.parquet." private val ParquetMrPropertyPrefix = "parquet." @@ -306,9 +306,11 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { .find(k => !IgnoredHadoopParquetConfKeys.contains(k)) .map(k => s"Hadoop configuration sets $k (reaches iceberg-java's writer but not native)") + // Verbatim, not lowercased: native opens the data location by its raw scheme and OpenDAL strips + // that prefix case-sensitively, so `S3://` must be declined here rather than fail at execution. private def storageScheme(location: String): String = if (location.contains("://")) { - location.substring(0, location.indexOf("://")).toLowerCase(Locale.ROOT) + location.substring(0, location.indexOf("://")) } else { "file" } diff --git a/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java index 02830e5408b..95ff00b818a 100644 --- a/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java +++ b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java @@ -25,9 +25,9 @@ /** * A local-disk-backed FileSystem that reports the {@code wasb} scheme, so a test can create and - * read an Iceberg table under a {@code wasb://} warehouse without Azure. Used to assert that {@code - * CometScanRule} declines a native Iceberg scan whose storage scheme object_store recognizes but - * the native Iceberg storage factory cannot build, instead of failing at execution. + * read an Iceberg table under a {@code wasb://} warehouse without Azure. Used to assert that the + * Iceberg scan and write gates decline a table whose storage scheme the native Iceberg storage + * factory has no arm for, instead of claiming it and failing at execution. */ public class FakeWasbSchemeFileSystem extends RawLocalFileSystem { diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index 37ddb1872e8..458098209f9 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -444,6 +444,20 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes } } + test("fall-back: mixed-case data location scheme (native opens the location verbatim)") { + // OpenDAL strips the scheme prefix from a path case-sensitively, so `S3://` cannot be opened + // natively even though `s3://` can. The gate must match the scheme verbatim and decline it. + withDetectionCatalog { dir => + createTable( + dir, + "mixed_case_scheme", + partitionSpec = "", + properties = + Some("'write.data.path'='S3://nonexistent-bucket/iceberg/db/mixed_case_scheme'")) + assertUnsupportedContainsAllowingWriteFailure("mixed_case_scheme", "storage scheme", "S3") + } + } + test("Compatible when the data location scheme is s3") { withDetectionCatalog { dir => createTable( diff --git a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala index 9de3c129184..b0e6714f203 100644 --- a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala +++ b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala @@ -27,10 +27,10 @@ import java.util.UUID import org.apache.commons.io.FileUtils import org.apache.spark.SparkConf import org.apache.spark.sql.{CometTestBase, SaveMode} -import org.apache.spark.sql.comet.{CometIcebergNativeScanExec, CometScanExec} +import org.apache.spark.sql.comet.{CometIcebergNativeScanExec, CometIcebergWriteExec, CometScanExec} import org.apache.spark.sql.execution.{FileSourceScanExec, SparkPlan} -import org.apache.comet.{CometConf, CometIcebergTestBase, NativeBase} +import org.apache.comet.{CometConf, CometIcebergTestBase, ExtendedExplainInfo, NativeBase} import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem, FakeWasbSchemeFileSystem} import org.apache.comet.iceberg.IcebergStorageSchemes @@ -124,6 +124,30 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa "gs://bucket/... must be admitted even after an authorityless gs:// URI was probed first") } + test("iceberg gate: the JNI scheme list parser trims, lowercases and falls back") { + val fallback = Set("fallback") + assert(IcebergStorageSchemes.parse("file, S3,,s3a,", fallback) == Set("file", "s3", "s3a")) + assert(IcebergStorageSchemes.parse("", fallback) == fallback) + assert(IcebergStorageSchemes.parse(" , ", fallback) == fallback) + assert(IcebergStorageSchemes.parse(null, fallback) == fallback) + } + + test("iceberg gate: JNI scheme lists round-trip to the JVM fallback constants") { + // Calls the JNI probe directly, bypassing the lazy vals, so this cannot pass on a JVM whose + // `read`/`write` quietly came from the fallback rather than from native. + assume(NativeBase.isLoaded, "Comet native library not loaded") + val read = IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(false), Set.empty) + val write = IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(true), Set.empty) + assert( + read == IcebergStorageSchemes.FallbackRead, + s"JNI read schemes ${read.toSeq.sorted} differ from the JVM fallback " + + s"${IcebergStorageSchemes.FallbackRead.toSeq.sorted}") + assert( + write == IcebergStorageSchemes.FallbackWrite, + s"JNI write schemes ${write.toSeq.sorted} differ from the JVM fallback " + + s"${IcebergStorageSchemes.FallbackWrite.toSeq.sorted}") + } + test("iceberg gate: native scheme lists match the JVM fallback constants") { // The gates load their scheme sets from the native storage factory over JNI. The JVM // fallback constants only stand in when the library cannot be loaded, so this pins them to @@ -138,14 +162,20 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa IcebergStorageSchemes.write == IcebergStorageSchemes.FallbackWrite, s"native write schemes ${IcebergStorageSchemes.write.toSeq.sorted} differ from the JVM " + s"fallback ${IcebergStorageSchemes.FallbackWrite.toSeq.sorted}") - Seq("memory", "oss").foreach { scheme => - assert( - IcebergStorageSchemes.read.contains(scheme), - s"$scheme must be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") - } + // `oss` is read-only (no `oss.*` property forwarding for writes) and `memory` is write-only + // (a fresh in-process store per FileIO, so a read finds nothing). + assert( + IcebergStorageSchemes.read.contains("oss"), + s"oss must be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") assert( !IcebergStorageSchemes.write.contains("oss"), s"oss must not be in the native write schemes ${IcebergStorageSchemes.write.toSeq.sorted}") + assert( + IcebergStorageSchemes.write.contains("memory"), + s"memory must be in the native write schemes ${IcebergStorageSchemes.write.toSeq.sorted}") + assert( + !IcebergStorageSchemes.read.contains("memory"), + s"memory must not be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") Seq("abfs", "abfss", "wasb", "wasbs", "gcs").foreach { scheme => assert( !IcebergStorageSchemes.read.contains(scheme), @@ -159,11 +189,12 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa test("iceberg gate: builtin allowlist admitted, unbuildable schemes rejected") { // The Iceberg gate is the allowlist the native storage factory publishes over JNI. Narrower - // than the Parquet gate: object_store recognizes http/abfs/wasb but iceberg-rust can't build - // them, so reject up-front rather than fail during native setup. + // than the Parquet gate: object_store recognizes http and the Azure schemes, but the native + // Iceberg storage factory has no arm for them, so reject up-front rather than fail at + // execution. `memory` is write-only (a fresh in-process store per FileIO holds no table), and + // the built-in set matches verbatim because native opens a location by its raw scheme. Seq( "file:///tmp/key.parquet", - "memory:///key.parquet", "s3://bucket/key.parquet", "s3a://bucket/key.parquet", "gs://bucket/key.parquet", @@ -173,6 +204,9 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa s"$u must be iceberg-readable; icebergReadableSchemes has regressed") } Seq( + "memory:///key.parquet", + "S3://bucket/key.parquet", + "File:///tmp/key.parquet", "http://bucket.example.com/key.parquet", "https://bucket.example.com/key.parquet", "abfs://container@acct/key.parquet", @@ -181,7 +215,7 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa "wasbs://container@acct/key.parquet").foreach { u => assert( !CometScanRule.isIcebergReadableScheme(new URI(u), Set.empty), - s"$u must not be iceberg-readable; storage_factory_for has no matching arm") + s"$u must not be iceberg-readable; storage_factory_for rejects it for reads") } } @@ -256,11 +290,11 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa "without opt-in, blob gets no bucket promotion, so a hostless blob location is unopenable") } - test("iceberg scan declines a wasb table instead of failing at execution") { - // End-to-end guard for the scheme gate: object_store recognizes `wasb`, but the native Iceberg - // storage factory has no arm for it. The planner must decline the scan with a reason naming - // the scheme so Spark reads the table, rather than claim it and fail at execution with - // "Unsupported storage scheme: wasb". + test("iceberg scan and write decline a wasb table instead of failing at execution") { + // Regression guard for both scheme gates: the native Iceberg storage factory has no arm for + // `wasb`, so the scan and the write must each be declined with a reason naming the scheme + // and run in Spark, rather than be claimed and fail at execution with "Unsupported storage + // scheme: wasb". assume(icebergAvailable, "Iceberg not available in classpath") withTempIcebergDir { dir => val warehouse = s"${FakeWasbSchemeFileSystem.PREFIX}${dir.getAbsolutePath}/warehouse" @@ -270,11 +304,31 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa "spark.sql.catalog.wasb_catalog.warehouse" -> warehouse, CometConf.COMET_ENABLED.key -> "true", CometConf.COMET_EXEC_ENABLED.key -> "true", - CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") { + CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "true", + // The VALUES source must be native, or the write is declined for its input before the + // write serde (and its scheme gate) is ever consulted. + CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key -> "true") { spark.sql("CREATE TABLE wasb_catalog.db.wasb_table (id INT, name STRING) USING iceberg") - spark.sql( - "INSERT INTO wasb_catalog.db.wasb_table VALUES (1, 'Alice'), (2, 'Bob'), (3, 'Carol')") try { + val insertPlans = capturePlans(spark) { + spark.sql( + "INSERT INTO wasb_catalog.db.wasb_table VALUES (1, 'Alice'), (2, 'Bob'), (3, 'Carol')") + } + assert(insertPlans.nonEmpty, "no INSERT plan was captured") + val nativeWrites = insertPlans.flatMap(_.collectWithSubqueries { + case w: CometIcebergWriteExec => w + }) + assert( + nativeWrites.isEmpty, + "`wasb://` has no native Iceberg storage factory arm; the write must fall back to " + + s"Spark, but Comet claimed it:\n${insertPlans.mkString("\n--\n")}") + val writeReasons = insertPlans.flatMap(new ExtendedExplainInfo().getFallbackReasons) + assert( + writeReasons.exists(_.contains("unsupported storage scheme: wasb")), + s"the write must be declined for its scheme, got reasons: $writeReasons") + val (_, cometPlan) = checkSparkAnswerAndFallbackReason( "SELECT * FROM wasb_catalog.db.wasb_table ORDER BY id", "wasb") @@ -283,6 +337,13 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa nativeScans.isEmpty, "`wasb://` has no native Iceberg storage factory arm; the scan must fall back to " + s"Spark, but Comet claimed it:\n$cometPlan") + // Nothing but the scheme may have caused the fallback: every reason names `wasb`, + // except the one the operator above records about its now-Spark BatchScanExec child. + val scanReasons = new ExtendedExplainInfo().getFallbackReasons(cometPlan) + val unrelated = scanReasons.filterNot(_.contains("wasb")) + assert( + unrelated.forall(_.contains("BatchScanExec")), + s"fallback reasons other than the scheme were recorded: $unrelated") } finally { spark.sql("DROP TABLE wasb_catalog.db.wasb_table") } From 1d3f95d6bb42858ff85012eb8d4f4f64f8077e7c Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Tue, 22 Sep 2026 23:25:23 +0700 Subject: [PATCH 4/4] fix: drop the fallback scheme lists and match Iceberg aliases verbatim The JVM scheme sets are only read after the native library has loaded, so the hand-maintained fallback lists never did work. `load` now returns the empty set when the library is not loaded, and a blank native list parses to the empty set, each with a warning. iceberg-rust builds the S3 prefix from the parsed scheme, which is lowercase, and compares it to the raw path, so a mixed-case alias location passed the gates and failed at open. Aliases are now matched as written on the Iceberg path, on the JVM gate and in the native factory, like the built-in schemes. The data sources guide says an Iceberg location must spell the alias in lowercase. --- docs/source/user-guide/latest/datasources.md | 5 +- .../src/execution/operators/iceberg_common.rs | 62 +++++++++- .../parquet/objectstore/s3_blob_fs_support.rs | 6 +- .../comet/iceberg/IcebergStorageSchemes.scala | 56 ++++----- .../apache/comet/rules/CometScanRule.scala | 13 ++- .../rules/CometScanSchemeFallbackSuite.scala | 109 ++++++++++-------- 6 files changed, 162 insertions(+), 89 deletions(-) diff --git a/docs/source/user-guide/latest/datasources.md b/docs/source/user-guide/latest/datasources.md index eff37b0a1ee..4c8482b566d 100644 --- a/docs/source/user-guide/latest/datasources.md +++ b/docs/source/user-guide/latest/datasources.md @@ -257,7 +257,10 @@ URLs. This is opt-in and disabled by default. Enable it by listing the schemes to treat as S3-compliant aliases in `spark.hadoop.fs.comet.s3Compliant.schemes` (Hadoop key `fs.comet.s3Compliant.schemes`), a comma-separated, case-insensitive list. This mirrors the -existing `fs.comet.libhdfs.schemes` config. +existing `fs.comet.libhdfs.schemes` config. The list entries are case-insensitive, but an Iceberg +table location must write the alias in lowercase (`blob://`, not `BLOB://`): the native Iceberg +reader opens each recorded location as written, and its S3 backend accepts only a lowercase +scheme prefix, so Comet declines such a location up front and leaves that scan to Spark. ```shell --conf spark.hadoop.fs.comet.s3Compliant.schemes=blob diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 5860f21216f..932fd5e5de9 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -39,9 +39,10 @@ const ICEBERG_PROVIDER_CLASS_PROPERTY: &str = "s3.comet.credential.provider.clas const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client."]; /// Pick an OpenDAL storage backend for a URI whose scheme `builtin_storage_schemes` lists for -/// `access_mode`, or that is an opted-in S3-compliant alias; anything else is rejected before the -/// match, so the list alone decides what each mode admits. For S3, the Comet credential bridge is -/// wired in when a provider class is configured and `access_mode` is forwarded to the JVM SPI. +/// `access_mode`, or that is an opted-in S3-compliant alias written in lowercase; anything else is +/// rejected before the match, so the list alone decides what each mode admits. For S3, the Comet +/// credential bridge is wired in when a provider class is configured and `access_mode` is +/// forwarded to the JVM SPI. pub(crate) fn storage_factory_for( path: &str, catalog_properties: &HashMap, @@ -49,7 +50,8 @@ pub(crate) fn storage_factory_for( access_mode: AccessMode, ) -> Result, DataFusionError> { // Verbatim match: OpenDAL strips the scheme prefix from every path case-sensitively at open - // time, so admitting `S3://` here would only defer the failure. The JVM gates match verbatim. + // time, so admitting `S3://` here would only defer the failure. Aliases are held to the same + // rule (`is_iceberg_alias_scheme`). The JVM gates match verbatim. let scheme = scheme_of(path); if !builtin_storage_schemes(access_mode).contains(&scheme) && !is_s3_family_scheme(scheme, catalog_properties) @@ -71,7 +73,7 @@ pub(crate) fn storage_factory_for( s if is_s3_family_scheme(s, catalog_properties) => { let customized_credential_load = build_s3_credential_loader(path, catalog_properties, catalog_name, access_mode)?; - if is_s3_compliant_alias_scheme(s, catalog_properties) { + if is_iceberg_alias_scheme(s, catalog_properties) { Ok(Arc::new(BlobHostPromotingS3StorageFactory::new( customized_credential_load, ))) @@ -235,7 +237,19 @@ fn scheme_of(path: &str) -> &str { /// iceberg-rust's storage factory has no `s3n` backend and the Scala Iceberg scheme gate /// (`isIcebergReadableScheme`) rejects `s3n` before a path ever reaches this operator. fn is_s3_family_scheme(scheme: &str, catalog_properties: &HashMap) -> bool { - matches!(scheme, "s3" | "s3a") || is_s3_compliant_alias_scheme(scheme, catalog_properties) + matches!(scheme, "s3" | "s3a") || is_iceberg_alias_scheme(scheme, catalog_properties) +} + +/// True if `scheme` is an opted-in S3-compliant alias written exactly as the lowercase form of its +/// list entry. The shared `is_s3_compliant_alias_scheme` is case-insensitive, which is safe for +/// the Parquet path because it rewrites an alias URL to `s3://` before anything opens it. The +/// Iceberg path opens the recorded location as written, and OpenDAL's S3 backend checks it against +/// a `scheme://bucket/` prefix whose scheme comes from `Url::parse` and so is lowercase, so +/// `BLOB://bucket/key` would pass a case-insensitive gate here and fail at open time. The JVM +/// Iceberg gate matches the same way. +fn is_iceberg_alias_scheme(scheme: &str, catalog_properties: &HashMap) -> bool { + !scheme.bytes().any(|b| b.is_ascii_uppercase()) + && is_s3_compliant_alias_scheme(scheme, catalog_properties) } #[cfg(test)] @@ -324,6 +338,42 @@ mod tests { } } + #[test] + fn alias_scheme_is_matched_verbatim() { + // An alias location is opened as written, and OpenDAL's S3 backend checks it against a + // lowercase `scheme://bucket/` prefix, so a mixed-case alias that passed a + // case-insensitive gate here would only fail at open time. The list entry may be written + // in any case; the location's scheme must match its lowercase form exactly. + for listed in ["blob", " BLOB ", "minio,Blob"] { + let props = HashMap::from([( + "fs.comet.s3Compliant.schemes".to_string(), + listed.to_string(), + )]); + for mode in [AccessMode::Read, AccessMode::Write] { + assert!( + storage_factory_for("blob://bucket/key", &props, "test_cat", mode).is_ok(), + "blob://bucket/key must be admitted for {mode:?} with list {listed:?}" + ); + for path in [ + "BLOB://bucket/key", + "Blob://bucket/key", + "BLOB:///bucket/key", + ] { + let err = storage_factory_for(path, &props, "test_cat", mode) + .map(|_| ()) + .expect_err(&format!( + "{path} must be rejected for {mode:?} with list {listed:?}" + )) + .to_string(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {path} in {mode:?}: {err}" + ); + } + } + } + } + #[test] fn unknown_scheme_is_rejected() { let err = factory_result("hdfs://nn/db/table", AccessMode::Read).unwrap_err(); diff --git a/native/core/src/parquet/objectstore/s3_blob_fs_support.rs b/native/core/src/parquet/objectstore/s3_blob_fs_support.rs index 0ec4c572519..87ba1d0be99 100644 --- a/native/core/src/parquet/objectstore/s3_blob_fs_support.rs +++ b/native/core/src/parquet/objectstore/s3_blob_fs_support.rs @@ -98,8 +98,10 @@ pub(crate) fn normalize_object_store_url( /// True if `scheme` is a configured s3-compliant alias (`fs.comet.s3Compliant.schemes`; empty or /// unset means none). These route to `AmazonS3` but are not recognized by /// `ObjectStoreScheme::parse`. `s3a` is excluded -- object_store knows it and callers special-case -/// it. Shared by the Parquet normalizer above and the Iceberg storage-factory gate -/// (`iceberg_common::storage_factory_for`) so both admit the same schemes. +/// it. Case-insensitive, which is safe for the Parquet normalizer above because it rewrites the +/// alias to `s3://` before anything opens the path. The Iceberg storage-factory gate opens the +/// recorded location as written, so it wraps this in `iceberg_common::is_iceberg_alias_scheme`, +/// which additionally requires the scheme to be lowercase. pub(crate) fn is_s3_compliant_alias_scheme( scheme: &str, object_store_configs: &HashMap, diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala index 8d0ec85b1f0..51f4b5f7019 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala @@ -27,47 +27,51 @@ import org.apache.comet.NativeBase /** * The storage schemes the native Iceberg storage factory publishes over JNI, so the JVM scan and - * write gates decline what native cannot open instead of failing at execution. Every caller sits - * behind `isCometLoaded`, so the fallback constants are only consulted in a JVM where nothing - * runs natively; the pinning test in `CometScanSchemeFallbackSuite` keeps them equal to the - * native lists. + * write gates decline what native cannot open instead of failing at execution. Native is the only + * source of these lists. Every caller sits behind `isCometLoaded`, so a JVM where the library is + * not loaded never plans a native scan or write, and both sets are simply empty there. */ private[comet] object IcebergStorageSchemes extends Logging { - private[comet] val FallbackRead: Set[String] = Set("file", "s3", "s3a", "gs", "oss") - private[comet] val FallbackWrite: Set[String] = Set("file", "memory", "s3", "s3a", "gs") + lazy val read: Set[String] = load(forWrite = false) + lazy val write: Set[String] = load(forWrite = true) - lazy val read: Set[String] = load(forWrite = false, FallbackRead) - lazy val write: Set[String] = load(forWrite = true, FallbackWrite) - - /** Splits the comma-joined native list; a null or blank list yields `fallback`. */ - private[comet] def parse(joined: String, fallback: Set[String]): Set[String] = { + /** + * Splits the comma-joined native list, trimming and lowercasing each entry. A null or blank + * list throws: `builtin_storage_schemes` is a fixed non-empty constant natively, so a loaded + * library that publishes nothing can only be a build or marshalling bug, and a warning here + * would silently disable every native Iceberg scan and write. + */ + private[comet] def parse(joined: String, forWrite: Boolean): Set[String] = { val schemes = Option(joined).toSeq .flatMap(_.split(",")) .map(_.trim.toLowerCase(Locale.ROOT)) .filter(_.nonEmpty) .toSet if (schemes.isEmpty) { - logWarning( - "Comet native library published an empty Iceberg scheme list; using the fallback list " + - fallback.toSeq.sorted.mkString(", ")) - fallback - } else { - schemes + val path = if (forWrite) "write" else "read" + throw new IllegalStateException( + s"Comet native library published no Iceberg $path scheme list") } + schemes } - // A native fault while answering the probe propagates: `isLoaded` true means the symbol - // resolves, and anything else is a build bug that must fail loudly rather than fall back. - private def load(forWrite: Boolean, fallback: Set[String]): Set[String] = { - if (NativeBase.isLoaded) { - parse(NativeBase.icebergStorageSchemes(forWrite), fallback) + /** + * Loads one list over JNI, or answers the empty set with a warning when the library is not + * loaded; that branch never runs a plan, because every caller sits behind `isCometLoaded`. + * `isLoaded` is a parameter so the unloaded answer can be tested without unloading the library. + * Once loaded, anything but a non-empty list is a build bug that must fail loudly: a native + * fault while answering the probe propagates, and a blank answer throws in `parse`. + */ + private[comet] def load( + forWrite: Boolean, + isLoaded: Boolean = NativeBase.isLoaded): Set[String] = { + if (isLoaded) { + parse(NativeBase.icebergStorageSchemes(forWrite), forWrite) } else { val path = if (forWrite) "write" else "read" - logWarning( - s"Comet native library is not loaded; using the fallback Iceberg $path scheme list " + - fallback.toSeq.sorted.mkString(", ")) - fallback + logWarning(s"Comet native library is not loaded; the Iceberg $path scheme list is empty") + Set.empty } } } diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index b3a6d8ae510..2dd04c25a5a 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -1201,13 +1201,16 @@ object CometScanRule extends Logging { private lazy val icebergReadableSchemes: Set[String] = IcebergStorageSchemes.read /** - * True when the Iceberg scan gate admits `scheme`. The built-in set matches verbatim because - * native opens a location by its raw scheme and OpenDAL strips that prefix case-sensitively, so - * `S3://` is not `s3://`; the opt-in alias list is matched case-insensitively on both sides. + * True when the Iceberg scan gate admits `scheme`, written exactly as recorded. Native opens a + * location by its raw scheme, and OpenDAL's S3 backend checks it against a `scheme://bucket/` + * prefix whose scheme comes from `Url::parse` and so is lowercase: `S3://` is not `s3://`, and + * `BLOB://` is not an opted-in `blob`. The alias set is lowercase + * (`NativeConfig.parseSchemeSet`), so an alias is admitted only when the location writes it in + * lowercase. The Parquet gate stays case-insensitive because it rewrites alias URLs to `s3://` + * before anything opens them. */ private def isAdmittedIcebergScheme(scheme: String, s3CompliantSchemes: Set[String]): Boolean = - icebergReadableSchemes.contains(scheme) || - s3CompliantSchemes.contains(scheme.toLowerCase(Locale.ROOT)) + icebergReadableSchemes.contains(scheme) || s3CompliantSchemes.contains(scheme) /** * "Supported schemes: ..." suffix shared by the Iceberg scheme-fallback messages. Lists the diff --git a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala index b0e6714f203..21011fae65d 100644 --- a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala +++ b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala @@ -124,66 +124,56 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa "gs://bucket/... must be admitted even after an authorityless gs:// URI was probed first") } - test("iceberg gate: the JNI scheme list parser trims, lowercases and falls back") { - val fallback = Set("fallback") - assert(IcebergStorageSchemes.parse("file, S3,,s3a,", fallback) == Set("file", "s3", "s3a")) - assert(IcebergStorageSchemes.parse("", fallback) == fallback) - assert(IcebergStorageSchemes.parse(" , ", fallback) == fallback) - assert(IcebergStorageSchemes.parse(null, fallback) == fallback) + test("iceberg gate: the JNI scheme list parser trims, lowercases and rejects a blank list") { + assert( + IcebergStorageSchemes.parse("file, S3,,s3a,", forWrite = false) == Set("file", "s3", "s3a")) + // A loaded library publishes a fixed non-empty list, so a blank answer is a build bug and + // must fail loudly rather than quietly disable every native Iceberg scan and write. + Seq(("", false), (" , ", true), (null, false)).foreach { case (joined, forWrite) => + val path = if (forWrite) "write" else "read" + val e = intercept[IllegalStateException](IcebergStorageSchemes.parse(joined, forWrite)) + assert( + e.getMessage == s"Comet native library published no Iceberg $path scheme list", + s"unexpected message for ${Option(joined)}: ${e.getMessage}") + } } - test("iceberg gate: JNI scheme lists round-trip to the JVM fallback constants") { - // Calls the JNI probe directly, bypassing the lazy vals, so this cannot pass on a JVM whose - // `read`/`write` quietly came from the fallback rather than from native. - assume(NativeBase.isLoaded, "Comet native library not loaded") - val read = IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(false), Set.empty) - val write = IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(true), Set.empty) - assert( - read == IcebergStorageSchemes.FallbackRead, - s"JNI read schemes ${read.toSeq.sorted} differ from the JVM fallback " + - s"${IcebergStorageSchemes.FallbackRead.toSeq.sorted}") - assert( - write == IcebergStorageSchemes.FallbackWrite, - s"JNI write schemes ${write.toSeq.sorted} differ from the JVM fallback " + - s"${IcebergStorageSchemes.FallbackWrite.toSeq.sorted}") + test("iceberg gate: an unloaded native library yields no schemes") { + // Every caller sits behind `isCometLoaded`, so this answer never reaches a plan; it only + // guarantees that nothing hand-maintained stands in for the native list. + assert(IcebergStorageSchemes.load(forWrite = false, isLoaded = false) == Set.empty) + assert(IcebergStorageSchemes.load(forWrite = true, isLoaded = false) == Set.empty) } - test("iceberg gate: native scheme lists match the JVM fallback constants") { - // The gates load their scheme sets from the native storage factory over JNI. The JVM - // fallback constants only stand in when the library cannot be loaded, so this pins them to - // what native publishes: extending `builtin_storage_schemes` without updating the fallback - // fails here rather than drifting silently. + test("iceberg gate: the scheme sets come from native and carry its mode-specific entries") { + // The lazy sets must be what the JNI probe answers, not an empty set from a load that ran + // before the library was ready. `oss` is read-only (no `oss.*` property forwarding for + // writes) and `memory` is write-only (a fresh in-process store per FileIO, so a read finds + // nothing); the Azure schemes and `gcs` have no storage factory arm at all. assume(NativeBase.isLoaded, "Comet native library not loaded") + val read = IcebergStorageSchemes.read + val write = IcebergStorageSchemes.write + assert(read == IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(false), false)) + assert(write == IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(true), true)) + assert(read.nonEmpty, "native published no read schemes") + assert(write.nonEmpty, "native published no write schemes") + assert(read.contains("oss"), s"oss must be in the native read schemes ${read.toSeq.sorted}") assert( - IcebergStorageSchemes.read == IcebergStorageSchemes.FallbackRead, - s"native read schemes ${IcebergStorageSchemes.read.toSeq.sorted} differ from the JVM " + - s"fallback ${IcebergStorageSchemes.FallbackRead.toSeq.sorted}") - assert( - IcebergStorageSchemes.write == IcebergStorageSchemes.FallbackWrite, - s"native write schemes ${IcebergStorageSchemes.write.toSeq.sorted} differ from the JVM " + - s"fallback ${IcebergStorageSchemes.FallbackWrite.toSeq.sorted}") - // `oss` is read-only (no `oss.*` property forwarding for writes) and `memory` is write-only - // (a fresh in-process store per FileIO, so a read finds nothing). - assert( - IcebergStorageSchemes.read.contains("oss"), - s"oss must be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") + !write.contains("oss"), + s"oss must not be in the native write schemes ${write.toSeq.sorted}") assert( - !IcebergStorageSchemes.write.contains("oss"), - s"oss must not be in the native write schemes ${IcebergStorageSchemes.write.toSeq.sorted}") + write.contains("memory"), + s"memory must be in the native write schemes ${write.toSeq.sorted}") assert( - IcebergStorageSchemes.write.contains("memory"), - s"memory must be in the native write schemes ${IcebergStorageSchemes.write.toSeq.sorted}") - assert( - !IcebergStorageSchemes.read.contains("memory"), - s"memory must not be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") + !read.contains("memory"), + s"memory must not be in the native read schemes ${read.toSeq.sorted}") Seq("abfs", "abfss", "wasb", "wasbs", "gcs").foreach { scheme => assert( - !IcebergStorageSchemes.read.contains(scheme), - s"$scheme must not be in the native read schemes ${IcebergStorageSchemes.read.toSeq.sorted}") + !read.contains(scheme), + s"$scheme must not be in the native read schemes ${read.toSeq.sorted}") assert( - !IcebergStorageSchemes.write.contains(scheme), - s"$scheme must not be in the native write schemes " + - s"${IcebergStorageSchemes.write.toSeq.sorted}") + !write.contains(scheme), + s"$scheme must not be in the native write schemes ${write.toSeq.sorted}") } } @@ -219,6 +209,27 @@ class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBa } } + test("iceberg gate: an opt-in alias is matched verbatim, like the built-in schemes") { + // Native opens an alias location as written and OpenDAL's S3 backend checks it against a + // lowercase `scheme://bucket/` prefix, so a `BLOB://` location admitted here would only fail + // at execution. The alias set is lowercase, so a scheme written any other way is declined. + val schemes = Set("blob") + assert( + CometScanRule.isIcebergReadableScheme(new URI("blob://bucket/key.parquet"), schemes), + "blob://bucket/... must be admitted once blob is opted in") + Seq("BLOB://bucket/key.parquet", "Blob://bucket/key.parquet", "BLOB:///bucket/key.parquet") + .foreach { u => + assert( + !CometScanRule.isIcebergReadableScheme(new URI(u), schemes), + s"$u must be declined: native opens the location as written and the S3 backend " + + "rejects a scheme prefix that is not lowercase") + } + // The fallback reason names the lowercase alias, which is the form the scan would admit. + assert( + CometScanRule.icebergSupportedSchemesMessage(schemes).contains("blob"), + "the supported-schemes message must list the opted-in alias in its admitted form") + } + test("parquet gate: mixed-bucket alias scan is declined (single object store per partition)") { // Native planning registers one object store per FilePartition and strips the authority from // every file's object key, so files in a second bucket would be read from the first. A scan