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 b83e0f26 fix(dlf): refresh expiring credentials (#714)
b83e0f26 is described below

commit b83e0f26f2a704dad81a9fea4b78bd120b1370cb
Author: XiaoHongbo <[email protected]>
AuthorDate: Sat Aug 15 16:15:01 2026 +0800

    fix(dlf): refresh expiring credentials (#714)
---
 crates/paimon/src/api/auth/dlf_provider.rs         |  82 +++++-
 .../paimon/src/catalog/rest/rest_token_file_io.rs  | 280 ++++++++++++---------
 crates/paimon/src/io/file_io.rs                    | 156 ++++++++----
 crates/paimon/src/table/rest_env.rs                |   5 +-
 crates/paimon/tests/rest_api_test.rs               |  16 ++
 docs/src/getting-started.md                        |  23 ++
 6 files changed, 396 insertions(+), 166 deletions(-)

diff --git a/crates/paimon/src/api/auth/dlf_provider.rs 
b/crates/paimon/src/api/auth/dlf_provider.rs
index a18a0143..7dd09ae2 100644
--- a/crates/paimon/src/api/auth/dlf_provider.rs
+++ b/crates/paimon/src/api/auth/dlf_provider.rs
@@ -167,10 +167,24 @@ impl DLFECSTokenLoader {
     /// Get the token from ECS metadata service.
     async fn get_token(&self, url: &str) -> Result<DLFToken> {
         let token_json = self.http_client.get(url).await?;
-        serde_json::from_str(&token_json).map_err(|e| Error::DataInvalid {
-            message: format!("Failed to parse token JSON: {e}"),
-            source: None,
-        })
+        let mut token: DLFToken =
+            serde_json::from_str(&token_json).map_err(|e| Error::DataInvalid {
+                message: format!("Failed to parse token JSON: {e}"),
+                source: None,
+            })?;
+        if token.expiration_at_millis.is_none() {
+            if let Some(expiration) = token.expiration.as_deref() {
+                token.expiration_at_millis = Some(
+                    
DLFToken::parse_expiration_to_millis(expiration).ok_or_else(|| {
+                        Error::DataInvalid {
+                            message: format!("Failed to parse token 
Expiration: {expiration}"),
+                            source: None,
+                        }
+                    })?,
+                );
+            }
+        }
+        Ok(token)
     }
 
     /// Build the token URL from base URL and role name.
@@ -473,4 +487,64 @@ mod tests {
         let millis = DLFToken::parse_expiration_to_millis(expiration);
         assert!(millis.is_some());
     }
+
+    struct RotatingTokenLoader {
+        requests: std::sync::atomic::AtomicUsize,
+    }
+
+    #[async_trait]
+    impl DLFTokenLoader for RotatingTokenLoader {
+        async fn load_token(&self) -> Result<DLFToken> {
+            use std::sync::atomic::Ordering;
+
+            let request = self.requests.fetch_add(1, Ordering::SeqCst);
+            let lifetime = if request == 0 {
+                TOKEN_EXPIRATION_SAFE_TIME_MILLIS / 2
+            } else {
+                TOKEN_EXPIRATION_SAFE_TIME_MILLIS * 2
+            };
+            Ok(DLFToken::new(
+                format!("key-{request}"),
+                "secret",
+                None,
+                Some(Utc::now().timestamp_millis() + lifetime),
+                None,
+            ))
+        }
+
+        fn description(&self) -> &str {
+            "test"
+        }
+    }
+
+    #[tokio::test]
+    async fn test_refreshes_expiring_loaded_token() {
+        use std::sync::atomic::Ordering;
+
+        let loader = Arc::new(RotatingTokenLoader {
+            requests: std::sync::atomic::AtomicUsize::new(0),
+        });
+        let provider = DLFAuthProvider::new(
+            "https://dlf.cn-hangzhou.aliyuncs.com";,
+            "cn-hangzhou",
+            "default",
+            None,
+            Some(loader.clone()),
+        )
+        .unwrap();
+
+        assert_eq!(
+            provider.get_or_refresh_token().await.unwrap().access_key_id,
+            "key-0"
+        );
+        assert_eq!(
+            provider.get_or_refresh_token().await.unwrap().access_key_id,
+            "key-1"
+        );
+        assert_eq!(
+            provider.get_or_refresh_token().await.unwrap().access_key_id,
+            "key-1"
+        );
+        assert_eq!(loader.requests.load(Ordering::SeqCst), 2);
+    }
 }
diff --git a/crates/paimon/src/catalog/rest/rest_token_file_io.rs 
b/crates/paimon/src/catalog/rest/rest_token_file_io.rs
index 8e2accfc..ee4212d2 100644
--- a/crates/paimon/src/catalog/rest/rest_token_file_io.rs
+++ b/crates/paimon/src/catalog/rest/rest_token_file_io.rs
@@ -24,14 +24,14 @@
 use std::collections::HashMap;
 use std::sync::Arc;
 
-use tokio::sync::{OnceCell, RwLock};
+use tokio::sync::{Mutex, RwLock};
 
 use crate::api::rest_api::RESTApi;
 use crate::api::rest_util::RESTUtil;
 use crate::catalog::Identifier;
 use crate::common::{CatalogOptions, Options};
 use crate::io::cache::LocalCache;
-use crate::io::FileIO;
+use crate::io::{FileIO, FileIOProvider};
 use crate::Result;
 
 use super::rest_token::RESTToken;
@@ -40,135 +40,95 @@ use super::rest_token::RESTToken;
 const TOKEN_EXPIRATION_SAFE_TIME_MILLIS: i64 = 3_600_000;
 const OSS_ENDPOINT: &str = "fs.oss.endpoint";
 
-/// A FileIO wrapper that supports getting data access tokens from a REST 
Server.
-///
-/// This struct handles:
-/// - Token caching with expiration detection
-/// - Automatic token refresh via `RESTApi::load_table_token`
-/// - Merging token credentials into catalog options to build the underlying 
`FileIO`
+/// A FileIO wrapper that refreshes data access tokens from the REST server.
+#[derive(Debug)]
+struct TokenState {
+    token: RESTToken,
+    file_io: FileIO,
+}
+
 pub struct RESTTokenFileIO {
-    /// Table identifier for token requests.
     identifier: Identifier,
-    /// Table path (e.g. "oss://bucket/warehouse/db.db/table").
     path: String,
-    /// Catalog options used to build FileIO and create RESTApi.
     catalog_options: Options,
-    /// Lazily-initialized REST API client for token refresh.
-    /// Created on first token refresh and reused for subsequent refreshes.
-    api: OnceCell<RESTApi>,
-    /// Cached token with RwLock for concurrent access.
-    token: RwLock<Option<RESTToken>>,
-    /// Catalog-scoped cache preserved across token-driven FileIO rebuilds.
+    api: Arc<RESTApi>,
+    state: RwLock<Option<TokenState>>,
+    refresh_lock: Mutex<()>,
     local_cache: Option<Arc<LocalCache>>,
 }
 
+impl std::fmt::Debug for RESTTokenFileIO {
+    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+        f.debug_struct("RESTTokenFileIO")
+            .field("identifier", &self.identifier)
+            .field("path", &self.path)
+            .finish_non_exhaustive()
+    }
+}
+
 impl RESTTokenFileIO {
-    /// Create a new RESTTokenFileIO.
-    ///
-    /// # Arguments
-    /// * `identifier` - Table identifier for token requests.
-    /// * `path` - Table path for FileIO construction.
-    /// * `catalog_options` - Catalog options for RESTApi and FileIO.
-    /// * `local_cache` - Catalog-scoped local cache shared across FileIO 
rebuilds.
     pub(crate) fn new(
         identifier: Identifier,
         path: String,
         catalog_options: Options,
+        api: Arc<RESTApi>,
         local_cache: Option<Arc<LocalCache>>,
     ) -> Self {
         Self {
             identifier,
             path,
             catalog_options,
-            api: OnceCell::new(),
-            token: RwLock::new(None),
+            api,
+            state: RwLock::new(None),
+            refresh_lock: Mutex::new(()),
             local_cache,
         }
     }
 
-    /// Build a `FileIO` instance with the current token merged into options.
-    ///
-    /// This method:
-    /// 1. Refreshes the token if expired or not yet obtained.
-    /// 2. Merges token credentials into catalog options.
-    /// 3. Builds a `FileIO` from the merged options.
-    ///
-    /// This method builds a FileIO with the current token,
-    /// which can be passed to `Table::new`. If the token expires, a new
-    /// `get_table` call is needed.
-    pub async fn build_file_io(&self) -> Result<FileIO> {
-        // Ensure token is fresh
-        self.try_to_refresh_token().await?;
-
-        let token_guard = self.token.read().await;
-        match token_guard.as_ref() {
-            Some(token) => {
-                // Merge catalog options (base) with token credentials 
(override)
-                let merged_props =
-                    RESTUtil::merge(Some(self.catalog_options.to_map()), 
Some(&token.token));
-                // Build FileIO with merged properties
-                let mut builder = FileIO::from_path(&self.path)?;
-                builder = builder.with_props(merged_props);
-                if let Some(local_cache) = &self.local_cache {
-                    builder = builder.with_local_cache(local_cache.clone());
-                }
-                builder.build()
-            }
-            None => {
-                // No token available, build FileIO from path only
-                let mut builder = FileIO::from_path(&self.path)?;
-                if let Some(local_cache) = &self.local_cache {
-                    builder = builder.with_local_cache(local_cache.clone());
-                }
-                builder.build()
-            }
-        }
+    pub(crate) async fn build_file_io(self: &Arc<Self>) -> Result<FileIO> {
+        let file_io = self.current_file_io().await?;
+        Ok(file_io.with_provider(self.clone()))
     }
 
-    /// Try to refresh the token if it is expired or not yet obtained.
-    async fn try_to_refresh_token(&self) -> Result<()> {
-        // Fast path: check if token is still valid under read lock
-        {
-            let token_guard = self.token.read().await;
-            if let Some(token) = token_guard.as_ref() {
-                if !Self::is_token_expired(token) {
-                    return Ok(());
-                }
-            }
+    async fn current_file_io(&self) -> Result<FileIO> {
+        if let Some(file_io) = self.valid_file_io().await {
+            return Ok(file_io);
         }
 
-        // Slow path: acquire write lock and check again
-        {
-            let token_guard = self.token.write().await;
-            if let Some(token) = token_guard.as_ref() {
-                if !Self::is_token_expired(token) {
-                    return Ok(());
-                }
-            }
+        let _refresh_guard = self.refresh_lock.lock().await;
+        if let Some(file_io) = self.valid_file_io().await {
+            return Ok(file_io);
         }
-        // Write lock released before .await to avoid potential deadlock
 
-        // Refresh the token WITHOUT holding the lock
-        let new_token = self.refresh_token().await?;
-
-        // Acquire write lock again to update
-        let mut token_guard = self.token.write().await;
-        *token_guard = Some(new_token);
-        Ok(())
+        let token = self.refresh_token().await?;
+        let file_io = self.build_static_file_io(&token)?;
+        *self.state.write().await = Some(TokenState {
+            token,
+            file_io: file_io.clone(),
+        });
+        Ok(file_io)
     }
 
-    /// Refresh the token by calling `RESTApi::load_table_token`.
-    ///
-    /// Lazily creates a `RESTApi` instance on first call and reuses it
-    /// for subsequent refreshes.
-    async fn refresh_token(&self) -> Result<RESTToken> {
-        let api = self
-            .api
-            .get_or_try_init(|| async { 
RESTApi::new(self.catalog_options.clone(), false).await })
-            .await?;
+    async fn valid_file_io(&self) -> Option<FileIO> {
+        self.state
+            .read()
+            .await
+            .as_ref()
+            .filter(|state| !Self::is_token_expired(&state.token))
+            .map(|state| state.file_io.clone())
+    }
 
-        let response = api.load_table_token(&self.identifier).await?;
+    fn build_static_file_io(&self, token: &RESTToken) -> Result<FileIO> {
+        let merged_props = 
RESTUtil::merge(Some(self.catalog_options.to_map()), Some(&token.token));
+        let mut builder = 
FileIO::from_path(&self.path)?.with_props(merged_props);
+        if let Some(local_cache) = &self.local_cache {
+            builder = builder.with_local_cache(local_cache.clone());
+        }
+        builder.build()
+    }
 
+    async fn refresh_token(&self) -> Result<RESTToken> {
+        let response = self.api.load_table_token(&self.identifier).await?;
         let expires_at_millis =
             response
                 .expires_at_millis
@@ -180,27 +140,23 @@ impl RESTTokenFileIO {
                     source: None,
                 })?;
 
-        // Merge token with catalog options (e.g. DLF OSS endpoint override)
         let merged_token = 
self.merge_token_with_catalog_options(response.token);
         Ok(RESTToken::new(merged_token, expires_at_millis))
     }
 
-    /// Check if a token is expired (within the safe time margin).
     fn is_token_expired(token: &RESTToken) -> bool {
         let current_time = std::time::SystemTime::now()
             .duration_since(std::time::UNIX_EPOCH)
             .unwrap_or_default()
             .as_millis() as i64;
-        (token.expire_at_millis - current_time) < 
TOKEN_EXPIRATION_SAFE_TIME_MILLIS
+        token.expire_at_millis - current_time < 
TOKEN_EXPIRATION_SAFE_TIME_MILLIS
     }
 
-    /// Merge token credentials with catalog options for DLF OSS endpoint 
override.
     fn merge_token_with_catalog_options(
         &self,
         token: HashMap<String, String>,
     ) -> HashMap<String, String> {
         let mut merged = token;
-        // If catalog options contain a DLF OSS endpoint, override the 
standard OSS endpoint
         if let Some(dlf_oss_endpoint) = 
self.catalog_options.get(CatalogOptions::DLF_OSS_ENDPOINT) {
             if !dlf_oss_endpoint.trim().is_empty() {
                 merged.insert(OSS_ENDPOINT.to_string(), 
dlf_oss_endpoint.clone());
@@ -210,37 +166,135 @@ impl RESTTokenFileIO {
     }
 }
 
+#[async_trait::async_trait]
+impl FileIOProvider for RESTTokenFileIO {
+    async fn create(&self, path: &str) -> Result<(opendal::Operator, String)> {
+        self.current_file_io().await?.create_static(path)
+    }
+}
+
 #[cfg(test)]
 mod tests {
+    use std::sync::atomic::{AtomicUsize, Ordering};
+
+    use axum::extract::State;
+    use axum::routing::get;
+    use axum::{Json, Router};
+    use bytes::Bytes;
+
     use super::*;
+    use crate::api::GetTableTokenResponse;
     use crate::io::cache::create_local_cache;
 
+    async fn token(State(requests): State<Arc<AtomicUsize>>) -> 
Json<GetTableTokenResponse> {
+        let request = requests.fetch_add(1, Ordering::SeqCst);
+        let now = std::time::SystemTime::now()
+            .duration_since(std::time::UNIX_EPOCH)
+            .unwrap()
+            .as_millis() as i64;
+        let lifetime = if request == 0 {
+            TOKEN_EXPIRATION_SAFE_TIME_MILLIS / 2
+        } else {
+            TOKEN_EXPIRATION_SAFE_TIME_MILLIS * 2
+        };
+        Json(GetTableTokenResponse {
+            token: HashMap::new(),
+            expires_at_millis: Some(now + lifetime),
+        })
+    }
+
+    async fn token_api() -> (
+        Options,
+        Arc<RESTApi>,
+        Arc<AtomicUsize>,
+        tokio::task::JoinHandle<()>,
+    ) {
+        let requests = Arc::new(AtomicUsize::new(0));
+        let app = Router::new()
+            .route("/v1/databases/database/tables/table/token", get(token))
+            .with_state(requests.clone());
+        let listener = 
tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
+        let address = listener.local_addr().unwrap();
+        let server = tokio::spawn(async move {
+            axum::serve(listener, app).await.unwrap();
+        });
+
+        let mut options = Options::new();
+        options.set(CatalogOptions::URI, format!("http://{address}";));
+        options.set(CatalogOptions::TOKEN_PROVIDER, "bear");
+        options.set(CatalogOptions::TOKEN, "test-token");
+        let api = Arc::new(RESTApi::new(options.clone(), 
false).await.unwrap());
+        (options, api, requests, server)
+    }
+
     #[tokio::test]
     async fn test_token_file_io_keeps_catalog_local_cache() {
         let cache_directory = tempfile::tempdir().unwrap();
         let table_directory = tempfile::tempdir().unwrap();
-        let mut options = Options::new();
+        let (mut options, api, _, server) = token_api().await;
         options.set(CatalogOptions::LOCAL_CACHE_ENABLED, "true");
         options.set(
             CatalogOptions::LOCAL_CACHE_DIR,
             cache_directory.path().to_string_lossy(),
         );
         let local_cache = create_local_cache(&options).unwrap();
-        let token_file_io = RESTTokenFileIO::new(
+        let token_file_io = Arc::new(RESTTokenFileIO::new(
             Identifier::new("database", "table"),
             table_directory.path().to_string_lossy().into_owned(),
             options,
+            api,
             local_cache,
-        );
-        let valid_until = std::time::SystemTime::now()
-            .duration_since(std::time::UNIX_EPOCH)
-            .unwrap()
-            .as_millis() as i64
-            + TOKEN_EXPIRATION_SAFE_TIME_MILLIS * 2;
-        *token_file_io.token.write().await = 
Some(RESTToken::new(HashMap::new(), valid_until));
+        ));
 
         let file_io = token_file_io.build_file_io().await.unwrap();
 
         assert!(file_io.has_local_cache());
+        server.abort();
+    }
+
+    #[tokio::test]
+    async fn test_file_io_refreshes_expiring_token() {
+        let table_directory = tempfile::tempdir().unwrap();
+        let file_path = table_directory.path().join("data");
+        let (options, api, requests, server) = token_api().await;
+        let token_file_io = Arc::new(RESTTokenFileIO::new(
+            Identifier::new("database", "table"),
+            table_directory.path().to_string_lossy().into_owned(),
+            options,
+            api,
+            None,
+        ));
+
+        let file_io = token_file_io.build_file_io().await.unwrap();
+        assert_eq!(requests.load(Ordering::SeqCst), 1);
+        let file_io = Arc::new(file_io);
+        let mut checks = Vec::new();
+        for _ in 0..8 {
+            let file_io = file_io.clone();
+            let path = file_path.to_string_lossy().into_owned();
+            checks.push(tokio::spawn(async move { file_io.exists(&path).await 
}));
+        }
+        for check in checks {
+            assert!(!check.await.unwrap().unwrap());
+        }
+        assert_eq!(requests.load(Ordering::SeqCst), 2);
+
+        file_io
+            .new_output(file_path.to_string_lossy().as_ref())
+            .unwrap()
+            .write(Bytes::from_static(b"data"))
+            .await
+            .unwrap();
+        assert_eq!(requests.load(Ordering::SeqCst), 2);
+
+        let bytes = file_io
+            .new_input(file_path.to_string_lossy().as_ref())
+            .unwrap()
+            .read()
+            .await
+            .unwrap();
+        assert_eq!(bytes, Bytes::from_static(b"data"));
+        assert_eq!(requests.load(Ordering::SeqCst), 2);
+        server.abort();
     }
 }
diff --git a/crates/paimon/src/io/file_io.rs b/crates/paimon/src/io/file_io.rs
index 55113721..1a4998a0 100644
--- a/crates/paimon/src/io/file_io.rs
+++ b/crates/paimon/src/io/file_io.rs
@@ -35,13 +35,46 @@ use url::Url;
 use super::cache::{CachedFileReader, LocalCache};
 use super::Storage;
 
-#[derive(Clone, Debug)]
+#[async_trait::async_trait]
+pub(crate) trait FileIOProvider: std::fmt::Debug + Send + Sync {
+    async fn create(&self, path: &str) -> crate::Result<(Operator, String)>;
+}
+
+#[derive(Clone)]
 pub struct FileIO {
     storage: Arc<Storage>,
     cache: Option<Arc<LocalCache>>,
+    provider: Option<Arc<dyn FileIOProvider>>,
+}
+
+impl std::fmt::Debug for FileIO {
+    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+        f.debug_struct("FileIO")
+            .field("storage", &self.storage)
+            .field("cache", &self.cache)
+            .field("provider", &self.provider)
+            .finish()
+    }
 }
 
 impl FileIO {
+    pub(crate) fn with_provider(mut self, provider: Arc<dyn FileIOProvider>) 
-> Self {
+        self.provider = Some(provider);
+        self
+    }
+
+    pub(crate) fn create_static(&self, path: &str) -> crate::Result<(Operator, 
String)> {
+        let (op, relative_path) = self.storage.create(path)?;
+        Ok((op, relative_path.into_owned()))
+    }
+
+    async fn create(&self, path: &str) -> crate::Result<(Operator, String)> {
+        match &self.provider {
+            Some(provider) => provider.create(path).await,
+            None => self.create_static(path),
+        }
+    }
+
     #[cfg(test)]
     pub(crate) fn has_local_cache(&self) -> bool {
         self.cache.is_some()
@@ -100,6 +133,7 @@ impl FileIO {
                 .as_ref()
                 .filter(|cache| cache.is_cacheable(path))
                 .cloned(),
+            provider: self.provider.clone(),
         })
     }
 
@@ -119,6 +153,7 @@ impl FileIO {
                 .as_ref()
                 .filter(|cache| cache.is_cacheable(path))
                 .cloned(),
+            provider: self.provider.clone(),
         })
     }
 
@@ -126,7 +161,7 @@ impl FileIO {
     ///
     /// Reference: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L97>
     pub async fn get_status(&self, path: &str) -> Result<FileStatus> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         let meta = op
             .stat(relative_path.as_ref())
             .await
@@ -150,7 +185,7 @@ impl FileIO {
     ///
     /// FIXME: how to handle large dir? Better to return a stream instead?
     pub async fn list_status(&self, path: &str) -> Result<Vec<FileStatus>> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         // `relative_path` is a byte-suffix of `path` for object stores and 
POSIX
         // local paths, so this recovers the scheme/root prefix. For a Windows
         // local path the relative form only swaps `\`->`/` 
(length-preserving),
@@ -188,7 +223,7 @@ impl FileIO {
 
     /// List all files recursively under the given directory path.
     pub async fn list_status_recursive(&self, path: &str) -> 
Result<Vec<FileStatus>> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         // See `list_status`: `relative_path` is a byte-suffix of `path` except
         // for Windows local paths, where it only swaps separators (same 
length).
         let base_path = &path[..path.len() - relative_path.len()];
@@ -230,7 +265,7 @@ impl FileIO {
     ///
     /// References: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L128>
     pub async fn exists(&self, path: &str) -> Result<bool> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
 
         op.exists(relative_path.as_ref())
             .await
@@ -241,7 +276,7 @@ impl FileIO {
 
     /// Check if a directory exists.
     pub async fn exists_dir(&self, path: &str) -> Result<bool> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         let dir_path = normalize_root(relative_path.as_ref());
 
         op.exists(&dir_path).await.context(IoUnexpectedSnafu {
@@ -253,7 +288,7 @@ impl FileIO {
     ///
     /// Reference: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L139>
     pub async fn delete_file(&self, path: &str) -> Result<()> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         let cache_path = cache_object_path(&op, relative_path.as_ref());
 
         op.delete(relative_path.as_ref())
@@ -272,7 +307,7 @@ impl FileIO {
     ///
     /// Reference: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L139>
     pub async fn delete_dir(&self, path: &str) -> Result<()> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         let cache_path = cache_object_path(&op, relative_path.as_ref());
 
         op.delete_with(relative_path.as_ref())
@@ -294,7 +329,7 @@ impl FileIO {
     ///
     /// Reference: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L150>
     pub async fn mkdirs(&self, path: &str) -> Result<()> {
-        let (op, relative_path) = self.storage.create(path)?;
+        let (op, relative_path) = self.create(path).await?;
         // Opendal create_dir expects the path to end with `/` to indicate a 
directory.
         let dir_path = normalize_root(relative_path.as_ref());
         op.create_dir(&dir_path).await.context(IoUnexpectedSnafu {
@@ -319,8 +354,8 @@ impl FileIO {
     ///
     /// Reference: 
<https://github.com/apache/paimon/blob/release-0.8.2/paimon-common/src/main/java/org/apache/paimon/fs/FileIO.java#L159>
     pub async fn rename(&self, src: &str, dst: &str) -> Result<()> {
-        let (op_src, relative_path_src) = self.storage.create(src)?;
-        let (op_dst, relative_path_dst) = self.storage.create(dst)?;
+        let (op_src, relative_path_src) = self.create(src).await?;
+        let (op_dst, relative_path_dst) = self.create(dst).await?;
         let cache_path_src = cache_object_path(&op_src, 
relative_path_src.as_ref());
         let cache_path_dst = cache_object_path(&op_dst, 
relative_path_dst.as_ref());
 
@@ -430,6 +465,7 @@ impl FileIOBuilder {
         Ok(FileIO {
             storage: Arc::new(storage),
             cache,
+            provider: None,
         })
     }
 }
@@ -573,19 +609,37 @@ pub struct InputFile {
     relative_path: String,
     cache_path: String,
     cache: Option<Arc<LocalCache>>,
+    provider: Option<Arc<dyn FileIOProvider>>,
 }
 
 impl InputFile {
+    async fn resolve(&self) -> crate::Result<(Operator, String, String)> {
+        match &self.provider {
+            Some(provider) => {
+                let (op, relative_path) = provider.create(&self.path).await?;
+                let cache_path = cache_object_path(&op, &relative_path);
+                Ok((op, relative_path, cache_path))
+            }
+            None => Ok((
+                self.op.clone(),
+                self.relative_path.clone(),
+                self.cache_path.clone(),
+            )),
+        }
+    }
+
     pub fn location(&self) -> &str {
         &self.path
     }
 
     pub async fn exists(&self) -> crate::Result<bool> {
-        Ok(self.op.exists(&self.relative_path).await?)
+        let (op, relative_path, _) = self.resolve().await?;
+        Ok(op.exists(&relative_path).await?)
     }
 
     pub async fn metadata(&self) -> crate::Result<FileStatus> {
-        let meta = self.op.stat(&self.relative_path).await?;
+        let (op, relative_path, _) = self.resolve().await?;
+        let meta = op.stat(&relative_path).await?;
 
         Ok(FileStatus {
             size: meta.content_length(),
@@ -598,49 +652,41 @@ impl InputFile {
     }
 
     pub async fn read(&self) -> crate::Result<Bytes> {
+        let (op, relative_path, cache_path) = self.resolve().await?;
         let Some(cache) = &self.cache else {
-            return Ok(self.op.read(&self.relative_path).await?.to_bytes());
+            return Ok(op.read(&relative_path).await?.to_bytes());
         };
-        let read_token = cache.read_token(&self.cache_path);
-        let size = if let Some(size) = cache.file_size(&self.cache_path, 
&read_token).await {
+        let read_token = cache.read_token(&cache_path);
+        let size = if let Some(size) = cache.file_size(&cache_path, 
&read_token).await {
             size
         } else {
-            let size = 
self.op.stat(&self.relative_path).await?.content_length();
-            cache
-                .put_file_size(&self.cache_path, size, &read_token)
-                .await;
+            let size = op.stat(&relative_path).await?.content_length();
+            cache.put_file_size(&cache_path, size, &read_token).await;
             size
         };
-        let delegate = Arc::new(self.op.reader(&self.relative_path).await?);
-        CachedFileReader::new_with_token(
-            delegate,
-            &self.cache_path,
-            size,
-            cache.clone(),
-            read_token,
-        )
-        .read_full()
-        .await
+        let delegate = Arc::new(op.reader(&relative_path).await?);
+        CachedFileReader::new_with_token(delegate, &cache_path, size, 
cache.clone(), read_token)
+            .read_full()
+            .await
     }
 
     pub async fn reader(&self) -> crate::Result<impl FileRead> {
-        let reader = self.op.reader(&self.relative_path).await?;
+        let (op, relative_path, cache_path) = self.resolve().await?;
+        let reader = op.reader(&relative_path).await?;
         let Some(cache) = &self.cache else {
             return Ok(InputFileReader::Direct(reader));
         };
-        let read_token = cache.read_token(&self.cache_path);
-        let size = if let Some(size) = cache.file_size(&self.cache_path, 
&read_token).await {
+        let read_token = cache.read_token(&cache_path);
+        let size = if let Some(size) = cache.file_size(&cache_path, 
&read_token).await {
             size
         } else {
-            let size = 
self.op.stat(&self.relative_path).await?.content_length();
-            cache
-                .put_file_size(&self.cache_path, size, &read_token)
-                .await;
+            let size = op.stat(&relative_path).await?.content_length();
+            cache.put_file_size(&cache_path, size, &read_token).await;
             size
         };
         Ok(InputFileReader::Cached(CachedFileReader::new_with_token(
             Arc::new(reader),
-            &self.cache_path,
+            &cache_path,
             size,
             cache.clone(),
             read_token,
@@ -657,15 +703,32 @@ pub struct OutputFile {
     relative_path: String,
     cache_path: String,
     cache: Option<Arc<LocalCache>>,
+    provider: Option<Arc<dyn FileIOProvider>>,
 }
 
 impl OutputFile {
+    async fn resolve(&self) -> crate::Result<(Operator, String, String)> {
+        match &self.provider {
+            Some(provider) => {
+                let (op, relative_path) = provider.create(&self.path).await?;
+                let cache_path = cache_object_path(&op, &relative_path);
+                Ok((op, relative_path, cache_path))
+            }
+            None => Ok((
+                self.op.clone(),
+                self.relative_path.clone(),
+                self.cache_path.clone(),
+            )),
+        }
+    }
+
     pub fn location(&self) -> &str {
         &self.path
     }
 
     pub async fn exists(&self) -> crate::Result<bool> {
-        Ok(self.op.exists(&self.relative_path).await?)
+        let (op, relative_path, _) = self.resolve().await?;
+        Ok(op.exists(&relative_path).await?)
     }
 
     pub fn to_input_file(self) -> InputFile {
@@ -676,6 +739,7 @@ impl OutputFile {
             relative_path: self.relative_path,
             cache_path: self.cache_path,
             cache,
+            provider: self.provider,
         }
     }
 
@@ -686,21 +750,23 @@ impl OutputFile {
     }
 
     pub async fn writer(&self) -> crate::Result<Box<dyn FileWrite>> {
-        let writer: Box<dyn FileWrite> = 
Box::new(self.opendal_writer().await?);
+        let (op, relative_path, cache_path) = self.resolve().await?;
+        let writer: Box<dyn FileWrite> = 
Box::new(op.writer(&relative_path).await?);
         let Some(cache) = &self.cache else {
             return Ok(writer);
         };
         Ok(Box::new(CacheInvalidatingWriter {
             delegate: writer,
             cache: cache.clone(),
-            path: self.cache_path.clone(),
+            path: cache_path,
         }))
     }
 
     /// Get an async streaming writer for format-level writes (e.g. parquet).
     pub(crate) async fn async_writer(&self) -> crate::Result<Box<dyn 
AsyncFileWrite>> {
+        let (op, relative_path, cache_path) = self.resolve().await?;
         let writer: Box<dyn AsyncFileWrite> = Box::new(
-            self.opendal_writer()
+            op.writer(&relative_path)
                 .await?
                 .into_futures_async_write()
                 .compat_write(),
@@ -711,15 +777,11 @@ impl OutputFile {
         Ok(Box::new(CacheInvalidatingAsyncWriter {
             delegate: writer,
             cache: cache.clone(),
-            path: self.cache_path.clone(),
+            path: cache_path,
             delegate_shutdown: false,
             invalidation: None,
         }))
     }
-
-    async fn opendal_writer(&self) -> crate::Result<opendal::Writer> {
-        Ok(self.op.writer(&self.relative_path).await?)
-    }
 }
 
 #[cfg(test)]
diff --git a/crates/paimon/src/table/rest_env.rs 
b/crates/paimon/src/table/rest_env.rs
index 4f2a9df7..b61370b2 100644
--- a/crates/paimon/src/table/rest_env.rs
+++ b/crates/paimon/src/table/rest_env.rs
@@ -153,12 +153,13 @@ impl RESTEnv {
         })?;
 
         let file_io = if data_token_enabled && !is_external {
-            RESTTokenFileIO::new(
+            Arc::new(RESTTokenFileIO::new(
                 identifier.clone(),
                 table_path.clone(),
                 options.clone(),
+                api.clone(),
                 local_cache.clone(),
-            )
+            ))
             .build_file_io()
             .await?
         } else {
diff --git a/crates/paimon/tests/rest_api_test.rs 
b/crates/paimon/tests/rest_api_test.rs
index 99b952ce..c7dff44b 100644
--- a/crates/paimon/tests/rest_api_test.rs
+++ b/crates/paimon/tests/rest_api_test.rs
@@ -673,6 +673,10 @@ async fn test_ecs_loader_token() {
         load_token.expiration,
         Some("2023-12-01T12:00:00Z".to_string())
     );
+    assert_eq!(
+        load_token.expiration_at_millis,
+        DLFToken::parse_expiration_to_millis("2023-12-01T12:00:00Z")
+    );
 
     // Test with role name
     let loader_with_role = DLFECSTokenLoader::new(&ecs_metadata_url, 
Some(role_name.to_string()));
@@ -685,4 +689,16 @@ async fn test_ecs_loader_token() {
         Some("AQoDYXdzEJr...<remainder of security token>".to_string())
     );
     assert_eq!(token.expiration, Some("2023-12-01T12:00:00Z".to_string()));
+
+    let invalid_token_json = json!({
+        "AccessKeyId": "AccessKeyId",
+        "AccessKeySecret": "AccessKeySecret",
+        "SecurityToken": "token",
+        "Expiration": "invalid-expiration"
+    });
+    server.set_ecs_metadata(role_name, invalid_token_json);
+    let error = loader_with_role.load_token().await.unwrap_err();
+    assert!(error
+        .to_string()
+        .contains("Failed to parse token Expiration"));
 }
diff --git a/docs/src/getting-started.md b/docs/src/getting-started.md
index 8d90895e..7a998910 100644
--- a/docs/src/getting-started.md
+++ b/docs/src/getting-started.md
@@ -129,6 +129,29 @@ options.set(CatalogOptions::WAREHOUSE, "my_warehouse");
 let catalog = CatalogFactory::create(options).await?;
 ```
 
+For a DLF REST catalog on ECS, RAM-role credentials can be rotated 
automatically:
+
+```rust
+let mut options = Options::new();
+options.set(CatalogOptions::METASTORE, "rest");
+options.set(CatalogOptions::URI, "https://your-dlf-endpoint";);
+options.set(CatalogOptions::WAREHOUSE, "your_catalog");
+options.set(CatalogOptions::TOKEN_PROVIDER, "dlf");
+options.set(CatalogOptions::DLF_REGION, "cn-hangzhou");
+options.set(CatalogOptions::DLF_TOKEN_LOADER, "ecs");
+options.set(CatalogOptions::DLF_TOKEN_ECS_ROLE_NAME, "your-ram-role");
+options.set(CatalogOptions::DATA_TOKEN_ENABLED, "true");
+let catalog = CatalogFactory::create(options).await?;
+```
+
+`dlf.token-loader=ecs` refreshes the credentials used to authenticate DLF
+catalog requests. The role name is optional; when omitted, it is read from the
+ECS metadata service. `data-token.enabled=true` separately enables temporary
+credentials returned by the REST server for table data access. A loaded table
+refreshes those credentials before expiration, so callers do not need to load
+the table again. Static `dlf.access-key-id`, `dlf.access-key-secret`, and
+`dlf.security-token` values are not rotated.
+
 Supported metastore types:
 
 | Metastore Type | Description                      |

Reply via email to