This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new a332ace6 feat(io): let FileIO wrap a caller-provided filesystem 
operator (#618)
a332ace6 is described below

commit a332ace6eed1d4ba35da72230a911b664a245f35
Author: Jordan Epstein <[email protected]>
AuthorDate: Wed Jul 29 20:11:18 2026 -0500

    feat(io): let FileIO wrap a caller-provided filesystem operator (#618)
---
 crates/paimon/src/io/file_io.rs | 22 +++++++++++++--
 crates/paimon/src/io/storage.rs | 62 +++++++++++++++++++++++++++++++++++++++--
 2 files changed, 80 insertions(+), 4 deletions(-)

diff --git a/crates/paimon/src/io/file_io.rs b/crates/paimon/src/io/file_io.rs
index ee5b5711..55113721 100644
--- a/crates/paimon/src/io/file_io.rs
+++ b/crates/paimon/src/io/file_io.rs
@@ -372,6 +372,7 @@ pub struct FileIOBuilder {
     scheme_str: Option<String>,
     props: HashMap<String, String>,
     cache: Option<Arc<LocalCache>>,
+    operator: Option<Operator>,
 }
 
 impl FileIOBuilder {
@@ -380,11 +381,28 @@ impl FileIOBuilder {
             scheme_str: Some(scheme_str.to_string()),
             props: HashMap::default(),
             cache: None,
+            operator: None,
         }
     }
 
-    pub(crate) fn into_parts(self) -> (String, HashMap<String, String>) {
-        (self.scheme_str.unwrap_or_default(), self.props)
+    pub(crate) fn into_parts(self) -> (String, HashMap<String, String>, 
Option<Operator>) {
+        (
+            self.scheme_str.unwrap_or_default(),
+            self.props,
+            self.operator,
+        )
+    }
+
+    /// Uses a caller-provided opendal operator as a **filesystem** backend 
instead of building
+    /// one from the scheme: embedders bring a customized local-filesystem 
service without
+    /// registering a scheme. Paths are resolved with the local-filesystem 
rules — absolute
+    /// paths, `file:` URLs, and Windows drive paths — and handed to the 
operator in relative
+    /// form, so the operator's root decides what they resolve against. 
Scheme'd paths
+    /// (`s3://…`, `oss://…`) are rejected rather than misresolved: an 
object-store operator
+    /// needs bucket/scheme resolution this hook deliberately does not provide.
+    pub fn with_fs_operator(mut self, operator: Operator) -> Self {
+        self.operator = Some(operator);
+        self
     }
 
     pub fn with_prop(mut self, key: impl ToString, value: impl ToString) -> 
Self {
diff --git a/crates/paimon/src/io/storage.rs b/crates/paimon/src/io/storage.rs
index fd696267..eb37b2af 100644
--- a/crates/paimon/src/io/storage.rs
+++ b/crates/paimon/src/io/storage.rs
@@ -68,6 +68,9 @@ use super::FileIOBuilder;
 /// The storage carries all supported storage services in paimon
 #[derive(Debug)]
 pub enum Storage {
+    /// A caller-provided opendal operator, explicitly scoped to filesystem 
semantics
+    /// (see `FileIOBuilder::with_fs_operator`).
+    CustomFs { op: Operator },
     #[cfg(feature = "storage-memory")]
     Memory { op: Operator },
     #[cfg(feature = "storage-fs")]
@@ -111,7 +114,10 @@ pub enum Storage {
 
 impl Storage {
     pub(crate) fn build(file_io_builder: FileIOBuilder) -> crate::Result<Self> 
{
-        let (scheme_str, props) = file_io_builder.into_parts();
+        let (scheme_str, props, operator) = file_io_builder.into_parts();
+        if let Some(op) = operator {
+            return Ok(Self::CustomFs { op });
+        }
         let scheme = scheme_str.to_ascii_lowercase();
         match scheme.as_str() {
             #[cfg(feature = "storage-memory")]
@@ -186,6 +192,22 @@ impl Storage {
 
     pub(crate) fn create<'a>(&self, path: &'a str) -> crate::Result<(Operator, 
Cow<'a, str>)> {
         match self {
+            Storage::CustomFs { op } => {
+                // The custom operator is a *filesystem* operator: paths 
resolve with the
+                // local-filesystem rules below. A scheme'd path would be 
silently mangled by
+                // those rules (`s3://bucket/a` → `3://bucket/a`), so reject 
it instead —
+                // object-store operators need the bucket/scheme resolution 
the built-in
+                // backends perform, which this hook deliberately does not 
reimplement.
+                if !path.starts_with("file:/") && path.contains("://") {
+                    return Err(error::Error::ConfigInvalid {
+                        message: format!(
+                            "the custom operator is a filesystem operator and 
cannot resolve \
+                             the scheme'd path: {path}"
+                        ),
+                    });
+                }
+                Ok((op.clone(), Self::fs_relative_path(path)?))
+            }
             #[cfg(feature = "storage-memory")]
             Storage::Memory { op } => {
                 Ok((op.clone(), 
Cow::Borrowed(Self::memory_relative_path(path)?)))
@@ -289,7 +311,6 @@ impl Storage {
     /// does `PathBuf::from("/").join("C:/dir")`, and because the argument
     /// carries a drive prefix `Path::join` replaces the base, yielding the 
real
     /// `C:\dir` on Windows.
-    #[cfg(feature = "storage-fs")]
     fn fs_relative_path(path: &str) -> crate::Result<Cow<'_, str>> {
         // A `file://` / `file:/` URL is already in scheme-relative form.
         if let Some(stripped) = path.strip_prefix("file:/") {
@@ -557,3 +578,40 @@ mod fs_relative_path_tests {
         assert_eq!(rel(r"C:\Users\wh/db.db/t"), "C:/Users/wh/db.db/t");
     }
 }
+
+#[cfg(all(test, feature = "storage-memory"))]
+mod tests {
+    use super::*;
+
+    fn memory_operator() -> Operator {
+        
Operator::from_config(opendal::services::MemoryConfig::default()).unwrap()
+    }
+
+    /// The filesystem operator resolves absolute paths and `file:` URLs with 
the
+    /// local-filesystem rules, against the operator's own root.
+    #[test]
+    fn custom_fs_operator_resolves_filesystem_paths() {
+        let storage = Storage::CustomFs {
+            op: memory_operator(),
+        };
+        let (_, relative) = storage.create("/warehouse/db/table").unwrap();
+        assert_eq!(relative, "warehouse/db/table");
+        let (_, relative) = 
storage.create("file:/warehouse/db/table").unwrap();
+        assert_eq!(relative, "warehouse/db/table");
+    }
+
+    /// A scheme'd path would be silently mangled by the filesystem rules
+    /// (`s3://bucket/a` → `3://bucket/a`), so it is rejected instead.
+    #[test]
+    fn custom_fs_operator_rejects_schemed_paths() {
+        let storage = Storage::CustomFs {
+            op: memory_operator(),
+        };
+        for path in ["s3://bucket/a", "oss://bucket/a", "hdfs://nn/a"] {
+            assert!(
+                storage.create(path).is_err(),
+                "expected {path} to be rejected by the filesystem operator"
+            );
+        }
+    }
+}

Reply via email to