laskoviymishka commented on code in PR #3288:
URL: https://github.com/apache/iceberg-rust/pull/3288#discussion_r4164784806


##########
crates/storage/common/tests/common/mod.rs:
##########
@@ -0,0 +1,276 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Shared test harness and helpers for storage integration suites.
+
+#![allow(dead_code)]
+
+use std::collections::HashMap;
+use std::sync::Arc;
+use std::time::Duration;
+
+use iceberg::io::{
+    FileIO, FileIOBuilder, GCS_NO_AUTH, GCS_SERVICE_HOST, S3_ACCESS_KEY_ID, 
S3_ENDPOINT,
+    S3_PATH_STYLE_ACCESS, S3_REGION, S3_SECRET_ACCESS_KEY,
+};
+use iceberg_storage_opendal::{OpenDalResolvingStorageFactory, 
OpenDalStorageFactory};
+use iceberg_test_utils::{
+    get_gcs_endpoint, get_object_store_endpoint, normalize_test_name, set_up,
+};
+use tempfile::TempDir;
+use tokio::time::sleep;
+
+static FAKE_GCS_BUCKET: &str = "test-bucket";
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum StorageKind {
+    OpenDalS3,
+    OpenDalGcs,
+    OpenDalFs,
+    OpenDalMemory,
+    OpenDalResolving,
+    // TODO: Wire ObjectStoreStorage::S3 once PR #3165 is merged 
(https://github.com/apache/iceberg-rust/pull/3165)
+}
+
+pub struct StorageHarness {
+    pub file_io: FileIO,
+    pub label: &'static str,
+    pub base_path: String,
+    pub _tempdirs: Option<Box<TempDir>>,
+}
+
+fn handle_unreachable_endpoint(kind: &'static str, endpoint: &str) -> 
Option<StorageHarness> {
+    if std::env::var("ICEBERG_REQUIRE_STORAGE")
+        .map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
+        .unwrap_or(false)
+    {
+        panic!(
+            "storage backed '{kind}' is required by ICEBERG_REQUIRE_STORAGE, 
but endpoint '{endpoint}' is unreachable"
+        );
+    }
+    eprintln!("Skipping {kind} storage test: {endpoint} not reachable");
+    None
+}
+
+impl StorageKind {
+    pub const fn as_str(&self) -> &'static str {
+        match self {
+            Self::OpenDalS3 => "opendal_s3",
+            Self::OpenDalGcs => "opendal_gcs",
+            Self::OpenDalFs => "opendal_fs",
+            Self::OpenDalMemory => "opendal_memory",
+            Self::OpenDalResolving => "opendal_resolving",
+        }
+    }
+}
+
+impl std::fmt::Display for StorageKind {
+    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
+        write!(f, "{}", self.as_str())
+    }
+}
+
+const DEFAULT_PROBE_TIMEOUT_MS: u64 = 1000;
+
+fn get_probe_timeout() -> Duration {
+    let ms = std::env::var("ICEBERG_PROBE_TIMEOUT_MS")
+        .ok()
+        .and_then(|v| v.parse().ok())
+        .unwrap_or(DEFAULT_PROBE_TIMEOUT_MS);
+    Duration::from_millis(ms)
+}
+
+/// Fast probe to check if an endpoint service is listening before entering 
retry loops.
+///
+/// Note: Any HTTP response from `.send().await.is_ok()` (including 4xx/5xx) 
is treated
+/// as reachable, as it proves the underlying server is up, listening on the 
port,
+/// and actively responding to HTTP requests.
+pub async fn is_endpoint_reachable(endpoint: &str) -> bool {
+    let Ok(client) = reqwest::Client::builder()
+        .timeout(get_probe_timeout())
+        .build()
+    else {
+        return false;
+    };
+    client.get(endpoint).send().await.is_ok()
+}
+
+async fn wait_until_ready(file_io: &FileIO, check_path: &str, kind: &'static 
str, endpoint: &str) {
+    let mut retries = 0;
+    while retries < 15 {
+        if file_io.exists(check_path).await.unwrap_or(false) {
+            return;
+        }
+        sleep(Duration::from_millis(500)).await;
+        retries += 1;
+    }
+
+    panic!(
+        "Storage backend '{kind}' was reachable at '{endpoint}', but failed 
readiness check on '{check_path}' after 15 retries"
+    );
+}
+
+pub async fn load_storage(kind: StorageKind) -> Option<StorageHarness> {
+    set_up();
+    match kind {
+        StorageKind::OpenDalS3 => load_opendal_s3().await,
+        StorageKind::OpenDalGcs => load_opendal_gcs().await,
+        StorageKind::OpenDalFs => load_opendal_fs().await,
+        StorageKind::OpenDalMemory => load_opendal_memory().await,
+        StorageKind::OpenDalResolving => load_opendal_resolving().await,
+    }
+}
+
+async fn load_opendal_s3() -> Option<StorageHarness> {
+    let object_store_endpoint = get_object_store_endpoint();
+
+    if !is_endpoint_reachable(&object_store_endpoint).await {
+        return handle_unreachable_endpoint("opendal_s3", 
&object_store_endpoint);
+    }
+
+    let file_io = FileIOBuilder::new(Arc::new(OpenDalStorageFactory::S3 {
+        customized_credential_load: None,
+    }))
+    .with_props(vec![
+        (S3_ENDPOINT, object_store_endpoint.clone()),
+        (S3_ACCESS_KEY_ID, "admin".to_string()),
+        (S3_SECRET_ACCESS_KEY, "password".to_string()),
+        (S3_REGION, "us-east-1".to_string()),
+        (S3_PATH_STYLE_ACCESS, "true".to_string()),
+    ])
+    .build();
+
+    wait_until_ready(
+        &file_io,
+        "s3://bucket1/",
+        "opendal_s3",
+        &object_store_endpoint,
+    )
+    .await;
+
+    Some(StorageHarness {
+        file_io,
+        label: "opendal_s3",
+        base_path: "s3://bucket1/".to_string(),
+        _tempdirs: None,
+    })
+}
+
+async fn load_opendal_gcs() -> Option<StorageHarness> {
+    let gcs_endpoint = get_gcs_endpoint();
+
+    if !is_endpoint_reachable(&gcs_endpoint).await {
+        return handle_unreachable_endpoint("opendal_gcs", &gcs_endpoint);
+    }
+
+    let mut bucket_data = HashMap::new();
+    bucket_data.insert("name", FAKE_GCS_BUCKET);
+
+    let client = reqwest::Client::new();
+    let endpoint = format!("{gcs_endpoint}/storage/v1/b");
+    let response = client
+        .post(&endpoint)
+        .json(&bucket_data)
+        .send()
+        .await
+        .unwrap_or_else(|e| {
+            panic!("Failed to send GCS Bucket creation request to 
'{endpoint}': {e}")
+        });
+
+    let status = response.status();
+    if !status.is_success() && status != reqwest::StatusCode::CONFLICT {
+        panic!("failed to create GCS test bucket '{FAKE_GCS_BUCKET}': HTTP 
status {status}");
+    }
+    let file_io = FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs))
+        .with_props(vec![
+            (GCS_SERVICE_HOST, gcs_endpoint.clone()),
+            (GCS_NO_AUTH, "true".to_string()),
+        ])
+        .build();
+
+    let base_path = format!("gs://{FAKE_GCS_BUCKET}/");
+
+    wait_until_ready(&file_io, &base_path, "opendal_gcs", &gcs_endpoint).await;
+
+    Some(StorageHarness {
+        file_io,
+        label: "opendal_gcs",
+        base_path,
+        _tempdirs: None,
+    })
+}
+
+async fn load_opendal_fs() -> Option<StorageHarness> {
+    let temp_dir = TempDir::new().ok()?;

Review Comment:
   `.ok()?` here turns a tempdir failure into `None`, which every caller treats 
as "skip" and returns `Ok(())` — but fs has no endpoint, so 
`ICEBERG_REQUIRE_STORAGE` can't catch it. A full or read-only `$TMPDIR` would 
green the whole fs matrix silently. I'd make the local backends infallible — 
`TempDir::new().expect("create temp dir for fs harness")` — so a setup failure 
is loud.



##########
crates/storage/common/tests/file_io_suite.rs:
##########
@@ -0,0 +1,566 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Shared FileIO integration tests parameterized over storage backends.
+
+mod common;
+
+use bytes::Bytes;
+use common::{StorageHarness, StorageKind, load_storage, unique_path};
+use futures::StreamExt;
+use iceberg::ErrorKind;
+use iceberg::io::FileIO;
+use rstest::rstest;
+
+// ---------------------------------------------------------------------------
+// Helpers
+// ---------------------------------------------------------------------------
+
+fn roundtrip_file_io(file_io: &FileIO) -> FileIO {
+    let serialized = file_io.serialize_all().unwrap();
+    FileIO::deserialize_all(&serialized).unwrap()
+}
+
+// ---------------------------------------------------------------------------
+// Shared Test Execution Bodies
+// ---------------------------------------------------------------------------
+
+async fn run_exists(harness: StorageHarness) -> iceberg::Result<()> {
+    let non_existent = unique_path(&harness, 
"non_existent_file_that_does_not_exist");
+    assert!(!harness.file_io.exists(&non_existent).await?);
+    assert!(harness.file_io.exists(&harness.base_path).await?);
+    Ok(())
+}
+
+async fn run_write(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_write");
+    let _ = harness.file_io.delete(&path).await;
+    assert!(!harness.file_io.exists(&path).await?);
+
+    let output_file = harness.file_io.new_output(&path)?;
+    output_file.write("123".into()).await?;
+    assert!(harness.file_io.exists(&path).await?);
+
+    let _ = harness.file_io.delete(&path).await;
+    Ok(())
+}
+
+async fn run_read(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_read");
+    let _ = harness.file_io.delete(&path).await;
+
+    let output_file = harness.file_io.new_output(&path)?;
+    output_file.write("test_input".into()).await?;
+
+    let input_file = harness.file_io.new_input(&path)?;
+    let buffer = input_file.read().await?;
+    assert_eq!(buffer, "test_input".as_bytes());
+
+    let _ = harness.file_io.delete(&path).await;
+    Ok(())
+}
+
+async fn run_delete(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_delete");
+    let _ = harness.file_io.delete(&path).await;
+
+    harness
+        .file_io
+        .new_output(&path)?
+        .write("delete_me".into())
+        .await?;
+    assert!(harness.file_io.exists(&path).await?);
+
+    harness.file_io.delete(&path).await?;
+    assert!(!harness.file_io.exists(&path).await?);
+    Ok(())
+}
+
+async fn run_delete_nonexistent(harness: StorageHarness) -> 
iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_delete_nonexistent");
+    harness.file_io.delete(&path).await?;
+    Ok(())
+}
+
+async fn run_delete_stream(harness: StorageHarness) -> iceberg::Result<()> {
+    let base = unique_path(&harness, "test_file_io_delete_stream");
+    let paths: Vec<String> = (0..5).map(|i| 
format!("{base}/file-{i}")).collect();
+    for path in &paths {
+        let _ = harness.file_io.delete(path).await;
+        harness
+            .file_io
+            .new_output(path)?
+            .write("delete-me".into())
+            .await?;
+        assert!(harness.file_io.exists(path).await?);
+    }
+    let stream = futures::stream::iter(paths.clone()).boxed();
+    harness.file_io.delete_stream(stream).await?;
+    for path in &paths {
+        assert!(!harness.file_io.exists(path).await?);
+    }
+    Ok(())
+}
+
+async fn run_delete_stream_empty(harness: StorageHarness) -> 
iceberg::Result<()> {
+    let stream = futures::stream::empty().boxed();
+    harness.file_io.delete_stream(stream).await?;
+    Ok(())
+}
+
+async fn run_metadata(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_metadata");
+    let _ = harness.file_io.delete(&path).await;
+    let content = "metadata_test_content";
+    harness
+        .file_io
+        .new_output(&path)?
+        .write(content.into())
+        .await?;
+    let input_file = harness.file_io.new_input(&path)?;
+    let metadata = input_file.metadata().await?;
+    assert_eq!(metadata.size, content.len() as u64);
+    let _ = harness.file_io.delete(&path).await;
+    Ok(())
+}
+
+async fn run_range_read(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_range_read");
+    let _ = harness.file_io.delete(&path).await;
+    let content = b"0123456789abcdef";
+    harness
+        .file_io
+        .new_output(&path)?
+        .write(Bytes::from_static(content))
+        .await?;
+    let input_file = harness.file_io.new_input(&path)?;
+    let reader = input_file.reader().await?;
+    let range_data = reader.read(4..10).await?;
+    assert_eq!(range_data.as_ref(), &content[4..10]);
+    let _ = harness.file_io.delete(&path).await;
+    Ok(())
+}
+
+async fn run_zero_byte_file(harness: StorageHarness) -> iceberg::Result<()> {
+    let path = unique_path(&harness, "test_file_io_zero_byte_file");
+    let _ = harness.file_io.delete(&path).await;
+
+    let output_file = harness.file_io.new_output(&path)?;
+    output_file.write(Bytes::new()).await?;
+
+    assert!(harness.file_io.exists(&path).await?);
+
+    let input_file = harness.file_io.new_input(&path)?;
+    let metadata = input_file.metadata().await?;
+    assert_eq!(metadata.size, 0);
+
+    let data = input_file.read().await?;
+    assert_eq!(data, Bytes::new());
+
+    let _ = harness.file_io.delete(&path).await;
+    Ok(())
+}
+
+async fn run_delete_stream_mixed(harness: StorageHarness) -> 
iceberg::Result<()> {
+    let base = unique_path(&harness, "test_file_io_delete_stream_mixed");
+    let existing_paths: Vec<String> = (0..3).map(|i| 
format!("{base}/exists-{i}")).collect();
+    let nonexistent_paths: Vec<String> = (0..3).map(|i| 
format!("{base}/missing-{i}")).collect();
+
+    for path in &existing_paths {
+        let _ = harness.file_io.delete(path).await;
+        harness
+            .file_io
+            .new_output(path)?
+            .write("data".into())
+            .await?;
+        assert!(harness.file_io.exists(path).await?);
+    }
+    for path in &nonexistent_paths {
+        let _ = harness.file_io.delete(path).await;
+        assert!(!harness.file_io.exists(path).await?);
+    }
+
+    let mut all_paths = existing_paths.clone();
+    all_paths.extend(nonexistent_paths);
+
+    let stream = futures::stream::iter(all_paths).boxed();
+    harness.file_io.delete_stream(stream).await?;
+
+    for path in &existing_paths {
+        assert!(!harness.file_io.exists(path).await?);
+    }
+    Ok(())
+}
+
+async fn run_concurrent_writes(harness: StorageHarness) -> iceberg::Result<()> 
{
+    let base = unique_path(&harness, "test_file_io_concurrent_writes");
+    let mut handles = Vec::new();
+
+    for i in 0..8 {
+        let file_io = harness.file_io.clone();
+        let path = format!("{base}/concurrent-{i}");
+        let payload = format!("payload-{i}");
+
+        handles.push(tokio::spawn(async move {
+            let output = file_io.new_output(&path)?;
+            output.write(payload.clone().into()).await?;
+
+            let input = file_io.new_input(&path)?;
+            let data = input.read().await?;
+            assert_eq!(data, payload.as_bytes());
+
+            let _ = file_io.delete(&path).await;
+            Ok::<(), iceberg::Error>(())
+        }));
+    }
+
+    for handle in handles {
+        handle.await.unwrap()?;
+    }
+
+    Ok(())
+}
+
+// ---------------------------------------------------------------------------
+// Matrix Tests
+// ---------------------------------------------------------------------------
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_exists(#[case] kind: StorageKind) -> iceberg::Result<()> 
{
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_exists(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_write(#[case] kind: StorageKind) -> iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_write(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_read(#[case] kind: StorageKind) -> iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_read(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete(#[case] kind: StorageKind) -> iceberg::Result<()> 
{
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_delete(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete_nonexistent(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_delete_nonexistent(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+// Note: fake-gcs-server emulator does not support batch delete 
(https://github.com/fsouza/fake-gcs-server/issues/1443)
+#[tokio::test]
+async fn test_file_io_delete_stream(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_delete_stream(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete_stream_empty(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_delete_stream_empty(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete_stream_mixed(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_delete_stream_mixed(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete_prefix(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    let prefix = unique_path(&harness, "test_file_io_delete_prefix");
+    let paths: Vec<String> = (0..3).map(|i| 
format!("{prefix}/file-{i}")).collect();
+    for path in &paths {
+        harness
+            .file_io
+            .new_output(path)?
+            .write("data".into())
+            .await?;
+        assert!(harness.file_io.exists(path).await?);
+    }
+    harness.file_io.delete_prefix(&prefix).await?;
+    for path in &paths {
+        assert!(!harness.file_io.exists(path).await?);
+    }
+    Ok(())
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_delete_prefix_nonexistent(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    let prefix = unique_path(&harness, 
"test_file_io_delete_prefix_nonexistent");
+    harness.file_io.delete_prefix(&prefix).await?;
+    Ok(())
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_metadata(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    run_metadata(harness).await
+}
+
+#[rstest]
+#[case::opendal_s3(StorageKind::OpenDalS3)]
+#[case::opendal_gcs(StorageKind::OpenDalGcs)]
+#[case::opendal_fs(StorageKind::OpenDalFs)]
+#[case::opendal_memory(StorageKind::OpenDalMemory)]
+#[tokio::test]
+async fn test_file_io_metadata_nonexistent(#[case] kind: StorageKind) -> 
iceberg::Result<()> {
+    let Some(harness) = load_storage(kind).await else {
+        return Ok(());
+    };
+    let path = unique_path(&harness, "test_file_io_metadata_nonexistent");
+    let input_file = harness.file_io.new_input(&path)?;
+    let Err(err) = input_file.metadata().await else {
+        panic!("expected metadata on nonexistent file to return error");
+    };
+    assert_eq!(err.kind(), ErrorKind::Unexpected);

Review Comment:
   `Unexpected` is OpenDAL's current mapping, not a `FileIO` contract — for 
metadata on a missing file `NotFound` is the semantically right kind, and the 
trait docs don't promise either. The object_store backend from #3165 returning 
`NotFound` would fail this, and a later fix to the OpenDAL mapping would too — 
both false failures. I'd assert just `is_err()` here (same for the double-close 
at `:536` and the credential case in `credential_suite.rs:126`), and if we want 
a specific kind, spec it in the trait docs first.



##########
crates/storage/common/tests/common/mod.rs:
##########
@@ -0,0 +1,276 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Shared test harness and helpers for storage integration suites.
+
+#![allow(dead_code)]
+
+use std::collections::HashMap;
+use std::sync::Arc;
+use std::time::Duration;
+
+use iceberg::io::{
+    FileIO, FileIOBuilder, GCS_NO_AUTH, GCS_SERVICE_HOST, S3_ACCESS_KEY_ID, 
S3_ENDPOINT,
+    S3_PATH_STYLE_ACCESS, S3_REGION, S3_SECRET_ACCESS_KEY,
+};
+use iceberg_storage_opendal::{OpenDalResolvingStorageFactory, 
OpenDalStorageFactory};
+use iceberg_test_utils::{
+    get_gcs_endpoint, get_object_store_endpoint, normalize_test_name, set_up,
+};
+use tempfile::TempDir;
+use tokio::time::sleep;
+
+static FAKE_GCS_BUCKET: &str = "test-bucket";
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum StorageKind {
+    OpenDalS3,
+    OpenDalGcs,
+    OpenDalFs,
+    OpenDalMemory,
+    OpenDalResolving,
+    // TODO: Wire ObjectStoreStorage::S3 once PR #3165 is merged 
(https://github.com/apache/iceberg-rust/pull/3165)
+}
+
+pub struct StorageHarness {
+    pub file_io: FileIO,
+    pub label: &'static str,
+    pub base_path: String,
+    pub _tempdirs: Option<Box<TempDir>>,
+}
+
+fn handle_unreachable_endpoint(kind: &'static str, endpoint: &str) -> 
Option<StorageHarness> {
+    if std::env::var("ICEBERG_REQUIRE_STORAGE")
+        .map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
+        .unwrap_or(false)
+    {
+        panic!(
+            "storage backed '{kind}' is required by ICEBERG_REQUIRE_STORAGE, 
but endpoint '{endpoint}' is unreachable"
+        );
+    }
+    eprintln!("Skipping {kind} storage test: {endpoint} not reachable");

Review Comment:
   The reachable-but-broken cases panic now, which handles the main half of 
last round's concern. But the skip is still the default and nothing sets 
`ICEBERG_REQUIRE_STORAGE` — it's not in the workflows or the Makefile. So if 
`docker-up` fails or the endpoint moves, every S3/GCS/resolving case lands 
here, returns `Ok(())`, and nextest reports it green: that's the "must pass → 
may silently vanish" drop from round 1, just narrowed to the unreachable path.
   
   I'd set `ICEBERG_REQUIRE_STORAGE=1` in the CI step right after `docker-up`, 
or invert the default so unreachable fails unless something like 
`ICEBERG_SKIP_STORAGE=1` is set — inverting is safer for a fork or a local run 
that forgets the flag.
   
   While we're here, the probe is a single 1s GET with the retry loop only 
running after it succeeds, so a slow container reads as absent (a false skip by 
default, a false failure under the gate); a short probe retry would smooth that.



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