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() {