laskoviymishka commented on code in PR #3263:
URL: https://github.com/apache/iceberg-rust/pull/3263#discussion_r4153164908
##########
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:
small thing while we're here — the doc says `10000` but the const is
`DEFAULT_IO_TIMEOUT_MS = 10_000`, so this can drift silently if the default
ever changes. I'd write it as `defaults to 10_000 ms` and let the const stay
the single source of truth.
##########
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:
`map_err(|_| ...)` drops the underlying parse error, so `"abc"`, `"0"`, and
`"18446744073709551616"` all surface the same message — we lose the std detail
that says which one it was. I'd thread it through:
```rust
value.parse().map_err(|e| {
Error::new(ErrorKind::DataInvalid,
format!("Expected a positive integer number of milliseconds: {e}"))
.with_context("value", format!("{value:?}"))
})
```
##########
crates/storage/opendal/src/resolving.rs:
##########
@@ -377,6 +388,18 @@ mod tests {
}
}
+ #[cfg(feature = "opendal-s3")]
+ #[test]
+ fn test_resolve_propagates_io_timeout() {
Review Comment:
`from_properties` runs once before the scheme match, so propagation is
scheme-independent — but this test only covers S3, and a build with just
`opendal-memory` exercises none of it. I'd add a `memory://` (or `file://`)
companion so a refactor that moves parsing into one match arm can't regress the
others silently.
##########
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:
Making `LocalFs` and `Memory` struct variants is a serialized-format break,
not just the API break the `!` marker covers. Externally-tagged serde wrote the
old unit `LocalFs` as the bare string `"LocalFs"` and the old fully-skipped
`Memory` newtype as `{"Memory": null}`; the new struct variants expect
`{"LocalFs": {...}}` / `{"Memory": {...}}`. Any `OpenDalStorage` that was
persisted to catalog properties, sent over RPC, or written to a config file
will fail to deserialize at runtime after upgrade — and typetag makes this part
of the `FileIO` serialization contract, so it's not hypothetical.
I'd either add back-compat deserialization that still accepts the old
bare-string / `null` forms, or, if we're comfortable requiring a regen, call
the wire-format break out explicitly in the description and release notes as
distinct from the API break.
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -391,10 +476,46 @@ 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(
Review Comment:
The propagation tests confirm the value lands in `client_config`, but
nothing exercises this line — that the timeout actually reaches
`TimeoutLayer::with_io_timeout`. A refactor back to a bare
`TimeoutLayer::new()` would pass every current test. It doesn't need a stall
fixture; even asserting the wiring or a note on what's covered would close the
gap.
##########
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:
`NonZeroU64: Copy`, so the generated getter handing back `&NonZeroU64` reads
a bit odd — callers end up writing `*config.io_timeout_ms()`. Since this lands
in `public-api.txt` on a crate that's about to publish, the reference shape is
hard to walk back later.
I'd either add `NonZeroU64` to the macro's `is_copy_type` allowlist so it
returns by value, or expose a `io_timeout(&self) -> Duration` accessor and keep
the millisecond representation internal to the crate.
--
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]