comphead commented on code in PR #3263:
URL: https://github.com/apache/iceberg-rust/pull/3263#discussion_r4169559300


##########
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:
   Added the 60s, pinned to OpenDAL's default by 
`test_default_timeouts_match_opendal`, and the doc now names `exists` and 
`metadata` as the calls that keep it. Deletes and listing are not control 
operations in OpenDAL 0.58, which wraps deleters and listers in the IO timeout, 
and a stalled `delete` with a 45000 ms setting fails with `timeout: 45`.
   



##########
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:
   Agreed, and the exposure is wider than absent payloads, since the container 
`#[serde(default)]` makes serde call `default()` on every `client_config` it 
deserializes. Reverted to the struct literal and added 
`test_default_matches_property_defaults`, which compares `default()` with 
`from_properties` on an empty map.
   



##########
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:
   Added `test_old_unit_variant_forms_are_rejected` next to the round trip. It 
checks that `"LocalFs"`, `"Memory"` and `{"Memory": null}` all fail to 
deserialize.
   



##########
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:
   OpenDAL renders the budget as plain seconds 
(`timeout.as_secs_f64().to_string()`), so the error reads `context: { timeout: 
45 }` and `"timeout: 45s"` would never match. The assertion now matches `{ 
timeout: 45 }`, whose closing brace rules out `450`.
   



-- 
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]

Reply via email to