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]