laskoviymishka commented on code in PR #3263:
URL: https://github.com/apache/iceberg-rust/pull/3263#discussion_r4168878631
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +103,65 @@ cfg_if! {
mod resolving;
pub use resolving::{OpenDalResolvingStorage, OpenDalResolvingStorageFactory};
+/// Deadline in milliseconds for one IO operation, and for every method call
on a returned
+/// reader, writer, lister or deleter. Honored by every [`OpenDalStorage`]
backend, where it
+/// defaults to [`OPENDAL_IO_TIMEOUT_MS_DEFAULT`].
+///
+/// Each retry attempt is bounded separately, so it is a per-attempt budget,
not a total one.
+/// Control operations such as `stat` and `rename` are bounded by a separate,
fixed budget.
Review Comment:
While we're documenting this, I'd give the "separate, fixed budget" line a
concrete number — it's OpenDAL's 60s control-op default, and it also governs
`delete` and the lister/deleter open, not just `stat`/`rename`. Otherwise
someone who sets `io-timeout-ms=5000` to bound all their I/O gets surprised
when metadata fetches and snapshot-expiry deletes still hang for up to a minute.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +103,65 @@ cfg_if! {
mod resolving;
pub use resolving::{OpenDalResolvingStorage, OpenDalResolvingStorageFactory};
+/// Deadline in milliseconds for one IO operation, and for every method call
on a returned
+/// reader, writer, lister or deleter. Honored by every [`OpenDalStorage`]
backend, where it
+/// defaults to [`OPENDAL_IO_TIMEOUT_MS_DEFAULT`].
+///
+/// Each retry attempt is bounded separately, so it is a per-attempt budget,
not a total one.
+/// Control operations such as `stat` and `rename` are bounded by a separate,
fixed budget.
+pub const OPENDAL_IO_TIMEOUT_MS: &str = "opendal.io-timeout-ms";
+
+/// Default for [`OPENDAL_IO_TIMEOUT_MS`]. Matches the IO timeout default of
OpenDAL's
+/// `TimeoutLayer`.
+pub const OPENDAL_IO_TIMEOUT_MS_DEFAULT: u64 = 10_000;
+
+/// [`OPENDAL_IO_TIMEOUT_MS_DEFAULT`] as the field type. A zero default fails
to compile.
+const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 =
NonZeroU64::new(OPENDAL_IO_TIMEOUT_MS_DEFAULT).unwrap();
+
+/// Backend-independent client settings, shared by every [`OpenDalStorage`]
variant.
+///
+/// Fields are private, so later settings are additive rather than breaking.
The
+/// container-level serde default lets an older payload deserialize as new
fields appear.
+#[derive(Clone, Debug, Properties, Serialize, Deserialize)]
+#[serde(default)]
+pub struct OpenDalClientConfig {
+ /// Per-attempt deadline for one IO operation, 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 {
+ // Reuse the `#[property]` defaults, so serde and `from_properties`
cannot disagree.
+ Self::from_properties(&HashMap::new()).expect("every client setting
has a default")
Review Comment:
I'd build this default directly (`Self { io_timeout_ms:
DEFAULT_IO_TIMEOUT_MS }`) rather than route through
`from_properties().expect()`. It's safe today, but the moment someone adds a
`#[property]` field without a `default =`, this becomes a production panic on
the deserialization path — serde calls `default()` for every absent
`client_config`, so a stored payload missing that field would panic instead of
erroring. If you want to keep the single-source guarantee, a `#[test]`
asserting `default() == from_properties(&empty)` turns the drift into a CI
failure instead.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -687,6 +825,122 @@ 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(OPENDAL_IO_TIMEOUT_MS_DEFAULT)
+ );
+
+ 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}");
+ // The parse error is kept as the source, so the message says why
it was rejected.
+ let reason =
invalid.parse::<NonZeroU64>().unwrap_err().to_string();
+ assert!(err.contains(&reason), "{err}");
+ }
+ }
+
+ #[test]
+ fn test_default_io_timeout_matches_opendal() {
+ // `TimeoutLayer` has no getters, so compare through `Debug`. An
OpenDAL upgrade that
+ // changes its default fails here instead of silently diverging from
it.
+ let opendal_default = format!("{:?}", TimeoutLayer::new());
+ let with_io_timeout = |ms| {
+ format!(
+ "{:?}",
+ TimeoutLayer::new().with_io_timeout(Duration::from_millis(ms))
+ )
+ };
+ assert_eq!(
+ opendal_default,
+ with_io_timeout(OPENDAL_IO_TIMEOUT_MS_DEFAULT)
+ );
+ // If `Debug` stopped printing `io_timeout`, the check above would
pass vacuously.
+ assert_ne!(
+ opendal_default,
+ with_io_timeout(OPENDAL_IO_TIMEOUT_MS_DEFAULT + 1)
+ );
+ }
+
+ #[cfg(feature = "opendal-s3")]
+ #[test]
+ fn test_client_config_serde_round_trip() {
Review Comment:
The "Breaking (serialized form)" section covers the wire break well and I'm
happy with the docs-only path we settled on last round — not reopening that.
One cheap addition while you're in here: a test asserting an old bare-string
`"LocalFs"` / `{"Memory": null}` payload now `is_err()` would pin the break as
intended behavior rather than an accident, sitting right next to this
round-trip test. Not blocking.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -687,6 +825,122 @@ 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(OPENDAL_IO_TIMEOUT_MS_DEFAULT)
+ );
+
+ 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}");
+ // The parse error is kept as the source, so the message says why
it was rejected.
+ let reason =
invalid.parse::<NonZeroU64>().unwrap_err().to_string();
+ assert!(err.contains(&reason), "{err}");
+ }
+ }
+
+ #[test]
+ fn test_default_io_timeout_matches_opendal() {
+ // `TimeoutLayer` has no getters, so compare through `Debug`. An
OpenDAL upgrade that
+ // changes its default fails here instead of silently diverging from
it.
+ let opendal_default = format!("{:?}", TimeoutLayer::new());
+ let with_io_timeout = |ms| {
+ format!(
+ "{:?}",
+ TimeoutLayer::new().with_io_timeout(Duration::from_millis(ms))
+ )
+ };
+ assert_eq!(
+ opendal_default,
+ with_io_timeout(OPENDAL_IO_TIMEOUT_MS_DEFAULT)
+ );
+ // If `Debug` stopped printing `io_timeout`, the check above would
pass vacuously.
+ assert_ne!(
+ opendal_default,
+ with_io_timeout(OPENDAL_IO_TIMEOUT_MS_DEFAULT + 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)
+ );
+
+ // Deserializing rejects zero too, not only `from_properties`.
+ let mut zero = value.clone();
+ zero["S3"]["client_config"]["io_timeout_ms"] = 0.into();
+ assert!(serde_json::from_value::<OpenDalStorage>(zero).is_err());
+
+ // A payload without `client_config` falls back to the default.
+ 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(OPENDAL_IO_TIMEOUT_MS_DEFAULT)
+ );
+ }
+
+ #[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));
+ }
+
+ #[cfg(feature = "opendal-memory")]
+ #[tokio::test(start_paused = true)]
+ async fn test_io_timeout_reaches_timeout_layer() {
+ use opendal::layers::ConcurrentLimitLayer;
+
+ // A concurrency limit of zero never grants a permit, so every IO call
stalls until
+ // `TimeoutLayer` gives up. Paused time skips the timeouts and the
retry backoff.
+ 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}");
Review Comment:
`"timeout: 45"` also matches `"timeout: 450"`, so this passes if the layer
ever reports 450s — `"timeout: 45s"` with the unit pins it to the value you
actually mean.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]