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


##########
crates/storage/opendal/src/lib.rs:
##########
@@ -687,6 +811,67 @@ 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_ms(), DEFAULT_IO_TIMEOUT_MS);
+        assert_eq!(client_config("45000").unwrap().io_timeout_ms(), 45_000);
+
+        for invalid in ["0", "-1", "12.5", "abc", ""] {
+            let err = client_config(invalid).unwrap_err().to_string();
+            assert!(err.contains(OPENDAL_IO_TIMEOUT_MS), "{invalid}");
+            assert!(err.contains(&format!("value: {invalid:?}")), "{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.
+        assert_eq!(
+            format!("{:?}", TimeoutLayer::new()),
+            format!(
+                "{:?}",
+                
TimeoutLayer::new().with_io_timeout(Duration::from_millis(DEFAULT_IO_TIMEOUT_MS))
+            ),
+        );

Review Comment:
   Kept it, since `TimeoutLayer` has no getters and both sides go through the 
same derived `Debug`, so a format change alone cannot make them differ. The 
`assert_ne!` from 5e9641c catches a `Debug` that stops printing `io_timeout`.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +102,55 @@ 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 10000 to match OpenDAL's `TimeoutLayer`.

Review Comment:
   Fixed. The doc now links to the new `OPENDAL_IO_TIMEOUT_MS_DEFAULT` instead 
of repeating the number.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +476,47 @@ 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(Duration::from_millis(self.client().io_timeout_ms())),
+            )
+            .layer(RetryLayer::new());

Review Comment:
   The comment just above this hunk, unchanged from `main`, already says 
`TimeoutLayer` must be inside `RetryLayer` so each retry attempt is 
independently bounded.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +103,54 @@ 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 10000 to match OpenDAL's `TimeoutLayer`.

Review Comment:
   Done. The default is now a public `OPENDAL_IO_TIMEOUT_MS_DEFAULT: u64`, 
following the `*_DEFAULT` consts in `iceberg`, and the key's doc links to it 
instead of repeating the number.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +103,54 @@ 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 10000 to match OpenDAL's `TimeoutLayer`.
+///
+/// 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";
+
+/// Matches the `opendal::layers::TimeoutLayer` default.
+const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 = NonZeroU64::new(10_000).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,
+        getter
+    )]
+    io_timeout_ms: NonZeroU64,

Review Comment:
   Went with the `Duration` accessor. The generated getter is gone, and 
`OpenDalClientConfig::io_timeout()` returns a `Duration` by value, so the 
millisecond `NonZeroU64` stays private.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -100,6 +103,54 @@ 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 10000 to match OpenDAL's `TimeoutLayer`.
+///
+/// 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";
+
+/// Matches the `opendal::layers::TimeoutLayer` default.
+const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 = NonZeroU64::new(10_000).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,
+        getter
+    )]
+    io_timeout_ms: NonZeroU64,
+}
+
+impl Default for OpenDalClientConfig {
+    fn default() -> Self {
+        Self {
+            io_timeout_ms: DEFAULT_IO_TIMEOUT_MS,
+        }
+    }
+}
+
+/// Parses one timeout value; the `Properties` derive adds the property-key 
context.
+/// Zero is rejected: it would time every operation out before it starts.
+fn parse_io_timeout_ms(value: &str) -> Result<NonZeroU64> {
+    value.parse().map_err(|_| {

Review Comment:
   Done, using `.with_source(error)` like `parse_pool_property` in the SQL 
catalog. `"0"` now ends with `source: number would be zero for non-zero type`, 
and the parsing test checks the reason for every rejected value.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -221,10 +279,21 @@ fn default_memory_operator() -> Operator {
 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 {

Review Comment:
   Confirmed for direct `OpenDalStorage` serialization, though the old `Memory` 
was the bare string `"Memory"`, not `{"Memory": null}`. `FileIO::serialize_all` 
is unaffected because it serializes the unchanged factory and properties, so I 
documented the break in the description and noted on `OpenDalStorage` that its 
serialized form is not stable across versions, as `FileIO::serialize_all` 
already says.
   



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