This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-3263-8e81f33daa31e495b94ef9e0bc5ff614355f451c in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
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() {
