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

github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new af1da4c5b feat(storage-opendal)!: make the per-IO-operation timeout 
configurable (#3263)
af1da4c5b is described below

commit af1da4c5bc86178c38c1c6db0578bdd9fb7021e3
Author: Oleks V <[email protected]>
AuthorDate: Sun Oct 4 16:06:52 2026 +0000

    feat(storage-opendal)!: make the per-IO-operation timeout configurable 
(#3263)
    
    * feat(storage-opendal)!: make the per-IO-operation timeout configurable
    
    `iceberg-storage-opendal` wraps every FileIO operator in 
`TimeoutLayer::new()`,
    whose 10s `io_timeout` bounds each `read`/`write` and each method call on a
    returned reader or writer. Nothing in the property or builder surface can
    override it, so an operation that legitimately needs longer fails
    deterministically: the `RetryLayer` above re-sends the same request, which
    cannot fit in the budget either.
    
    Add `client.io-timeout-ms`, parsed into an `OpenDalClientConfig` built with
    `#[derive(Properties)]` and handed to `TimeoutLayer::with_io_timeout`. The
    config's fields are private, so later client settings such as retry are
    additive rather than breaking. The default is unchanged.
    
    BREAKING CHANGE: `OpenDalStorage` variants carry a `client` field, and
    `OpenDalStorage::Memory` is now a struct variant.
    
    * refactor(storage-opendal): address review feedback on the IO timeout
    
    Quote the rejected value in the parse error, so an empty input shows as
    `value: ""`. `TimeoutLayer` exposes no getters, so a new test compares
    its `Debug` output to pin `DEFAULT_IO_TIMEOUT_MS` to OpenDAL's default.
    The parsing test now uses the constant instead of repeating the literal.
    
    Add a serde round-trip test for the `client` field, including a payload
    that omits it. Drop the `DEFAULT` associated const, which only fed a
    fallback arm that cannot run. Update the upstream
    `test_writer_close_returns_stored_size` for the struct variant, and reuse
    `StorageConfig::with_prop` and `empty_resolving_storage` in the tests.
    
    * refactor(storage-opendal): rename the IO timeout key to 
opendal.io-timeout-ms
    
    Only `iceberg-storage-opendal` honors the key, so namespace it under
    `opendal.` and define it as `OPENDAL_IO_TIMEOUT_MS` next to
    `OpenDalClientConfig`, rather than in `iceberg::io`. This matches how
    `ADLS_SAS_TOKEN` and `GCS_ALLOW_ANONYMOUS` sit next to the configs that
    parse them. With the key moved, the `iceberg` crate needs no change.
    
    * refactor(storage-opendal): rename the client field to client_config
    
    `client` reads like an HTTP client handle. Rename the field on every
    `OpenDalStorage` variant, and the `OpenDalStorage::client()` accessor,
    to `client_config` to match `OpenDalClientConfig`. `opendal_config`
    would be ambiguous next to the backend `config` field, which holds an
    OpenDAL config too.
    
    `test_default_io_timeout_matches_opendal` compares `Debug` output. Both
    sides go through the same derived impl, so a format change cannot make
    them differ. A `Debug` that stopped printing `io_timeout` would make the
    check pass vacuously, so also assert that a different `io_timeout`
    produces a different string.
    
    * fix(storage-opendal): reject a zero IO timeout when deserializing
    
    `parse_io_timeout_ms` rejected zero, but the derived `Deserialize` did
    not go through it. A serialized `OpenDalStorage` carrying
    `"io_timeout_ms": 0` was rebuilt without error, and the zero reached
    `TimeoutLayer`.
    
    Store the timeout as a `NonZeroU64`. `str::parse` and serde's
    `NonZeroU64` impl both reject zero, so both construction paths share
    one check. The serialized form is still a plain integer. The generated
    getter now returns `&NonZeroU64`, because the `Properties` derive only
    returns by value for the primitive types it recognizes.
    
    The parsing test also covers the edges of the valid range. `1` and
    `u64::MAX` are accepted, and `u64::MAX + 1` is rejected.
    
    * refactor(storage-opendal): return the IO timeout as a Duration
    
    `OpenDalClientConfig::io_timeout_ms()` returned `&NonZeroU64`, because the
    `Properties` derive only returns primitives by value. Replace the generated
    getter with `io_timeout()`, which returns a `Duration`, and keep the
    millisecond field private.
    
    Expose the default as `OPENDAL_IO_TIMEOUT_MS_DEFAULT`, following the
    `*_DEFAULT` consts in `iceberg`, so the key's doc links to it instead of
    repeating the number. Keep the `ParseIntError` as the source of a rejected
    value, as `parse_pool_property` does, so the message says why it failed.
    
    Document that the serialized form of `OpenDalStorage` is not stable across
    crate versions, as `FileIO::serialize_all` already does for its own.
    
    Test that the configured timeout reaches `TimeoutLayer`. A memory operator
    behind `ConcurrentLimitLayer::new(0)` stalls every IO call, and paused tokio
    time skips the wait and the retry backoff. Check propagation through
    `OpenDalResolvingStorage` for every enabled scheme, not only S3.
    
    * refactor(storage-opendal): take the client config defaults from the derive
    
    Build `OpenDalClientConfig::default()` from `from_properties` with no
    properties, so each default is declared once in `#[property]` and serde
    and `from_properties` cannot disagree as settings are added.
    
    Say in the `parse_io_timeout_ms` doc that it exists to quote the value
    and name the unit, and that `NonZeroU64` is what rejects zero. Fold the
    valid-value asserts in the parsing test into one loop.
    
    * refactor(storage-opendal): build the client config default directly
    
    `Default` delegated to `from_properties` with no properties and called
    `expect`. A later setting without a `#[property]` default would make it
    panic, and because of the container-level `#[serde(default)]`, serde
    calls `default()` on every `client_config` it deserializes. Build the
    struct directly, and keep the two defaults equal with
    `test_default_matches_property_defaults`.
    
    Name OpenDAL's 60-second control-operation budget in the
    `OPENDAL_IO_TIMEOUT_MS` doc, along with the calls it covers here,
    `exists` and `metadata`. `test_default_timeouts_match_opendal` now pins
    both defaults that the doc states.
    
    Pin the serialized-form break with a test that rejects the old
    `"LocalFs"`, `"Memory"` and `{"Memory": null}` forms. Match the stall
    test's error on `{ timeout: 45 }`, so a 450-second budget cannot pass.
    
    * refactor(storage-opendal): keep the client config accessors crate-private
    
    `OpenDalStorage::client_config()`, `OpenDalClientConfig::io_timeout()` and
    `OPENDAL_IO_TIMEOUT_MS_DEFAULT` have no callers outside the crate. Keep them
    out of the public API so they can be exposed later without a breaking 
change.
    
    Co-authored-by: Copilot App <[email protected]>
    
    * refactor(storage-opendal): trim redundant docs, comments and test
    
    Fold the `u64` default into the `NonZeroU64` const, drop
    `test_default_matches_property_defaults` (one field, one shared const), 
shorten
    docs and comments that restated the code, and keep the original `op` 
binding in
    the `Memory` arm of `create_operator`.
    
    Co-authored-by: Copilot App <[email protected]>
    
    * docs(storage-opendal): clarify the IO timeout and serialization docs
    
    Name OpenDAL's `io_timeout` and `timeout` settings so it is clear the 
property
    bounds each IO call per retry attempt and leaves control operations such as
    `stat` at 60 seconds. Align the `OpenDalStorage` serialization note with the
    wording on `FileIO::serialize_all`.
    
    Co-authored-by: Copilot App <[email protected]>
    
    * test(storage-opendal): check OpenDalStorageFactory propagates the IO 
timeout
    
    Every `OpenDalStorageFactory::build` arm could fall back to the default
    `OpenDalClientConfig` without any test failing. Build each enabled backend 
with
    a 45000 ms timeout and read it back from the serialized storage.
    
    Co-authored-by: Copilot App <[email protected]>
    
    ---------
    
    Co-authored-by: ovoievodin <[email protected]>
    Co-authored-by: Kevin Liu <[email protected]>
    Co-authored-by: Copilot App <[email protected]>
---
 Cargo.lock                              |   2 +
 crates/storage/opendal/Cargo.toml       |   4 +-
 crates/storage/opendal/public-api.txt   |  24 ++-
 crates/storage/opendal/src/lib.rs       | 330 ++++++++++++++++++++++++++++++--
 crates/storage/opendal/src/resolving.rs |  51 ++++-
 5 files changed, 387 insertions(+), 24 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock
index fcf5dd449..7b2115a71 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -2935,12 +2935,14 @@ dependencies = [
  "cfg-if 1.0.4",
  "futures",
  "iceberg",
+ "iceberg-property-macro",
  "iceberg_test_utils",
  "opendal",
  "reqsign-aws-v4",
  "reqsign-core",
  "reqwest 0.12.28",
  "serde",
+ "serde_json",
  "tempfile",
  "tokio",
  "typetag",
diff --git a/crates/storage/opendal/Cargo.toml 
b/crates/storage/opendal/Cargo.toml
index c6ea9f9b6..d3403598e 100644
--- a/crates/storage/opendal/Cargo.toml
+++ b/crates/storage/opendal/Cargo.toml
@@ -54,6 +54,7 @@ bytes = { workspace = true }
 cfg-if = { workspace = true }
 futures = { workspace = true }
 iceberg = { workspace = true }
+iceberg-property-macro = { workspace = true }
 opendal = { workspace = true }
 reqsign-aws-v4 = { version = "3.0.0", optional = true }
 reqsign-core = { version = "3.0.0", optional = true }
@@ -65,8 +66,9 @@ url = { workspace = true }
 async-trait = { workspace = true }
 iceberg_test_utils = { path = "../../test_utils", features = ["tests"] }
 reqwest = { workspace = true }
+serde_json = { workspace = true }
 tempfile = { workspace = true }
-tokio = { workspace = true, features = ["macros"] }
+tokio = { workspace = true, features = ["macros", "test-util"] }
 
 [lints]
 workspace = true
diff --git a/crates/storage/opendal/public-api.txt 
b/crates/storage/opendal/public-api.txt
index d8c4ecdb3..c2e310335 100644
--- a/crates/storage/opendal/public-api.txt
+++ b/crates/storage/opendal/public-api.txt
@@ -3,16 +3,24 @@ pub use iceberg_storage_opendal::AwsCredential
 pub use iceberg_storage_opendal::ProvideCredential
 pub enum iceberg_storage_opendal::OpenDalStorage
 pub iceberg_storage_opendal::OpenDalStorage::Azdls
+pub iceberg_storage_opendal::OpenDalStorage::Azdls::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
 pub iceberg_storage_opendal::OpenDalStorage::Azdls::config: 
alloc::sync::Arc<opendal_service_azdls::config::AzdlsConfig>
 pub iceberg_storage_opendal::OpenDalStorage::Gcs
+pub iceberg_storage_opendal::OpenDalStorage::Gcs::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
 pub iceberg_storage_opendal::OpenDalStorage::Gcs::config: 
alloc::sync::Arc<opendal_service_gcs::config::GcsConfig>
 pub iceberg_storage_opendal::OpenDalStorage::Hf
+pub iceberg_storage_opendal::OpenDalStorage::Hf::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
 pub iceberg_storage_opendal::OpenDalStorage::Hf::config: 
alloc::sync::Arc<opendal_service_hf::config::HfConfig>
 pub iceberg_storage_opendal::OpenDalStorage::LocalFs
-pub 
iceberg_storage_opendal::OpenDalStorage::Memory(opendal_core::types::operator::operator::Operator)
+pub iceberg_storage_opendal::OpenDalStorage::LocalFs::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
+pub iceberg_storage_opendal::OpenDalStorage::Memory
+pub iceberg_storage_opendal::OpenDalStorage::Memory::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
+pub iceberg_storage_opendal::OpenDalStorage::Memory::operator: 
opendal_core::types::operator::operator::Operator
 pub iceberg_storage_opendal::OpenDalStorage::Oss
+pub iceberg_storage_opendal::OpenDalStorage::Oss::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
 pub iceberg_storage_opendal::OpenDalStorage::Oss::config: 
alloc::sync::Arc<opendal_service_oss::config::OssConfig>
 pub iceberg_storage_opendal::OpenDalStorage::S3
+pub iceberg_storage_opendal::OpenDalStorage::S3::client_config: 
iceberg_storage_opendal::OpenDalClientConfig
 pub iceberg_storage_opendal::OpenDalStorage::S3::config: 
alloc::sync::Arc<opendal_service_s3::config::S3Config>
 pub iceberg_storage_opendal::OpenDalStorage::S3::customized_credential_load: 
core::option::Option<iceberg_storage_opendal::CustomAwsCredentialLoader>
 impl core::clone::Clone for iceberg_storage_opendal::OpenDalStorage
@@ -61,6 +69,19 @@ impl core::clone::Clone for 
iceberg_storage_opendal::CustomAwsCredentialLoader
 pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::clone(&self) -> Self
 impl core::fmt::Debug for iceberg_storage_opendal::CustomAwsCredentialLoader
 pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::fmt(&self, f: &mut 
core::fmt::Formatter<'_>) -> core::fmt::Result
+pub struct iceberg_storage_opendal::OpenDalClientConfig
+impl iceberg_storage_opendal::OpenDalClientConfig
+pub fn 
iceberg_storage_opendal::OpenDalClientConfig::from_properties(properties: 
&std::collections::hash::map::HashMap<alloc::string::String, 
alloc::string::String>) -> iceberg::error::Result<Self>
+impl core::clone::Clone for iceberg_storage_opendal::OpenDalClientConfig
+pub fn iceberg_storage_opendal::OpenDalClientConfig::clone(&self) -> 
iceberg_storage_opendal::OpenDalClientConfig
+impl core::default::Default for iceberg_storage_opendal::OpenDalClientConfig
+pub fn iceberg_storage_opendal::OpenDalClientConfig::default() -> Self
+impl core::fmt::Debug for iceberg_storage_opendal::OpenDalClientConfig
+pub fn iceberg_storage_opendal::OpenDalClientConfig::fmt(&self, f: &mut 
core::fmt::Formatter<'_>) -> core::fmt::Result
+impl serde_core::ser::Serialize for 
iceberg_storage_opendal::OpenDalClientConfig
+pub fn iceberg_storage_opendal::OpenDalClientConfig::serialize<__S>(&self, 
__serializer: __S) -> core::result::Result<<__S as 
serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> 
where __S: serde_core::ser::Serializer
+impl<'de> serde_core::de::Deserialize<'de> for 
iceberg_storage_opendal::OpenDalClientConfig where 
iceberg_storage_opendal::OpenDalClientConfig: core::default::Default
+pub fn 
iceberg_storage_opendal::OpenDalClientConfig::deserialize<__D>(__deserializer: 
__D) -> core::result::Result<Self, <__D as 
serde_core::de::Deserializer>::Error> where __D: 
serde_core::de::Deserializer<'de>
 pub struct iceberg_storage_opendal::OpenDalResolvingStorage
 impl core::fmt::Debug for iceberg_storage_opendal::OpenDalResolvingStorage
 pub fn iceberg_storage_opendal::OpenDalResolvingStorage::fmt(&self, f: &mut 
core::fmt::Formatter<'_>) -> core::fmt::Result
@@ -96,3 +117,4 @@ impl serde_core::ser::Serialize for 
iceberg_storage_opendal::OpenDalResolvingSto
 pub fn 
iceberg_storage_opendal::OpenDalResolvingStorageFactory::serialize<__S>(&self, 
__serializer: __S) -> core::result::Result<<__S as 
serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> 
where __S: serde_core::ser::Serializer
 impl<'de> serde_core::de::Deserialize<'de> for 
iceberg_storage_opendal::OpenDalResolvingStorageFactory
 pub fn 
iceberg_storage_opendal::OpenDalResolvingStorageFactory::deserialize<__D>(__deserializer:
 __D) -> core::result::Result<Self, <__D as 
serde_core::de::Deserializer>::Error> where __D: 
serde_core::de::Deserializer<'de>
+pub const iceberg_storage_opendal::OPENDAL_IO_TIMEOUT_MS: &str
diff --git a/crates/storage/opendal/src/lib.rs 
b/crates/storage/opendal/src/lib.rs
index fe485cb91..f58d6f534 100644
--- a/crates/storage/opendal/src/lib.rs
+++ b/crates/storage/opendal/src/lib.rs
@@ -26,7 +26,9 @@ mod utils;
 
 use std::collections::HashMap;
 use std::collections::hash_map::Entry;
+use std::num::NonZeroU64;
 use std::sync::Arc;
+use std::time::Duration;
 
 use async_trait::async_trait;
 use bytes::Bytes;
@@ -38,6 +40,7 @@ use iceberg::io::{
     StorageFactory,
 };
 use iceberg::{Error, ErrorKind, Result};
+use iceberg_property_macro::Properties;
 use opendal::Operator;
 use opendal::layers::{RetryLayer, TimeoutLayer};
 use serde::{Deserialize, Serialize};
@@ -100,6 +103,54 @@ cfg_if! {
 mod resolving;
 pub use resolving::{OpenDalResolvingStorage, OpenDalResolvingStorageFactory};
 
+/// Timeout in milliseconds for each IO call on a reader, writer, lister or 
deleter, applied per
+/// retry attempt via OpenDAL's `TimeoutLayer::with_io_timeout`. Defaults to 
10 seconds.
+///
+/// OpenDAL's separate `timeout`, which bounds whole control operations such 
as `stat`, is not
+/// affected and stays at its 60-second default.
+pub const OPENDAL_IO_TIMEOUT_MS: &str = "opendal.io-timeout-ms";
+
+/// Matches OpenDAL's `TimeoutLayer` default.
+const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 = NonZeroU64::new(10_000).unwrap();
+
+/// Backend-independent client settings shared by every [`OpenDalStorage`] 
variant.
+#[derive(Clone, Debug, Properties, Serialize, Deserialize)]
+#[serde(default)]
+pub struct OpenDalClientConfig {
+    /// IO timeout in milliseconds.
+    #[property(
+        key = OPENDAL_IO_TIMEOUT_MS,
+        default = DEFAULT_IO_TIMEOUT_MS,
+        parse_with = parse_io_timeout_ms
+    )]
+    io_timeout_ms: NonZeroU64,
+}
+
+impl Default for OpenDalClientConfig {
+    fn default() -> Self {
+        Self {
+            io_timeout_ms: DEFAULT_IO_TIMEOUT_MS,
+        }
+    }
+}
+
+impl OpenDalClientConfig {
+    pub(crate) fn io_timeout(&self) -> Duration {
+        Duration::from_millis(self.io_timeout_ms.get())
+    }
+}
+
+fn parse_io_timeout_ms(value: &str) -> Result<NonZeroU64> {
+    value.parse().map_err(|error| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            "Expected a positive integer number of milliseconds",
+        )
+        .with_context("value", format!("{value:?}"))
+        .with_source(error)
+    })
+}
+
 /// OpenDAL-based storage factory.
 ///
 /// Maps scheme to the corresponding OpenDalStorage storage variant.
@@ -163,35 +214,42 @@ where
 impl StorageFactory for OpenDalStorageFactory {
     #[allow(unused_variables)]
     fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>> {
+        let client_config = 
OpenDalClientConfig::from_properties(config.props())?;
         match self {
             #[cfg(feature = "opendal-memory")]
-            OpenDalStorageFactory::Memory => {
-                Ok(Arc::new(OpenDalStorage::Memory(memory_config_build()?)))
-            }
+            OpenDalStorageFactory::Memory => 
Ok(Arc::new(OpenDalStorage::Memory {
+                operator: memory_config_build()?,
+                client_config,
+            })),
             #[cfg(feature = "opendal-fs")]
-            OpenDalStorageFactory::Fs => Ok(Arc::new(OpenDalStorage::LocalFs)),
+            OpenDalStorageFactory::Fs => Ok(Arc::new(OpenDalStorage::LocalFs { 
client_config })),
             #[cfg(feature = "opendal-s3")]
             OpenDalStorageFactory::S3 {
                 customized_credential_load,
             } => Ok(Arc::new(OpenDalStorage::S3 {
                 config: s3_config_parse(config.props().clone())?.into(),
                 customized_credential_load: customized_credential_load.clone(),
+                client_config,
             })),
             #[cfg(feature = "opendal-gcs")]
             OpenDalStorageFactory::Gcs => Ok(Arc::new(OpenDalStorage::Gcs {
                 config: gcs_config_parse(config.props().clone())?.into(),
+                client_config,
             })),
             #[cfg(feature = "opendal-oss")]
             OpenDalStorageFactory::Oss => Ok(Arc::new(OpenDalStorage::Oss {
                 config: oss_config_parse(config.props().clone())?.into(),
+                client_config,
             })),
             #[cfg(feature = "opendal-azdls")]
             OpenDalStorageFactory::Azdls => Ok(Arc::new(OpenDalStorage::Azdls {
                 config: azdls_config_parse(config.props().clone())?.into(),
+                client_config,
             })),
             #[cfg(feature = "opendal-hf")]
             OpenDalStorageFactory::Hf => Ok(Arc::new(OpenDalStorage::Hf {
                 config: hf_config_parse(config.props().clone())?.into(),
+                client_config,
             })),
             #[cfg(all(
                 not(feature = "opendal-memory"),
@@ -217,14 +275,27 @@ fn default_memory_operator() -> Operator {
 }
 
 /// OpenDAL-based storage implementation.
+///
+/// The serialized representation is not a stable format and may change 
between crate versions.
 #[derive(Clone, Debug, Serialize, Deserialize)]
 pub enum OpenDalStorage {
     /// Memory storage variant.
     #[cfg(feature = "opendal-memory")]
-    Memory(#[serde(skip, default = "self::default_memory_operator")] Operator),
+    Memory {
+        /// Pre-built memory operator.
+        #[serde(skip, default = "self::default_memory_operator")]
+        operator: Operator,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
+    },
     /// Local filesystem storage variant.
     #[cfg(feature = "opendal-fs")]
-    LocalFs,
+    LocalFs {
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
+    },
     /// S3 storage variant.
     ///
     /// Accepts any S3-family URL (`s3://`, `s3a://`, `s3n://`); the scheme is
@@ -236,18 +307,27 @@ pub enum OpenDalStorage {
         /// Custom AWS credential loader.
         #[serde(skip)]
         customized_credential_load: Option<CustomAwsCredentialLoader>,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
     },
     /// GCS storage variant.
     #[cfg(feature = "opendal-gcs")]
     Gcs {
         /// GCS configuration.
         config: Arc<GcsConfig>,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
     },
     /// OSS storage variant.
     #[cfg(feature = "opendal-oss")]
     Oss {
         /// OSS configuration.
         config: Arc<OssConfig>,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
     },
     /// Azure Data Lake Storage variant.
     ///
@@ -259,6 +339,9 @@ pub enum OpenDalStorage {
     Azdls {
         /// Azure DLS configuration.
         config: Arc<AzdlsConfig>,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
     },
     /// HuggingFace Hub storage variant.
     ///
@@ -269,6 +352,9 @@ pub enum OpenDalStorage {
     Hf {
         /// HuggingFace Hub configuration (token + endpoint).
         config: Arc<HfConfig>,
+        /// Backend-independent client settings.
+        #[serde(default)]
+        client_config: OpenDalClientConfig,
     },
 }
 
@@ -293,7 +379,7 @@ impl OpenDalStorage {
         let path = path.as_ref();
         let (operator, relative_path): (Operator, &str) = match self {
             #[cfg(feature = "opendal-memory")]
-            OpenDalStorage::Memory(op) => {
+            OpenDalStorage::Memory { operator: op, .. } => {
                 if let Some(stripped) = path.strip_prefix("memory:/") {
                     (op.clone(), stripped)
                 } else {
@@ -301,7 +387,7 @@ impl OpenDalStorage {
                 }
             }
             #[cfg(feature = "opendal-fs")]
-            OpenDalStorage::LocalFs => {
+            OpenDalStorage::LocalFs { .. } => {
                 let op = fs_config_build()?;
                 if let Some(stripped) = path.strip_prefix("file:/") {
                     (op, stripped)
@@ -313,6 +399,7 @@ impl OpenDalStorage {
             OpenDalStorage::S3 {
                 config,
                 customized_credential_load,
+                ..
             } => {
                 let op = s3_config_build(config, customized_credential_load, 
path)?;
                 let op_info = op.info();
@@ -336,7 +423,7 @@ impl OpenDalStorage {
                 }
             }
             #[cfg(feature = "opendal-gcs")]
-            OpenDalStorage::Gcs { config } => {
+            OpenDalStorage::Gcs { config, .. } => {
                 let operator = gcs_config_build(config, path)?;
                 let prefix = format!("gs://{}/", operator.info().name());
                 if path.starts_with(&prefix) {
@@ -349,7 +436,7 @@ impl OpenDalStorage {
                 }
             }
             #[cfg(feature = "opendal-oss")]
-            OpenDalStorage::Oss { config } => {
+            OpenDalStorage::Oss { config, .. } => {
                 let op = oss_config_build(config, path)?;
                 let prefix = format!("oss://{}/", op.info().name());
                 if path.starts_with(&prefix) {
@@ -362,9 +449,9 @@ impl OpenDalStorage {
                 }
             }
             #[cfg(feature = "opendal-azdls")]
-            OpenDalStorage::Azdls { config } => azdls_create_operator(path, 
config)?,
+            OpenDalStorage::Azdls { config, .. } => 
azdls_create_operator(path, config)?,
             #[cfg(feature = "opendal-hf")]
-            OpenDalStorage::Hf { config } => hf_config_build(config, path)?,
+            OpenDalStorage::Hf { config, .. } => hf_config_build(config, 
path)?,
             #[cfg(all(
                 not(feature = "opendal-s3"),
                 not(feature = "opendal-fs"),
@@ -391,10 +478,41 @@ impl OpenDalStorage {
         // Transient errors are common for object stores; we retry temporary
         // failures with exponential backoff. The retry behavior also
         // benefits non-object-store backends.
-        let operator = 
operator.layer(TimeoutLayer::new()).layer(RetryLayer::new());
+        let operator = operator
+            
.layer(TimeoutLayer::new().with_io_timeout(self.client_config().io_timeout()))
+            .layer(RetryLayer::new());
         Ok((operator, relative_path))
     }
 
+    pub(crate) fn client_config(&self) -> &OpenDalClientConfig {
+        match self {
+            #[cfg(feature = "opendal-memory")]
+            OpenDalStorage::Memory { client_config, .. } => client_config,
+            #[cfg(feature = "opendal-fs")]
+            OpenDalStorage::LocalFs { client_config } => client_config,
+            #[cfg(feature = "opendal-s3")]
+            OpenDalStorage::S3 { client_config, .. } => client_config,
+            #[cfg(feature = "opendal-gcs")]
+            OpenDalStorage::Gcs { client_config, .. } => client_config,
+            #[cfg(feature = "opendal-oss")]
+            OpenDalStorage::Oss { client_config, .. } => client_config,
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorage::Azdls { client_config, .. } => client_config,
+            #[cfg(feature = "opendal-hf")]
+            OpenDalStorage::Hf { client_config, .. } => client_config,
+            #[cfg(all(
+                not(feature = "opendal-memory"),
+                not(feature = "opendal-s3"),
+                not(feature = "opendal-fs"),
+                not(feature = "opendal-gcs"),
+                not(feature = "opendal-oss"),
+                not(feature = "opendal-azdls"),
+                not(feature = "opendal-hf"),
+            ))]
+            _ => unreachable!(),
+        }
+    }
+
     /// Returns a cache key used by `delete_stream` to group paths by storage 
operator.
     ///
     /// For most backends the URL host (bucket name) is sufficient. For HF the 
host
@@ -419,9 +537,11 @@ impl OpenDalStorage {
     pub(crate) fn relativize_path<'a>(&self, path: &'a str) -> Result<&'a str> 
{
         match self {
             #[cfg(feature = "opendal-memory")]
-            OpenDalStorage::Memory(_) => 
Ok(path.strip_prefix("memory:/").unwrap_or(&path[1..])),
+            OpenDalStorage::Memory { .. } => {
+                Ok(path.strip_prefix("memory:/").unwrap_or(&path[1..]))
+            }
             #[cfg(feature = "opendal-fs")]
-            OpenDalStorage::LocalFs => 
Ok(path.strip_prefix("file:/").unwrap_or(&path[1..])),
+            OpenDalStorage::LocalFs { .. } => 
Ok(path.strip_prefix("file:/").unwrap_or(&path[1..])),
             #[cfg(feature = "opendal-s3")]
             OpenDalStorage::S3 { .. } => {
                 let url = url::Url::parse(path)?;
@@ -480,7 +600,7 @@ impl OpenDalStorage {
                 }
             }
             #[cfg(feature = "opendal-azdls")]
-            OpenDalStorage::Azdls { config } => {
+            OpenDalStorage::Azdls { config, .. } => {
                 let azure_path = path.parse::<AzureStoragePath>()?;
                 match_path_with_config(&azure_path, config)?;
                 let relative_path_len = azure_path.path.len();
@@ -687,6 +807,164 @@ impl FileWrite for OpenDalWriter {
 mod tests {
     use super::*;
 
+    fn client_config(value: &str) -> Result<OpenDalClientConfig> {
+        OpenDalClientConfig::from_properties(&HashMap::from([(
+            OPENDAL_IO_TIMEOUT_MS.to_string(),
+            value.to_string(),
+        )]))
+    }
+
+    #[test]
+    fn test_io_timeout_parsing() {
+        let unset = 
OpenDalClientConfig::from_properties(&HashMap::new()).unwrap();
+        assert_eq!(
+            unset.io_timeout(),
+            Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get())
+        );
+
+        let max = u64::MAX.to_string();
+        for (valid, ms) in [("45000", 45_000), ("1", 1), (max.as_str(), 
u64::MAX)] {
+            assert_eq!(
+                client_config(valid).unwrap().io_timeout(),
+                Duration::from_millis(ms),
+                "{valid}"
+            );
+        }
+
+        for invalid in ["0", "-1", "12.5", "abc", "", "18446744073709551616"] {
+            let err = client_config(invalid).unwrap_err().to_string();
+            assert!(err.contains(OPENDAL_IO_TIMEOUT_MS), "{invalid}");
+            assert!(err.contains(&format!("value: {invalid:?}")), "{err}");
+            let reason = 
invalid.parse::<NonZeroU64>().unwrap_err().to_string();
+            assert!(err.contains(&reason), "{err}");
+        }
+    }
+
+    #[test]
+    fn test_default_timeouts_match_opendal() {
+        // `TimeoutLayer` has no getters, so compare through `Debug`.
+        let opendal_default = format!("{:?}", TimeoutLayer::new());
+        let layer = |timeout, io_timeout| {
+            format!(
+                "{:?}",
+                TimeoutLayer::new()
+                    .with_timeout(timeout)
+                    .with_io_timeout(io_timeout)
+            )
+        };
+        let control = Duration::from_secs(60);
+        let io = Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get());
+        assert_eq!(opendal_default, layer(control, io));
+        assert_ne!(opendal_default, layer(control + Duration::from_secs(1), 
io));
+        assert_ne!(
+            opendal_default,
+            layer(control, io + Duration::from_millis(1))
+        );
+    }
+
+    #[cfg(feature = "opendal-s3")]
+    #[test]
+    fn test_client_config_serde_round_trip() {
+        let storage = OpenDalStorage::S3 {
+            config: Arc::new(S3Config::default()),
+            customized_credential_load: None,
+            client_config: client_config("45000").unwrap(),
+        };
+
+        let mut value = serde_json::to_value(&storage).unwrap();
+        let restored: OpenDalStorage = 
serde_json::from_value(value.clone()).unwrap();
+        assert_eq!(
+            restored.client_config().io_timeout(),
+            Duration::from_secs(45)
+        );
+
+        let mut zero = value.clone();
+        zero["S3"]["client_config"]["io_timeout_ms"] = 0.into();
+        assert!(serde_json::from_value::<OpenDalStorage>(zero).is_err());
+
+        value["S3"].as_object_mut().unwrap().remove("client_config");
+        let restored: OpenDalStorage = serde_json::from_value(value).unwrap();
+        assert_eq!(
+            restored.client_config().io_timeout(),
+            Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get())
+        );
+    }
+
+    #[cfg(all(feature = "opendal-fs", feature = "opendal-memory"))]
+    #[test]
+    fn test_old_unit_variant_forms_are_rejected() {
+        for old in [r#""LocalFs""#, r#""Memory""#, r#"{"Memory":null}"#] {
+            assert!(
+                serde_json::from_str::<OpenDalStorage>(old).is_err(),
+                "{old}"
+            );
+        }
+    }
+
+    #[cfg(feature = "opendal-memory")]
+    #[test]
+    fn test_factory_rejects_invalid_io_timeout() {
+        let config = StorageConfig::new().with_prop(OPENDAL_IO_TIMEOUT_MS, 
"nope");
+
+        let err = OpenDalStorageFactory::Memory.build(&config).unwrap_err();
+        assert!(err.to_string().contains(OPENDAL_IO_TIMEOUT_MS));
+    }
+
+    #[test]
+    fn test_factory_propagates_io_timeout() {
+        let config = StorageConfig::new().with_prop(OPENDAL_IO_TIMEOUT_MS, 
"45000");
+        let factories = [
+            #[cfg(feature = "opendal-memory")]
+            OpenDalStorageFactory::Memory,
+            #[cfg(feature = "opendal-fs")]
+            OpenDalStorageFactory::Fs,
+            #[cfg(feature = "opendal-s3")]
+            OpenDalStorageFactory::S3 {
+                customized_credential_load: None,
+            },
+            #[cfg(feature = "opendal-gcs")]
+            OpenDalStorageFactory::Gcs,
+            #[cfg(feature = "opendal-oss")]
+            OpenDalStorageFactory::Oss,
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorageFactory::Azdls,
+            #[cfg(feature = "opendal-hf")]
+            OpenDalStorageFactory::Hf,
+        ];
+        for factory in factories {
+            // `build` returns `dyn Storage`, so read the config back from its 
serialized form.
+            let storage = factory.build(&config).unwrap();
+            let json = serde_json::to_value(&*storage).unwrap();
+            let client_config = json
+                .as_object()
+                .unwrap()
+                .values()
+                .find_map(|variant| variant.get("client_config"))
+                .unwrap_or_else(|| panic!("{json}"));
+            assert_eq!(client_config["io_timeout_ms"], 45_000, "{factory:?}");
+        }
+    }
+
+    #[cfg(feature = "opendal-memory")]
+    #[tokio::test(start_paused = true)]
+    async fn test_io_timeout_reaches_timeout_layer() {
+        use opendal::layers::ConcurrentLimitLayer;
+
+        // A zero concurrency limit stalls every IO call; paused time skips 
the waits.
+        let storage = OpenDalStorage::Memory {
+            operator: 
default_memory_operator().layer(ConcurrentLimitLayer::new(0)),
+            client_config: client_config("45000").unwrap(),
+        };
+
+        let err = storage
+            .read("memory:/stalled")
+            .await
+            .unwrap_err()
+            .to_string();
+        assert!(err.contains("io operation timeout reached"), "{err}");
+        assert!(err.contains("{ timeout: 45 }"), "{err}");
+    }
+
     #[cfg(feature = "opendal-s3")]
     #[derive(Debug)]
     struct EmptyCredentialLoader;
@@ -732,7 +1010,10 @@ mod tests {
 
         // Note: the memory service does report a content length, so this only 
pins the happy
         // path. The counter in `OpenDalWriter` is what covers services that 
don't, such as S3.
-        let storage = 
Arc::new(OpenDalStorage::Memory(default_memory_operator()));
+        let storage = Arc::new(OpenDalStorage::Memory {
+            operator: default_memory_operator(),
+            client_config: OpenDalClientConfig::default(),
+        });
         let path = "memory:///stored-size";
         for plaintext in [
             Bytes::new(),
@@ -760,7 +1041,10 @@ mod tests {
     #[cfg(feature = "opendal-memory")]
     #[test]
     fn test_relativize_path_memory() {
-        let storage = OpenDalStorage::Memory(default_memory_operator());
+        let storage = OpenDalStorage::Memory {
+            operator: default_memory_operator(),
+            client_config: OpenDalClientConfig::default(),
+        };
 
         assert_eq!(
             storage.relativize_path("memory:/path/to/file").unwrap(),
@@ -776,7 +1060,9 @@ mod tests {
     #[cfg(feature = "opendal-fs")]
     #[test]
     fn test_relativize_path_fs() {
-        let storage = OpenDalStorage::LocalFs;
+        let storage = OpenDalStorage::LocalFs {
+            client_config: OpenDalClientConfig::default(),
+        };
 
         assert_eq!(
             storage
@@ -796,6 +1082,7 @@ mod tests {
         let storage = OpenDalStorage::S3 {
             config: Arc::new(S3Config::default()),
             customized_credential_load: None,
+            client_config: OpenDalClientConfig::default(),
         };
 
         // All S3-family schemes are accepted by the same storage instance.
@@ -816,6 +1103,7 @@ mod tests {
     fn test_relativize_path_gcs() {
         let storage = OpenDalStorage::Gcs {
             config: Arc::new(GcsConfig::default()),
+            client_config: OpenDalClientConfig::default(),
         };
 
         assert_eq!(
@@ -831,6 +1119,7 @@ mod tests {
     fn test_relativize_path_gcs_invalid_scheme() {
         let storage = OpenDalStorage::Gcs {
             config: Arc::new(GcsConfig::default()),
+            client_config: OpenDalClientConfig::default(),
         };
 
         assert!(
@@ -845,6 +1134,7 @@ mod tests {
     fn test_relativize_path_oss() {
         let storage = OpenDalStorage::Oss {
             config: Arc::new(OssConfig::default()),
+            client_config: OpenDalClientConfig::default(),
         };
 
         assert_eq!(
@@ -860,6 +1150,7 @@ mod tests {
     fn test_relativize_path_oss_invalid_scheme() {
         let storage = OpenDalStorage::Oss {
             config: Arc::new(OssConfig::default()),
+            client_config: OpenDalClientConfig::default(),
         };
 
         assert!(
@@ -878,6 +1169,7 @@ mod tests {
                 endpoint: 
Some("https://myaccount.dfs.core.windows.net".to_string()),
                 ..Default::default()
             }),
+            client_config: OpenDalClientConfig::default(),
         };
 
         assert_eq!(
diff --git a/crates/storage/opendal/src/resolving.rs 
b/crates/storage/opendal/src/resolving.rs
index 3b99b08fc..51e76da47 100644
--- a/crates/storage/opendal/src/resolving.rs
+++ b/crates/storage/opendal/src/resolving.rs
@@ -33,9 +33,9 @@ use iceberg::{Error, ErrorKind, Result};
 use serde::{Deserialize, Serialize};
 use url::Url;
 
-use crate::OpenDalStorage;
 #[cfg(feature = "opendal-s3")]
 use crate::s3::CustomAwsCredentialLoader;
+use crate::{OpenDalClientConfig, OpenDalStorage};
 
 /// Schemes supported by OpenDalResolvingStorage
 pub const SCHEME_MEMORY: &str = "memory";
@@ -86,6 +86,7 @@ fn build_storage_for_scheme(
     props: &HashMap<String, String>,
     #[cfg(feature = "opendal-s3")] customized_credential_load: 
&Option<CustomAwsCredentialLoader>,
 ) -> Result<OpenDalStorage> {
+    let client_config = OpenDalClientConfig::from_properties(props)?;
     match scheme {
         #[cfg(feature = "opendal-s3")]
         "s3" => {
@@ -93,6 +94,7 @@ fn build_storage_for_scheme(
             Ok(OpenDalStorage::S3 {
                 config: Arc::new(config),
                 customized_credential_load: customized_credential_load.clone(),
+                client_config,
             })
         }
         #[cfg(feature = "opendal-gcs")]
@@ -100,6 +102,7 @@ fn build_storage_for_scheme(
             let config = crate::gcs::gcs_config_parse(props.clone())?;
             Ok(OpenDalStorage::Gcs {
                 config: Arc::new(config),
+                client_config,
             })
         }
         #[cfg(feature = "opendal-oss")]
@@ -107,6 +110,7 @@ fn build_storage_for_scheme(
             let config = crate::oss::oss_config_parse(props.clone())?;
             Ok(OpenDalStorage::Oss {
                 config: Arc::new(config),
+                client_config,
             })
         }
         #[cfg(feature = "opendal-azdls")]
@@ -114,17 +118,22 @@ fn build_storage_for_scheme(
             let config = crate::azdls::azdls_config_parse(props.clone())?;
             Ok(OpenDalStorage::Azdls {
                 config: Arc::new(config),
+                client_config,
             })
         }
         #[cfg(feature = "opendal-fs")]
-        "file" => Ok(OpenDalStorage::LocalFs),
+        "file" => Ok(OpenDalStorage::LocalFs { client_config }),
         #[cfg(feature = "opendal-memory")]
-        "memory" => 
Ok(OpenDalStorage::Memory(crate::memory::memory_config_build()?)),
+        "memory" => Ok(OpenDalStorage::Memory {
+            operator: crate::memory::memory_config_build()?,
+            client_config,
+        }),
         #[cfg(feature = "opendal-hf")]
         "hf" => {
             let config = crate::hf::hf_config_parse(props.clone())?;
             Ok(OpenDalStorage::Hf {
                 config: Arc::new(config),
+                client_config,
             })
         }
         unsupported => Err(Error::new(
@@ -334,7 +343,10 @@ impl Storage for OpenDalResolvingStorage {
 
 #[cfg(test)]
 mod tests {
+    use std::time::Duration;
+
     use super::*;
+    use crate::OPENDAL_IO_TIMEOUT_MS;
 
     #[cfg(feature = "opendal-s3")]
     #[derive(Debug)]
@@ -377,6 +389,39 @@ mod tests {
         }
     }
 
+    #[test]
+    fn test_resolve_propagates_io_timeout() {
+        let mut storage = empty_resolving_storage();
+        storage
+            .props
+            .insert(OPENDAL_IO_TIMEOUT_MS.to_string(), "45000".to_string());
+
+        let paths: &[&str] = &[
+            #[cfg(feature = "opendal-memory")]
+            "memory:/key",
+            #[cfg(feature = "opendal-fs")]
+            "file:/key",
+            #[cfg(feature = "opendal-s3")]
+            "s3://bucket/key",
+            #[cfg(feature = "opendal-gcs")]
+            "gs://bucket/key",
+            #[cfg(feature = "opendal-oss")]
+            "oss://bucket/key",
+            #[cfg(feature = "opendal-azdls")]
+            "abfss://[email protected]/key",
+            #[cfg(feature = "opendal-hf")]
+            "hf://datasets/user/repo/key",
+        ];
+        for path in paths {
+            let resolved = storage.resolve(path).unwrap();
+            assert_eq!(
+                resolved.client_config().io_timeout(),
+                Duration::from_secs(45),
+                "{path}"
+            );
+        }
+    }
+
     #[cfg(feature = "opendal-s3")]
     #[test]
     fn test_resolve_s3_aliases_share_instance() {

Reply via email to