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


##########
crates/storage/opendal/src/lib.rs:
##########
@@ -163,35 +164,42 @@ where
 impl StorageFactory for OpenDalStorageFactory {
     #[allow(unused_variables)]
     fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>> {
+        let io_timeout_ms = io_timeout_ms_parse(config.props())?;

Review Comment:
   Thanks both. Done in 18a0407. `OpenDalClientConfig` uses 
`#[derive(Properties)]`, and every `OpenDalStorage` variant carries it as 
`client`. Its field is private behind the generated getter, and a075964 adds a 
serde round-trip test for it.
   
   One change from the sketch: the field is `io_timeout_ms: u64` rather than 
`io_timeout: Duration`. The derive's getter returns by value only for types it 
recognizes as `Copy`, which are mostly primitives, so a `Duration` field would 
get a getter returning `&Duration`. A `u64` also serializes as a plain integer 
in the key's own unit rather than as `{secs, nanos}`. Happy to switch to 
`Duration` if you prefer it.
   



##########
crates/iceberg/src/io/storage/config/mod.rs:
##########
@@ -45,6 +45,12 @@ pub use oss::*;
 pub use s3::*;
 use serde::{Deserialize, Serialize};
 
+/// Deadline in milliseconds for one IO operation, and for every method call 
on a returned
+/// reader, writer, lister or deleter. Applies to all backends. Defaults to 
10000.

Review Comment:
   Applied your wording in 18a0407. I also added a short paragraph saying the 
budget is per attempt, and that control operations such as `stat` and `rename` 
have a separate fixed budget.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +426,35 @@ 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.io_timeout()))
+            .layer(RetryLayer::new());
         Ok((operator, relative_path))
     }
 
+    /// Per-IO-operation deadline, from 
[`CLIENT_IO_TIMEOUT_MS`](iceberg::io::CLIENT_IO_TIMEOUT_MS).
+    #[allow(unreachable_patterns)]
+    fn io_timeout(&self) -> Duration {
+        let ms = match self {
+            #[cfg(feature = "opendal-memory")]
+            OpenDalStorage::Memory { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-fs")]
+            OpenDalStorage::LocalFs { io_timeout_ms } => *io_timeout_ms,
+            #[cfg(feature = "opendal-s3")]
+            OpenDalStorage::S3 { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-gcs")]
+            OpenDalStorage::Gcs { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-oss")]
+            OpenDalStorage::Oss { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorage::Azdls { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-hf")]
+            OpenDalStorage::Hf { io_timeout_ms, .. } => *io_timeout_ms,
+            _ => default_io_timeout_ms(),
+        };
+        Duration::from_millis(ms)

Review Comment:
   Done in 18a0407. The method is now `OpenDalStorage::client()`, which returns 
`&OpenDalClientConfig`. Its `_` arm uses the same all-features-off gate as 
`create_operator`, with `opendal-memory` added, and the `allow` is gone. The 
enum has no variants in that build, so a075964 makes the arm `unreachable!()`.
   



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +429,35 @@ 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.io_timeout()))
+            .layer(RetryLayer::new());
         Ok((operator, relative_path))
     }
 
+    /// Per-IO-operation deadline, from 
[`CLIENT_IO_TIMEOUT_MS`](iceberg::io::CLIENT_IO_TIMEOUT_MS).
+    #[allow(unreachable_patterns)]
+    fn io_timeout(&self) -> Duration {
+        let ms = match self {
+            #[cfg(feature = "opendal-memory")]
+            OpenDalStorage::Memory { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-fs")]
+            OpenDalStorage::LocalFs { io_timeout_ms } => *io_timeout_ms,
+            #[cfg(feature = "opendal-s3")]
+            OpenDalStorage::S3 { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-gcs")]
+            OpenDalStorage::Gcs { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-oss")]
+            OpenDalStorage::Oss { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorage::Azdls { io_timeout_ms, .. } => *io_timeout_ms,
+            #[cfg(feature = "opendal-hf")]
+            OpenDalStorage::Hf { io_timeout_ms, .. } => *io_timeout_ms,
+            _ => default_io_timeout_ms(),

Review Comment:
   Fixed in 18a0407. The match now lives in `OpenDalStorage::client()`, which 
has no `#[allow(unreachable_patterns)]`. Its `_` arm only compiles when every 
backend feature is off, so a new variant without an arm is a compile error.
   



##########
crates/storage/opendal/src/utils.rs:
##########
@@ -15,10 +15,34 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::collections::HashMap;
+
+use iceberg::io::CLIENT_IO_TIMEOUT_MS;
+
 pub(crate) fn is_truthy(value: &str) -> bool {
     ["true", "t", "1", "on"].contains(&value.to_lowercase().as_str())
 }
 
+/// Matches the `opendal::layers::TimeoutLayer` default.
+pub(crate) fn default_io_timeout_ms() -> u64 {
+    10_000
+}

Review Comment:
   Kept a named `DEFAULT_IO_TIMEOUT_MS`, which is the shape suggested in 
review. OpenDAL's `TimeoutLayer` has no getter for its default, so a075964 adds 
`test_default_io_timeout_matches_opendal`. It compares the layer's `Debug` 
output and fails if an OpenDAL upgrade changes the default.
   



##########
crates/storage/opendal/src/utils.rs:
##########
@@ -15,10 +15,34 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::collections::HashMap;
+
+use iceberg::io::CLIENT_IO_TIMEOUT_MS;
+
 pub(crate) fn is_truthy(value: &str) -> bool {
     ["true", "t", "1", "on"].contains(&value.to_lowercase().as_str())
 }
 
+/// Matches the `opendal::layers::TimeoutLayer` default.
+pub(crate) fn default_io_timeout_ms() -> u64 {
+    10_000
+}
+
+/// Parse iceberg props to the per-IO-operation timeout.
+pub(crate) fn io_timeout_ms_parse(m: &HashMap<String, String>) -> 
iceberg::Result<u64> {
+    let Some(value) = m.get(CLIENT_IO_TIMEOUT_MS) else {
+        return Ok(default_io_timeout_ms());
+    };
+    // Zero would make every operation time out before it starts.
+    match value.parse::<u64>() {
+        Ok(ms) if ms > 0 => Ok(ms),
+        _ => Err(iceberg::Error::new(
+            iceberg::ErrorKind::DataInvalid,
+            format!("Invalid {CLIENT_IO_TIMEOUT_MS}: {value}, expected a 
positive integer"),
+        )),

Review Comment:
   Addressed in a075964. The rejected value now goes into the error context as 
`{value:?}`, so an empty input renders as `value: ""`. The error also names the 
property key.
   



##########
crates/storage/opendal/src/utils.rs:
##########
@@ -27,3 +51,23 @@ pub(crate) fn from_opendal_error(e: opendal::Error) -> 
iceberg::Error {
     )
     .with_source(e)
 }
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+
+    fn props(value: &str) -> HashMap<String, String> {
+        HashMap::from([(CLIENT_IO_TIMEOUT_MS.to_string(), value.to_string())])
+    }
+
+    #[test]
+    fn test_io_timeout_ms_parse() {
+        assert_eq!(io_timeout_ms_parse(&HashMap::new()).unwrap(), 10_000);

Review Comment:
   Fixed in a075964. The test now asserts against `DEFAULT_IO_TIMEOUT_MS`, and 
`test_default_io_timeout_matches_opendal` pins that constant to OpenDAL's 
actual default.
   



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