From 0a1b3ef36287a2f9178a00c20f5effe74a48b608 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Mon, 21 Sep 2026 07:32:03 +0000 Subject: [PATCH 1/8] first draft: Test Suite 1 first 5 test cases --- tests/storage/src/bidi_read.rs | 314 ++++++++++++++++++++++++++++++++- 1 file changed, 313 insertions(+), 1 deletion(-) diff --git a/tests/storage/src/bidi_read.rs b/tests/storage/src/bidi_read.rs index a886a5eaa1..993cdea2dc 100644 --- a/tests/storage/src/bidi_read.rs +++ b/tests/storage/src/bidi_read.rs @@ -14,9 +14,23 @@ use google_cloud_storage::client::Storage; use google_cloud_storage::model_ext::ReadRange; +use google_cloud_storage::read_object::ReadObjectResponse; pub async fn run(bucket_name: &str) -> anyhow::Result<()> { - let client = Storage::builder().build().await?; + let mut builder = Storage::builder(); + if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { + builder = builder.with_endpoint(endpoint); + } + let client = builder.build().await?; + + // Suite 1: Bidirectional Read live cloud integration tests + test_multiple_ranged_read(&client, bucket_name).await?; + test_read_post_stream_close(&client, bucket_name).await?; + test_zero_copy_read(&client, bucket_name).await?; + test_non_existent_bucket_read(&client).await?; + test_out_of_range(&client, bucket_name).await?; + + // Pre-existing Bidirectional Read tests send(&client, bucket_name).await?; send_and_read(&client, bucket_name).await?; send_and_read_full(&client, bucket_name).await?; @@ -25,6 +39,304 @@ pub async fn run(bucket_name: &str) -> anyhow::Result<()> { Ok(()) } +/// Test Suite 1 - Test 1: Multiple Ranged Read +/// +/// Tests reading an Object across multiple concurrent range read streams over the +/// bidirectional gRPC stream session. Validates that concurrent streams drain properly +/// without deadlock, all received bytes match the expected slices, total length matches, +/// and CRC32C checksum integrity across all ranges matches. +pub async fn test_multiple_ranged_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Test 1/5] Testing Multiple Ranged Read ---"); + const TOTAL_SIZE: usize = 512 * 1024; + let payload = String::from_iter(('a'..='z').cycle().take(TOTAL_SIZE)); + let object_name = "bidi_read/multi_range_source.txt"; + + let write = client + .write_object(bucket_name, object_name, payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let descriptor = client.open_object(bucket_name, &write.name).send().await?; + + // Define 4 non-overlapping segments covering the entire 512 KiB object + // Range 0: [0..64 KiB] (64 KiB) + // Range 1: [64 KiB..192 KiB] (128 KiB) + // Range 2: [192 KiB..384 KiB] (192 KiB) + // Range 3: [384 KiB..512 KiB] (128 KiB) + let range0 = ReadRange::segment(0, 64 * 1024); + let range1 = ReadRange::segment(64 * 1024, 128 * 1024); + let range2 = ReadRange::segment(192 * 1024, 192 * 1024); + let range3 = ReadRange::segment(384 * 1024, 128 * 1024); + + let r0 = descriptor.read_range(range0).await; + let r1 = descriptor.read_range(range1).await; + let r2 = descriptor.read_range(range2).await; + let r3 = descriptor.read_range(range3).await; + + // Concurrent draining is essential: all ranges share the underlying gRPC stream, + // so sequential awaiting would cause buffer starvation and backpressure deadlocks. + let (buf0, buf1, buf2, buf3) = tokio::try_join!( + drain_reader(r0), + drain_reader(r1), + drain_reader(r2), + drain_reader(r3), + )?; + + let payload_bytes = payload.as_bytes(); + assert_eq!(buf0, &payload_bytes[0..64 * 1024]); + assert_eq!(buf1, &payload_bytes[64 * 1024..192 * 1024]); + assert_eq!(buf2, &payload_bytes[192 * 1024..384 * 1024]); + assert_eq!(buf3, &payload_bytes[384 * 1024..512 * 1024]); + + let total_len = buf0.len() + buf1.len() + buf2.len() + buf3.len(); + assert_eq!(total_len, TOTAL_SIZE); + + let crc0 = crc32c::crc32c(&buf0); + assert_eq!(crc0, crc32c::crc32c(&payload_bytes[0..64 * 1024])); + let crc1 = crc32c::crc32c(&buf1); + assert_eq!(crc1, crc32c::crc32c(&payload_bytes[64 * 1024..192 * 1024])); + let crc2 = crc32c::crc32c(&buf2); + assert_eq!(crc2, crc32c::crc32c(&payload_bytes[192 * 1024..384 * 1024])); + let crc3 = crc32c::crc32c(&buf3); + assert_eq!(crc3, crc32c::crc32c(&payload_bytes[384 * 1024..512 * 1024])); + + println!("SUCCESS on Test 1: Multiple Ranged Read (512 KiB across 4 concurrent ranges)"); + Ok(()) +} + +/// Test Suite 1 - Test 2: Read Post Stream Close +/// +/// Verifies stream lifecycle and session isolation: +/// 1. Verifies that once a range reader reaches EOF, subsequent calls to next() idempotently return None. +/// 2. Verifies that dropping an in-flight reader (aborting the range) leaves the underlying ObjectDescriptor +/// healthy and able to issue and read new ranges successfully. +pub async fn test_read_post_stream_close( + client: &Storage, + bucket_name: &str, +) -> anyhow::Result<()> { + println!("--- [Test 2/5] Testing Read Post Stream Close ---"); + let payload = String::from_iter(('a'..='z').cycle().take(100_000)); + let object_name = "bidi_read/post_close_source.txt"; + + let write = client + .write_object(bucket_name, object_name, payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let descriptor = client.open_object(bucket_name, &write.name).send().await?; + + // 1. Read small range to completion (EOF) + let mut reader = descriptor.read_range(ReadRange::head(100)).await; + let mut data = Vec::new(); + while let Some(chunk) = reader.next().await.transpose()? { + data.extend_from_slice(&chunk); + } + assert_eq!(data.len(), 100); + + // Verify idempotency of EOF: next() must consistently return None + assert!(reader.next().await.is_none()); + assert!(reader.next().await.is_none()); + + // 2. Cancellation / abort: drop an in-flight reader without reading it to completion + let unconsumed_reader = descriptor.read_range(ReadRange::segment(500, 10_000)).await; + drop(unconsumed_reader); + + // 3. Verify descriptor session remains fully functional for subsequent reads + let mut subsequent_reader = descriptor.read_range(ReadRange::segment(200, 50)).await; + let mut subsequent_data = Vec::new(); + while let Some(chunk) = subsequent_reader.next().await.transpose()? { + subsequent_data.extend_from_slice(&chunk); + } + assert_eq!(subsequent_data.len(), 50); + assert_eq!(subsequent_data, &payload.as_bytes()[200..250]); + + println!("SUCCESS on Test 2: Read Post Stream Close"); + Ok(()) +} + +/// Test Suite 1 - Test 3: Zero-Copy Read +/// +/// Tests concurrent zero-copy range reads, validating bytes::Bytes buffer access and memory safety. +pub async fn test_zero_copy_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Test 3/5] Testing Zero Copy Read ---"); + const SIZE: usize = 100_000; + let payload = String::from_iter(('a'..='z').cycle().take(SIZE)); + let object_name = "bidi_read/zero_copy_source.txt"; + + let write = client + .write_object(bucket_name, object_name, payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + // Open connection and read first range with Fast Open (send_and_read) + let (descriptor, reader1) = client + .open_object(bucket_name, &write.name) + .send_and_read(ReadRange::segment(0, 50_000)) + .await?; + + // Initiate second range on the open descriptor concurrently + let reader2 = descriptor + .read_range(ReadRange::segment(50_000, 50_000)) + .await; + + // Concurrently collect zero-copy bytes::Bytes chunks + let (chunks1, chunks2) = tokio::try_join!( + collect_zero_copy_chunks(reader1), + collect_zero_copy_chunks(reader2), + )?; + + // Verify chunk properties and reconstruct + let mut combined1 = Vec::new(); + for chunk in &chunks1 { + assert!(!chunk.is_empty(), "chunks should not be empty"); + combined1.extend_from_slice(chunk); + } + assert_eq!(combined1, &payload.as_bytes()[0..50_000]); + + let mut combined2 = Vec::new(); + for chunk in &chunks2 { + assert!(!chunk.is_empty(), "chunks should not be empty"); + combined2.extend_from_slice(chunk); + } + assert_eq!(combined2, &payload.as_bytes()[50_000..100_000]); + + println!("SUCCESS on Test 3: Zero Copy Read"); + Ok(()) +} + +/// Test Suite 1 - Test 4: Non-Existent Bucket Read +/// +/// Tests opening a stream on a non-existent bucket. Verifies that an appropriate +/// error with NotFound status (HTTP 404) is returned. +pub async fn test_non_existent_bucket_read(client: &Storage) -> anyhow::Result<()> { + println!("--- [Test 4/5] Testing Non Existent Bucket Read ---"); + let non_existent_bucket = format!( + "projects/_/buckets/non-existent-bucket-{}", + google_cloud_test_utils::resource_names::random_bucket_id() + ); + + let result = client + .open_object(&non_existent_bucket, "non_existent_object.txt") + .send() + .await; + + match result { + Ok(descriptor) => { + let mut reader = descriptor.read_range(ReadRange::head(100)).await; + let read_res = reader.next().await; + match read_res { + Some(Err(err)) => { + assert_is_not_found(&err); + } + other => anyhow::bail!("expected NotFound error on read_range, got {other:?}"), + } + } + Err(err) => { + assert_is_not_found(&err); + } + } + + println!("SUCCESS on Test 4: Non Existent Bucket Read"); + Ok(()) +} + +fn assert_is_not_found(err: &google_cloud_gax::error::Error) { + if let Some(status) = err.status() { + assert!( + status.code == google_cloud_gax::error::rpc::Code::NotFound + || status.code == google_cloud_gax::error::rpc::Code::PermissionDenied, + "expected NotFound or PermissionDenied rpc code, got {status:?}" + ); + } else if let Some(code) = err.http_status_code() { + assert!( + code == 404 || code == 403, + "expected 404 or 403 HTTP status code, got {code}" + ); + } else { + panic!("expected NotFound or PermissionDenied status, got error: {err:?}"); + } +} + +/// Test Suite 1 - Test 5: Out Of Range Read +/// +/// Tests out-of-bounds range reads beyond object size (offset > size). +/// Ensures appropriate exception/EOF is returned for the invalid range while valid range reads +/// on the same session succeed. +pub async fn test_out_of_range(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Test 5/5] Testing Out Of Range Read ---"); + let payload = String::from_iter(('a'..='z').cycle().take(10_000)); + let object_name = "bidi_read/out_of_range_source.txt"; + + let write = client + .write_object(bucket_name, object_name, payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let descriptor = client.open_object(bucket_name, &write.name).send().await?; + + // 1. Session verification: Verify that a valid range read on this descriptor succeeds + let mut valid_reader = descriptor.read_range(ReadRange::head(50)).await; + let mut valid_data = Vec::new(); + while let Some(chunk) = valid_reader.next().await.transpose()? { + valid_data.extend_from_slice(&chunk); + } + assert_eq!(valid_data.len(), 50); + assert_eq!(valid_data, &payload.as_bytes()[0..50]); + + // 2. Request an out-of-bounds range: offset 50,000 when object size is only 10,000 bytes + let mut oob_reader = descriptor + .read_range(ReadRange::segment(50_000, 1_000)) + .await; + + let oob_res = oob_reader.next().await; + match oob_res { + None => { + println!("Out-of-range read returned immediate EOF"); + } + Some(Err(err)) => { + println!("Out-of-range read returned error as expected: {err:?}"); + let err_str = format!("{err:?}"); + assert!( + err_str.contains("OUT_OF_RANGE") + || err_str.contains("OutOfRange") + || err_str.contains("InvalidArgument"), + "unexpected error message for out of range: {err_str}" + ); + } + Some(Ok(data)) => { + anyhow::bail!( + "unexpected data returned for out of range read: {} bytes", + data.len() + ); + } + } + + println!("SUCCESS on Test 5: Out Of Range Read"); + Ok(()) +} + +async fn drain_reader(mut reader: ReadObjectResponse) -> anyhow::Result> { + let mut buf = Vec::new(); + while let Some(chunk) = reader.next().await.transpose()? { + buf.extend_from_slice(&chunk); + } + Ok(buf) +} + +async fn collect_zero_copy_chunks( + mut reader: ReadObjectResponse, +) -> anyhow::Result> { + let mut chunks = Vec::new(); + while let Some(chunk) = reader.next().await.transpose()? { + chunks.push(chunk); + } + Ok(chunks) +} + async fn send(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { let write = client .write_object( From 2890f01375cc3578a239e2d1fe66bf23031417f8 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Mon, 21 Sep 2026 08:18:36 +0000 Subject: [PATCH 2/8] Change mod struct --- .../conformance.rs} | 249 ++---------------- tests/storage/src/bidi_read/features.rs | 230 ++++++++++++++++ tests/storage/src/bidi_read/mod.rs | 41 +++ 3 files changed, 297 insertions(+), 223 deletions(-) rename tests/storage/src/{bidi_read.rs => bidi_read/conformance.rs} (57%) create mode 100644 tests/storage/src/bidi_read/features.rs create mode 100644 tests/storage/src/bidi_read/mod.rs diff --git a/tests/storage/src/bidi_read.rs b/tests/storage/src/bidi_read/conformance.rs similarity index 57% rename from tests/storage/src/bidi_read.rs rename to tests/storage/src/bidi_read/conformance.rs index 993cdea2dc..035f6d9212 100644 --- a/tests/storage/src/bidi_read.rs +++ b/tests/storage/src/bidi_read/conformance.rs @@ -1,4 +1,4 @@ -// Copyright 2025 Google LLC +// Copyright 2026 Google LLC // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. @@ -12,41 +12,33 @@ // See the License for the specific language governing permissions and // limitations under the License. +//! Cross-SDK Conformance Tests for Bidirectional Read (Test Suite 1). +//! +//! Implements the formal test cases specified in the GCS Bidirectional Read +//! specification and the Rapid Cache Ultra (RCU) integration testing matrix. + use google_cloud_storage::client::Storage; use google_cloud_storage::model_ext::ReadRange; use google_cloud_storage::read_object::ReadObjectResponse; -pub async fn run(bucket_name: &str) -> anyhow::Result<()> { - let mut builder = Storage::builder(); - if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { - builder = builder.with_endpoint(endpoint); - } - let client = builder.build().await?; - - // Suite 1: Bidirectional Read live cloud integration tests - test_multiple_ranged_read(&client, bucket_name).await?; - test_read_post_stream_close(&client, bucket_name).await?; - test_zero_copy_read(&client, bucket_name).await?; - test_non_existent_bucket_read(&client).await?; - test_out_of_range(&client, bucket_name).await?; - - // Pre-existing Bidirectional Read tests - send(&client, bucket_name).await?; - send_and_read(&client, bucket_name).await?; - send_and_read_full(&client, bucket_name).await?; - send_and_read_md5(&client, bucket_name).await?; - send_and_read_gzip(&client, bucket_name).await?; +/// Runs all 5 live cloud conformance tests for Bidirectional Read. +pub async fn run(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + test_multiple_ranged_read(client, bucket_name).await?; + test_read_post_stream_close(client, bucket_name).await?; + test_zero_copy_read(client, bucket_name).await?; + test_non_existent_bucket_read(client).await?; + test_out_of_range(client, bucket_name).await?; Ok(()) } /// Test Suite 1 - Test 1: Multiple Ranged Read /// -/// Tests reading an Object across multiple concurrent range read streams over the +/// Tests reading an object across multiple concurrent range read streams over the /// bidirectional gRPC stream session. Validates that concurrent streams drain properly /// without deadlock, all received bytes match the expected slices, total length matches, /// and CRC32C checksum integrity across all ranges matches. pub async fn test_multiple_ranged_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - println!("--- [Test 1/5] Testing Multiple Ranged Read ---"); + println!("--- [Conformance 1/5] Testing Multiple Ranged Read ---"); const TOTAL_SIZE: usize = 512 * 1024; let payload = String::from_iter(('a'..='z').cycle().take(TOTAL_SIZE)); let object_name = "bidi_read/multi_range_source.txt"; @@ -59,7 +51,7 @@ pub async fn test_multiple_ranged_read(client: &Storage, bucket_name: &str) -> a let descriptor = client.open_object(bucket_name, &write.name).send().await?; - // Define 4 non-overlapping segments covering the entire 512 KiB object + // Define 4 non-overlapping segments covering the entire 512 KiB object: // Range 0: [0..64 KiB] (64 KiB) // Range 1: [64 KiB..192 KiB] (128 KiB) // Range 2: [192 KiB..384 KiB] (192 KiB) @@ -101,7 +93,7 @@ pub async fn test_multiple_ranged_read(client: &Storage, bucket_name: &str) -> a let crc3 = crc32c::crc32c(&buf3); assert_eq!(crc3, crc32c::crc32c(&payload_bytes[384 * 1024..512 * 1024])); - println!("SUCCESS on Test 1: Multiple Ranged Read (512 KiB across 4 concurrent ranges)"); + println!("SUCCESS on Conformance 1: Multiple Ranged Read (512 KiB across 4 concurrent ranges)"); Ok(()) } @@ -115,7 +107,7 @@ pub async fn test_read_post_stream_close( client: &Storage, bucket_name: &str, ) -> anyhow::Result<()> { - println!("--- [Test 2/5] Testing Read Post Stream Close ---"); + println!("--- [Conformance 2/5] Testing Read Post Stream Close ---"); let payload = String::from_iter(('a'..='z').cycle().take(100_000)); let object_name = "bidi_read/post_close_source.txt"; @@ -152,7 +144,7 @@ pub async fn test_read_post_stream_close( assert_eq!(subsequent_data.len(), 50); assert_eq!(subsequent_data, &payload.as_bytes()[200..250]); - println!("SUCCESS on Test 2: Read Post Stream Close"); + println!("SUCCESS on Conformance 2: Read Post Stream Close"); Ok(()) } @@ -160,7 +152,7 @@ pub async fn test_read_post_stream_close( /// /// Tests concurrent zero-copy range reads, validating bytes::Bytes buffer access and memory safety. pub async fn test_zero_copy_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - println!("--- [Test 3/5] Testing Zero Copy Read ---"); + println!("--- [Conformance 3/5] Testing Zero Copy Read ---"); const SIZE: usize = 100_000; let payload = String::from_iter(('a'..='z').cycle().take(SIZE)); let object_name = "bidi_read/zero_copy_source.txt"; @@ -203,16 +195,16 @@ pub async fn test_zero_copy_read(client: &Storage, bucket_name: &str) -> anyhow: } assert_eq!(combined2, &payload.as_bytes()[50_000..100_000]); - println!("SUCCESS on Test 3: Zero Copy Read"); + println!("SUCCESS on Conformance 3: Zero Copy Read"); Ok(()) } /// Test Suite 1 - Test 4: Non-Existent Bucket Read /// /// Tests opening a stream on a non-existent bucket. Verifies that an appropriate -/// error with NotFound status (HTTP 404) is returned. +/// error with NotFound status (HTTP 404) or PermissionDenied (allowlist check) is returned. pub async fn test_non_existent_bucket_read(client: &Storage) -> anyhow::Result<()> { - println!("--- [Test 4/5] Testing Non Existent Bucket Read ---"); + println!("--- [Conformance 4/5] Testing Non Existent Bucket Read ---"); let non_existent_bucket = format!( "projects/_/buckets/non-existent-bucket-{}", google_cloud_test_utils::resource_names::random_bucket_id() @@ -239,7 +231,7 @@ pub async fn test_non_existent_bucket_read(client: &Storage) -> anyhow::Result<( } } - println!("SUCCESS on Test 4: Non Existent Bucket Read"); + println!("SUCCESS on Conformance 4: Non Existent Bucket Read"); Ok(()) } @@ -266,7 +258,7 @@ fn assert_is_not_found(err: &google_cloud_gax::error::Error) { /// Ensures appropriate exception/EOF is returned for the invalid range while valid range reads /// on the same session succeed. pub async fn test_out_of_range(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - println!("--- [Test 5/5] Testing Out Of Range Read ---"); + println!("--- [Conformance 5/5] Testing Out Of Range Read ---"); let payload = String::from_iter(('a'..='z').cycle().take(10_000)); let object_name = "bidi_read/out_of_range_source.txt"; @@ -315,7 +307,7 @@ pub async fn test_out_of_range(client: &Storage, bucket_name: &str) -> anyhow::R } } - println!("SUCCESS on Test 5: Out Of Range Read"); + println!("SUCCESS on Conformance 5: Out Of Range Read"); Ok(()) } @@ -336,192 +328,3 @@ async fn collect_zero_copy_chunks( } Ok(chunks) } - -async fn send(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - let write = client - .write_object( - bucket_name, - "basic/source.txt", - String::from_iter((0..100_000).map(|_| 'a')), - ) - .set_if_generation_match(0) - .send_unbuffered() - .await?; - - let open = client.open_object(bucket_name, &write.name).send().await?; - tracing::info!("open returns: {open:?}"); - let got = open.object(); - let mut want = write.clone(); - // This field is a mismatch, but both `Some(false)` and `None` represent - // the same value. - want.event_based_hold = want.event_based_hold.or(Some(false)); - // There is a submillisecond difference, maybe rounding? - want.finalize_time = got.finalize_time; - assert_eq!(got, want); - - let mut reader = open.read_range(ReadRange::head(100)).await; - let mut count = 0_usize; - while let Some(r) = reader.next().await.transpose()? { - tracing::info!("received {} bytes", r.len()); - count += r.len(); - } - assert_eq!(count, 100_usize); - - Ok(()) -} - -pub async fn send_and_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - let payload = String::from_iter(('a'..='z').cycle().take(100_000)); - let write = client - .write_object(bucket_name, "open_and_read/source.txt", payload.clone()) - .set_if_generation_match(0) - .send_unbuffered() - .await?; - - let (descriptor, mut reader) = client - .open_object(bucket_name, &write.name) - .send_and_read(ReadRange::tail(100)) - .await?; - tracing::info!("object: {:?}", descriptor.object()); - tracing::info!("headers: {:?}", descriptor.headers()); - tracing::info!("reader: {:?}", reader); - let got = descriptor.object(); - let mut want = write.clone(); - // This field is a mismatch, but both `Some(false)` and `None` represent - // the same value. - want.event_based_hold = want.event_based_hold.or(Some(false)); - // There is a submillisecond difference, maybe rounding? - want.finalize_time = got.finalize_time; - assert_eq!(got, want); - - let mut data = Vec::new(); - while let Some(r) = reader.next().await.transpose()? { - tracing::info!("received {} bytes", r.len()); - data.extend_from_slice(&r); - } - assert_eq!(data, &payload.as_bytes()[(payload.len() - 100)..]); - - Ok(()) -} - -pub async fn send_and_read_md5(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - let payload = String::from_iter(('a'..='z').cycle().take(100_000)); - let write = client - .write_object(bucket_name, "open_and_read_md5/source.txt", payload.clone()) - .set_if_generation_match(0) - .send_unbuffered() - .await?; - - let (descriptor, mut reader) = client - .open_object(bucket_name, &write.name) - .compute_md5() - .send_and_read(ReadRange::all()) - .await?; - tracing::info!("object: {:?}", descriptor.object()); - tracing::info!("headers: {:?}", descriptor.headers()); - tracing::info!("reader: {:?}", reader); - let got = descriptor.object(); - let mut want = write.clone(); - want.event_based_hold = want.event_based_hold.or(Some(false)); - want.finalize_time = got.finalize_time; - assert_eq!(got, want); - - let mut data = Vec::new(); - while let Some(r) = reader.next().await.transpose()? { - tracing::info!("received {} bytes", r.len()); - data.extend_from_slice(&r); - } - assert_eq!(data, payload.as_bytes()); - - Ok(()) -} - -pub async fn send_and_read_full(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - let payload = String::from_iter(('a'..='z').cycle().take(100_000)); - let write = client - .write_object( - bucket_name, - "open_and_read_full/source.txt", - payload.clone(), - ) - .set_if_generation_match(0) - .send_unbuffered() - .await?; - - let (descriptor, mut reader) = client - .open_object(bucket_name, &write.name) - .send_and_read(ReadRange::all()) - .await?; - tracing::info!("object: {:?}", descriptor.object()); - tracing::info!("headers: {:?}", descriptor.headers()); - tracing::info!("reader: {:?}", reader); - let got = descriptor.object(); - let mut want = write.clone(); - want.event_based_hold = want.event_based_hold.or(Some(false)); - want.finalize_time = got.finalize_time; - assert_eq!(got, want); - - let mut data = Vec::new(); - while let Some(r) = reader.next().await.transpose()? { - tracing::info!("received {} bytes", r.len()); - data.extend_from_slice(&r); - } - assert_eq!(data, payload.as_bytes()); - - Ok(()) -} - -/// This test verifies the checksum validation behavior for gzip-encoded objects -/// over the gRPC Bidi read stream. -/// -/// Unlike the JSON REST API, which often transcodes (decompresses) gzip objects -/// on the fly, the gRPC Bidi read stream delivers the raw, compressed bytes directly. -/// Because no on-the-fly decompression occurs, the CRC32C checksum of the received -/// chunks will naturally match the server's stored checksum of the compressed object. -/// -/// We explicitly expect `RangeReader`'s automatic checksum validation to succeed -/// without throwing a `ChecksumMismatch` error, proving that we do not need to -/// bypass checksum validation for `content-encoding: gzip` objects in gRPC. -pub async fn send_and_read_gzip(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { - use std::io::Write; - let payload = String::from_iter(('a'..='z').cycle().take(100_000)); - - // Compress the payload - let mut e = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default()); - e.write_all(payload.as_bytes())?; - let compressed_payload = e.finish()?; - - let write = client - .write_object( - bucket_name, - "open_and_read_gzip/source.txt", - bytes::Bytes::from_owner(compressed_payload.clone()), - ) - .set_if_generation_match(0) - .set_content_encoding("gzip") - .send_unbuffered() - .await?; - - let (descriptor, mut reader) = client - .open_object(bucket_name, &write.name) - .send_and_read(ReadRange::all()) - .await?; - tracing::info!("object: {:?}", descriptor.object()); - tracing::info!("headers: {:?}", descriptor.headers()); - tracing::info!("reader: {:?}", reader); - let got = descriptor.object(); - let mut want = write.clone(); - want.event_based_hold = want.event_based_hold.or(Some(false)); - want.finalize_time = got.finalize_time; - assert_eq!(got, want); - - let mut data = Vec::new(); - while let Some(r) = reader.next().await.transpose()? { - tracing::info!("received {} bytes", r.len()); - data.extend_from_slice(&r); - } - // Verify we received the EXACT compressed payload, meaning gRPC did not decompress it. - assert_eq!(data, compressed_payload); - - Ok(()) -} diff --git a/tests/storage/src/bidi_read/features.rs b/tests/storage/src/bidi_read/features.rs new file mode 100644 index 0000000000..4706799b9b --- /dev/null +++ b/tests/storage/src/bidi_read/features.rs @@ -0,0 +1,230 @@ +// Copyright 2025 Google LLC +// +// Licensed 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 +// +// https://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. + +//! Client Feature Tests for Bidirectional Read. +//! +//! Validates client builder options, encodings, and request parameters +//! such as `.compute_md5()`, `content-encoding: gzip`, and metadata matching. + +use google_cloud_storage::client::Storage; +use google_cloud_storage::model_ext::ReadRange; + +/// Runs all client feature regression tests for Bidirectional Read. +pub async fn run(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + send(client, bucket_name).await?; + send_and_read(client, bucket_name).await?; + send_and_read_full(client, bucket_name).await?; + send_and_read_md5(client, bucket_name).await?; + send_and_read_gzip(client, bucket_name).await?; + Ok(()) +} + +pub async fn send(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Features 1/5] Testing Basic Send & Metadata ---"); + let write = client + .write_object( + bucket_name, + "basic/source.txt", + String::from_iter((0..100_000).map(|_| 'a')), + ) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let open = client.open_object(bucket_name, &write.name).send().await?; + tracing::info!("open returns: {open:?}"); + let got = open.object(); + let mut want = write.clone(); + // This field is a mismatch, but both `Some(false)` and `None` represent + // the same value. + want.event_based_hold = want.event_based_hold.or(Some(false)); + // There is a submillisecond difference, maybe rounding? + want.finalize_time = got.finalize_time; + assert_eq!(got, want); + + let mut reader = open.read_range(ReadRange::head(100)).await; + let mut count = 0_usize; + while let Some(r) = reader.next().await.transpose()? { + tracing::info!("received {} bytes", r.len()); + count += r.len(); + } + assert_eq!(count, 100_usize); + + println!("SUCCESS on Features 1: Basic Send & Metadata"); + Ok(()) +} + +pub async fn send_and_read(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Features 2/5] Testing Send & Read Tail Range ---"); + let payload = String::from_iter(('a'..='z').cycle().take(100_000)); + let write = client + .write_object(bucket_name, "open_and_read/source.txt", payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let (descriptor, mut reader) = client + .open_object(bucket_name, &write.name) + .send_and_read(ReadRange::tail(100)) + .await?; + tracing::info!("object: {:?}", descriptor.object()); + tracing::info!("headers: {:?}", descriptor.headers()); + tracing::info!("reader: {:?}", reader); + let got = descriptor.object(); + let mut want = write.clone(); + // This field is a mismatch, but both `Some(false)` and `None` represent + // the same value. + want.event_based_hold = want.event_based_hold.or(Some(false)); + // There is a submillisecond difference, maybe rounding? + want.finalize_time = got.finalize_time; + assert_eq!(got, want); + + let mut data = Vec::new(); + while let Some(r) = reader.next().await.transpose()? { + tracing::info!("received {} bytes", r.len()); + data.extend_from_slice(&r); + } + assert_eq!(data, &payload.as_bytes()[(payload.len() - 100)..]); + + println!("SUCCESS on Features 2: Send & Read Tail Range"); + Ok(()) +} + +pub async fn send_and_read_md5(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Features 3/5] Testing Send & Read with MD5 ---"); + let payload = String::from_iter(('a'..='z').cycle().take(100_000)); + let write = client + .write_object(bucket_name, "open_and_read_md5/source.txt", payload.clone()) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let (descriptor, mut reader) = client + .open_object(bucket_name, &write.name) + .compute_md5() + .send_and_read(ReadRange::all()) + .await?; + tracing::info!("object: {:?}", descriptor.object()); + tracing::info!("headers: {:?}", descriptor.headers()); + tracing::info!("reader: {:?}", reader); + let got = descriptor.object(); + let mut want = write.clone(); + want.event_based_hold = want.event_based_hold.or(Some(false)); + want.finalize_time = got.finalize_time; + assert_eq!(got, want); + + let mut data = Vec::new(); + while let Some(r) = reader.next().await.transpose()? { + tracing::info!("received {} bytes", r.len()); + data.extend_from_slice(&r); + } + assert_eq!(data, payload.as_bytes()); + + println!("SUCCESS on Features 3: Send & Read with MD5"); + Ok(()) +} + +pub async fn send_and_read_full(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Features 4/5] Testing Send & Read Full Object ---"); + let payload = String::from_iter(('a'..='z').cycle().take(100_000)); + let write = client + .write_object( + bucket_name, + "open_and_read_full/source.txt", + payload.clone(), + ) + .set_if_generation_match(0) + .send_unbuffered() + .await?; + + let (descriptor, mut reader) = client + .open_object(bucket_name, &write.name) + .send_and_read(ReadRange::all()) + .await?; + tracing::info!("object: {:?}", descriptor.object()); + tracing::info!("headers: {:?}", descriptor.headers()); + tracing::info!("reader: {:?}", reader); + let got = descriptor.object(); + let mut want = write.clone(); + want.event_based_hold = want.event_based_hold.or(Some(false)); + want.finalize_time = got.finalize_time; + assert_eq!(got, want); + + let mut data = Vec::new(); + while let Some(r) = reader.next().await.transpose()? { + tracing::info!("received {} bytes", r.len()); + data.extend_from_slice(&r); + } + assert_eq!(data, payload.as_bytes()); + + println!("SUCCESS on Features 4: Send & Read Full Object"); + Ok(()) +} + +/// This test verifies the checksum validation behavior for gzip-encoded objects +/// over the gRPC Bidi read stream. +/// +/// Unlike the JSON REST API, which often transcodes (decompresses) gzip objects +/// on the fly, the gRPC Bidi read stream delivers the raw, compressed bytes directly. +/// Because no on-the-fly decompression occurs, the CRC32C checksum of the received +/// chunks will naturally match the server's stored checksum of the compressed object. +/// +/// We explicitly expect `RangeReader`'s automatic checksum validation to succeed +/// without throwing a `ChecksumMismatch` error, proving that we do not need to +/// bypass checksum validation for `content-encoding: gzip` objects in gRPC. +pub async fn send_and_read_gzip(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + println!("--- [Features 5/5] Testing Send & Read Gzip Encoded Object ---"); + use std::io::Write; + let payload = String::from_iter(('a'..='z').cycle().take(100_000)); + + // Compress the payload + let mut e = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default()); + e.write_all(payload.as_bytes())?; + let compressed_payload = e.finish()?; + + let write = client + .write_object( + bucket_name, + "open_and_read_gzip/source.txt", + bytes::Bytes::from_owner(compressed_payload.clone()), + ) + .set_if_generation_match(0) + .set_content_encoding("gzip") + .send_unbuffered() + .await?; + + let (descriptor, mut reader) = client + .open_object(bucket_name, &write.name) + .send_and_read(ReadRange::all()) + .await?; + tracing::info!("object: {:?}", descriptor.object()); + tracing::info!("headers: {:?}", descriptor.headers()); + tracing::info!("reader: {:?}", reader); + let got = descriptor.object(); + let mut want = write.clone(); + want.event_based_hold = want.event_based_hold.or(Some(false)); + want.finalize_time = got.finalize_time; + assert_eq!(got, want); + + let mut data = Vec::new(); + while let Some(r) = reader.next().await.transpose()? { + tracing::info!("received {} bytes", r.len()); + data.extend_from_slice(&r); + } + // Verify we received the EXACT compressed payload, meaning gRPC did not decompress it. + assert_eq!(data, compressed_payload); + + println!("SUCCESS on Features 5: Send & Read Gzip Encoded Object"); + Ok(()) +} diff --git a/tests/storage/src/bidi_read/mod.rs b/tests/storage/src/bidi_read/mod.rs new file mode 100644 index 0000000000..c0f0b0bd17 --- /dev/null +++ b/tests/storage/src/bidi_read/mod.rs @@ -0,0 +1,41 @@ +// Copyright 2025 Google LLC +// +// Licensed 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 +// +// https://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. + +//! Bidirectional Read Integration Tests. +//! +//! Separated into two distinct test suites: +//! - [conformance]: Formal cross-SDK conformance tests for Bidi Read sessions (Suite 1). +//! - [features]: Client builder options, encodings, and backward-compatibility regression tests. + +pub mod conformance; +pub mod features; + +use google_cloud_storage::client::Storage; + +/// Default runner that executes both conformance and feature test suites. +pub async fn run(bucket_name: &str) -> anyhow::Result<()> { + let mut builder = Storage::builder(); + if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { + builder = builder.with_endpoint(endpoint); + } + let client = builder.build().await?; + + println!("\n=== Running Bidirectional Read Conformance Suite (Suite 1) ==="); + conformance::run(&client, bucket_name).await?; + + println!("\n=== Running Bidirectional Read Features Suite ==="); + features::run(&client, bucket_name).await?; + + Ok(()) +} From 22d2849bc50ea0739ac9273b22d4dcfaa805c5b1 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Mon, 21 Sep 2026 09:18:16 +0000 Subject: [PATCH 3/8] feat(storage): parameterize bidi read conformance tests into scenario groups --- tests/storage/src/bidi_read/conformance.rs | 11 +++ tests/storage/src/lib.rs | 48 ++++++++- tests/storage/tests/driver.rs | 110 +++++++++++++++++++++ 3 files changed, 168 insertions(+), 1 deletion(-) diff --git a/tests/storage/src/bidi_read/conformance.rs b/tests/storage/src/bidi_read/conformance.rs index 035f6d9212..9554b359a9 100644 --- a/tests/storage/src/bidi_read/conformance.rs +++ b/tests/storage/src/bidi_read/conformance.rs @@ -23,11 +23,22 @@ use google_cloud_storage::read_object::ReadObjectResponse; /// Runs all 5 live cloud conformance tests for Bidirectional Read. pub async fn run(client: &Storage, bucket_name: &str) -> anyhow::Result<()> { + run_with_scenario(client, bucket_name, "Default").await +} + +/// Runs all 5 live cloud conformance tests for Bidirectional Read with a specific scenario name. +pub async fn run_with_scenario( + client: &Storage, + bucket_name: &str, + scenario_name: &str, +) -> anyhow::Result<()> { + println!("\n=== [Scenario: {scenario_name}] Running Bidi Read Conformance Suite ==="); test_multiple_ranged_read(client, bucket_name).await?; test_read_post_stream_close(client, bucket_name).await?; test_zero_copy_read(client, bucket_name).await?; test_non_existent_bucket_read(client).await?; test_out_of_range(client, bucket_name).await?; + println!("=== [Scenario: {scenario_name}] Conformance Suite Completed Successfully ===\n"); Ok(()) } diff --git a/tests/storage/src/lib.rs b/tests/storage/src/lib.rs index eee7fa75de..165db67fe9 100644 --- a/tests/storage/src/lib.rs +++ b/tests/storage/src/lib.rs @@ -29,7 +29,7 @@ use google_cloud_lro::Poller; pub use google_cloud_storage::builder::storage::ClientBuilder as StorageBuilder; use google_cloud_storage::builder::storage::SignedUrlBuilder; pub use google_cloud_storage::builder::storage_control::ClientBuilder as StorageControlBuilder; -use google_cloud_storage::client::StorageControl; +use google_cloud_storage::client::{Storage, StorageControl}; use google_cloud_storage::model::Bucket; use google_cloud_storage::model::bucket::iam_config::UniformBucketLevelAccess; use google_cloud_storage::model::bucket::{HierarchicalNamespace, IamConfig}; @@ -42,6 +42,52 @@ pub use storage_samples::{ cleanup_stale_buckets, create_test_bucket, create_test_hns_bucket, create_test_rapid_bucket, }; +pub async fn build_storage_client() -> Result { + let mut builder = Storage::builder(); + if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { + builder = builder.with_endpoint(endpoint); + } + Ok(builder.build().await?) +} + +pub async fn build_non_colocated_storage_client(off_zone: &str) -> Result { + let mut builder = Storage::builder(); + if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { + builder = builder.with_endpoint(endpoint); + } else { + builder = builder.with_endpoint(format!("https://{off_zone}-storage.googleapis.com")); + } + Ok(builder.build().await?) +} + +pub async fn create_test_regional_rapid_bucket() -> Result<(StorageControl, Bucket)> { + let project_id = project_id()?; + let control = StorageControl::builder().build().await?; + cleanup_stale_buckets(&control, &project_id).await; + + let bucket_id = random_bucket_id(); + let create = control + .create_bucket() + .set_parent("projects/_") + .set_bucket_id(bucket_id) + .set_bucket( + Bucket::new() + .set_project(format!("projects/{project_id}")) + .set_location("us-central1") + .set_storage_class("RAPID") + .set_labels([("integration-test", "true")]) + .set_hierarchical_namespace(HierarchicalNamespace::new().set_enabled(true)) + .set_iam_config(IamConfig::new().set_uniform_bucket_level_access( + UniformBucketLevelAccess::new().set_enabled(true), + )), + ) + .with_idempotency(true) + .send() + .await?; + println!("create_test_regional_rapid_bucket(): {create:?}"); + Ok((control, create)) +} + pub async fn objects(builder: StorageBuilder, bucket_name: &str, prefix: &str) -> Result<()> { let client = builder.build().await?; tracing::info!("testing insert_object()"); diff --git a/tests/storage/tests/driver.rs b/tests/storage/tests/driver.rs index b1f99f6de5..e12c3356ad 100644 --- a/tests/storage/tests/driver.rs +++ b/tests/storage/tests/driver.rs @@ -129,6 +129,116 @@ mod storage { result } + #[tokio::test(flavor = "multi_thread")] + async fn run_storage_bidi_conformance_regional_standard() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Standard (HNS)", + ) + .await + .inspect_err(anydump); + let _ = + storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn run_storage_bidi_conformance_regional_rapid() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_regional_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Rapid (HNS)", + ) + .await + .inspect_err(anydump); + let _ = + storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn run_storage_bidi_conformance_zonal_rapid_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Co-located, us-central1-a)", + ) + .await + .inspect_err(anydump); + let _ = + storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn run_storage_bidi_conformance_zonal_rapid_non_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = + integration_tests_storage::build_non_colocated_storage_client("us-central1-b").await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Non Co-located, off-zone endpoint)", + ) + .await + .inspect_err(anydump); + let _ = + storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + async fn run_storage_bidi_features() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::features::run(&client, &bucket.name) + .await + .inspect_err(anydump); + let _ = + storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + #[tokio::test(flavor = "multi_thread")] #[cfg(google_cloud_unstable_storage_bidi)] async fn run_storage_bidi_write() -> anyhow::Result<()> { From d861f0e6cfe5b8f6f7fa5517dc892eccc1f70351 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Tue, 22 Sep 2026 06:09:53 +0000 Subject: [PATCH 4/8] cleaner test code design structure --- .../src/{bidi_read/mod.rs => bidi_read.rs} | 19 -- tests/storage/tests/driver.rs | 227 +++++++++--------- 2 files changed, 115 insertions(+), 131 deletions(-) rename tests/storage/src/{bidi_read/mod.rs => bidi_read.rs} (57%) diff --git a/tests/storage/src/bidi_read/mod.rs b/tests/storage/src/bidi_read.rs similarity index 57% rename from tests/storage/src/bidi_read/mod.rs rename to tests/storage/src/bidi_read.rs index c0f0b0bd17..1372309345 100644 --- a/tests/storage/src/bidi_read/mod.rs +++ b/tests/storage/src/bidi_read.rs @@ -20,22 +20,3 @@ pub mod conformance; pub mod features; - -use google_cloud_storage::client::Storage; - -/// Default runner that executes both conformance and feature test suites. -pub async fn run(bucket_name: &str) -> anyhow::Result<()> { - let mut builder = Storage::builder(); - if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { - builder = builder.with_endpoint(endpoint); - } - let client = builder.build().await?; - - println!("\n=== Running Bidirectional Read Conformance Suite (Suite 1) ==="); - conformance::run(&client, bucket_name).await?; - - println!("\n=== Running Bidirectional Read Features Suite ==="); - features::run(&client, bucket_name).await?; - - Ok(()) -} diff --git a/tests/storage/tests/driver.rs b/tests/storage/tests/driver.rs index e12c3356ad..9bca7e1020 100644 --- a/tests/storage/tests/driver.rs +++ b/tests/storage/tests/driver.rs @@ -112,131 +112,134 @@ mod storage { result } - #[tokio::test(flavor = "multi_thread")] - async fn run_storage_bidi() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_hns_bucket() - .await - .inspect_err(anydump)?; - let result = integration_tests_storage::bidi_read::run(&bucket.name) - .await - .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result - } + mod bidi_read { + use super::*; - #[tokio::test(flavor = "multi_thread")] - async fn run_storage_bidi_conformance_regional_standard() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_hns_bucket() - .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Standard (HNS)", - ) - .await - .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + #[tokio::test(flavor = "multi_thread")] + async fn features() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_hns_bucket() .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result - } - - #[tokio::test(flavor = "multi_thread")] - #[cfg(google_cloud_unstable_storage_bidi)] - async fn run_storage_bidi_conformance_regional_rapid() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_regional_rapid_bucket() - .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Rapid (HNS)", - ) - .await - .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::features::run(&client, &bucket.name) .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) .inspect_err(anydump); - result - } - - #[tokio::test(flavor = "multi_thread")] - #[cfg(google_cloud_unstable_storage_bidi)] - async fn run_storage_bidi_conformance_zonal_rapid_colocated() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Co-located, us-central1-a)", - ) - .await - .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result - } + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } - #[tokio::test(flavor = "multi_thread")] - #[cfg(google_cloud_unstable_storage_bidi)] - async fn run_storage_bidi_conformance_zonal_rapid_non_colocated() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + #[tokio::test(flavor = "multi_thread")] + async fn conformance_regional_standard() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Standard (HNS)", + ) .await - .inspect_err(anydump)?; - let client = - integration_tests_storage::build_non_colocated_storage_client("us-central1-b").await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Non Co-located, off-zone endpoint)", - ) - .await - .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn conformance_regional_rapid() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_regional_rapid_bucket() .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result - } + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Rapid (HNS)", + ) + .await + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } - #[tokio::test(flavor = "multi_thread")] - async fn run_storage_bidi_features() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn conformance_zonal_rapid_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Co-located, us-central1-a)", + ) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::features::run(&client, &bucket.name) + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) .inspect_err(anydump); - let _ = - storage_samples::cleanup_bucket(control, bucket.name.clone(), bucket.project.clone()) + result + } + + #[tokio::test(flavor = "multi_thread")] + #[cfg(google_cloud_unstable_storage_bidi)] + async fn conformance_zonal_rapid_non_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump)?; + let client = + integration_tests_storage::build_non_colocated_storage_client("us-central1-b") + .await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Non Co-located, off-zone endpoint)", + ) + .await + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } } #[tokio::test(flavor = "multi_thread")] From 533f868b1b22d558b0b09eabe1e0e7478eb19032 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Tue, 22 Sep 2026 07:12:34 +0000 Subject: [PATCH 5/8] feat(storage): structure bidi_read tests into regional_standard, zonal_rapid, and regional_rapid submodules --- tests/storage/src/lib.rs | 29 ++-- tests/storage/tests/driver.rs | 247 +++++++++++++++++++++------------- 2 files changed, 171 insertions(+), 105 deletions(-) diff --git a/tests/storage/src/lib.rs b/tests/storage/src/lib.rs index 165db67fe9..d158ea45c4 100644 --- a/tests/storage/src/lib.rs +++ b/tests/storage/src/lib.rs @@ -60,31 +60,34 @@ pub async fn build_non_colocated_storage_client(off_zone: &str) -> Result Result<(StorageControl, Bucket)> { +pub async fn create_test_regional_rapid_bucket(hns: bool) -> Result<(StorageControl, Bucket)> { let project_id = project_id()?; let control = StorageControl::builder().build().await?; cleanup_stale_buckets(&control, &project_id).await; let bucket_id = random_bucket_id(); + let mut bucket = Bucket::new() + .set_project(format!("projects/{project_id}")) + .set_location("us-central1") + .set_storage_class("RAPID") + .set_labels([("integration-test", "true")]) + .set_iam_config( + IamConfig::new() + .set_uniform_bucket_level_access(UniformBucketLevelAccess::new().set_enabled(true)), + ); + if hns { + bucket = bucket.set_hierarchical_namespace(HierarchicalNamespace::new().set_enabled(true)); + } + let create = control .create_bucket() .set_parent("projects/_") .set_bucket_id(bucket_id) - .set_bucket( - Bucket::new() - .set_project(format!("projects/{project_id}")) - .set_location("us-central1") - .set_storage_class("RAPID") - .set_labels([("integration-test", "true")]) - .set_hierarchical_namespace(HierarchicalNamespace::new().set_enabled(true)) - .set_iam_config(IamConfig::new().set_uniform_bucket_level_access( - UniformBucketLevelAccess::new().set_enabled(true), - )), - ) + .set_bucket(bucket) .with_idempotency(true) .send() .await?; - println!("create_test_regional_rapid_bucket(): {create:?}"); + println!("create_test_regional_rapid_bucket(hns={hns}): {create:?}"); Ok((control, create)) } diff --git a/tests/storage/tests/driver.rs b/tests/storage/tests/driver.rs index 9bca7e1020..af9d0a0919 100644 --- a/tests/storage/tests/driver.rs +++ b/tests/storage/tests/driver.rs @@ -136,109 +136,172 @@ mod storage { result } - #[tokio::test(flavor = "multi_thread")] - async fn conformance_regional_standard() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + mod regional_standard { + use super::*; + + #[tokio::test(flavor = "multi_thread")] + async fn hns() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Standard (HNS)", + ) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Standard (HNS)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result - } + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } - #[tokio::test(flavor = "multi_thread")] - #[cfg(google_cloud_unstable_storage_bidi)] - async fn conformance_regional_rapid() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_regional_rapid_bucket() + #[tokio::test(flavor = "multi_thread")] + async fn flat() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Standard (Flat)", + ) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Rapid (HNS)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } } - #[tokio::test(flavor = "multi_thread")] #[cfg(google_cloud_unstable_storage_bidi)] - async fn conformance_zonal_rapid_colocated() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + mod zonal_rapid { + use super::*; + + #[tokio::test(flavor = "multi_thread")] + async fn colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Co-located, us-central1-a)", + ) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Co-located, us-central1-a)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + async fn non_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + .await + .inspect_err(anydump)?; + let client = + integration_tests_storage::build_non_colocated_storage_client("us-central1-b") + .await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Zonal Rapid (Non Co-located, off-zone endpoint)", + ) + .await + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } } - #[tokio::test(flavor = "multi_thread")] #[cfg(google_cloud_unstable_storage_bidi)] - async fn conformance_zonal_rapid_non_colocated() -> anyhow::Result<()> { - let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + mod regional_rapid { + use super::*; + + #[tokio::test(flavor = "multi_thread")] + async fn hns() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = + integration_tests_storage::create_test_regional_rapid_bucket(true) + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Rapid (HNS)", + ) .await - .inspect_err(anydump)?; - let client = - integration_tests_storage::build_non_colocated_storage_client("us-central1-b") - .await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Non Co-located, off-zone endpoint)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } + + #[tokio::test(flavor = "multi_thread")] + async fn flat() -> anyhow::Result<()> { + let _guard = enable_tracing(); + let (control, bucket) = + integration_tests_storage::create_test_regional_rapid_bucket(false) + .await + .inspect_err(anydump)?; + let client = integration_tests_storage::build_storage_client().await?; + let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( + &client, + &bucket.name, + "Regional Rapid (Flat)", + ) + .await + .inspect_err(anydump); + let _ = storage_samples::cleanup_bucket( + control, + bucket.name.clone(), + bucket.project.clone(), + ) + .await + .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(anydump); + result + } } } From d983077efbb55f07847170fc340bbe1e365c46a0 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Tue, 22 Sep 2026 07:31:57 +0000 Subject: [PATCH 6/8] feat(storage): dedicated create and cleanup for rapid regional buckets with full RCU lifecycle --- tests/storage/src/lib.rs | 50 +++++++++++++++++++++++++++++++---- tests/storage/tests/driver.rs | 18 ++++++++++--- 2 files changed, 59 insertions(+), 9 deletions(-) diff --git a/tests/storage/src/lib.rs b/tests/storage/src/lib.rs index d158ea45c4..16bf0958b3 100644 --- a/tests/storage/src/lib.rs +++ b/tests/storage/src/lib.rs @@ -30,9 +30,9 @@ pub use google_cloud_storage::builder::storage::ClientBuilder as StorageBuilder; use google_cloud_storage::builder::storage::SignedUrlBuilder; pub use google_cloud_storage::builder::storage_control::ClientBuilder as StorageControlBuilder; use google_cloud_storage::client::{Storage, StorageControl}; -use google_cloud_storage::model::Bucket; use google_cloud_storage::model::bucket::iam_config::UniformBucketLevelAccess; use google_cloud_storage::model::bucket::{HierarchicalNamespace, IamConfig}; +use google_cloud_storage::model::{Bucket, RapidCache}; use google_cloud_storage::read_object::ReadObjectResponse; use google_cloud_test_utils::resource_names::random_bucket_id; use google_cloud_test_utils::runtime_config::{project_id, test_service_account}; @@ -69,7 +69,6 @@ pub async fn create_test_regional_rapid_bucket(hns: bool) -> Result<(StorageCont let mut bucket = Bucket::new() .set_project(format!("projects/{project_id}")) .set_location("us-central1") - .set_storage_class("RAPID") .set_labels([("integration-test", "true")]) .set_iam_config( IamConfig::new() @@ -79,7 +78,7 @@ pub async fn create_test_regional_rapid_bucket(hns: bool) -> Result<(StorageCont bucket = bucket.set_hierarchical_namespace(HierarchicalNamespace::new().set_enabled(true)); } - let create = control + let created_bucket = control .create_bucket() .set_parent("projects/_") .set_bucket_id(bucket_id) @@ -87,8 +86,49 @@ pub async fn create_test_regional_rapid_bucket(hns: bool) -> Result<(StorageCont .with_idempotency(true) .send() .await?; - println!("create_test_regional_rapid_bucket(hns={hns}): {create:?}"); - Ok((control, create)) + println!( + "create_test_regional_rapid_bucket(hns={hns}) created base bucket: {:?}", + created_bucket.name + ); + + let rapid_cache = RapidCache::new() + .set_name(format!("{}/rapidCaches/us-central1-a", created_bucket.name)) + .set_zone("us-central1-a") + .set_cache_type("rapid-cache-ultra"); + + let _op = control + .create_rapid_cache() + .set_parent(&created_bucket.name) + .set_rapid_cache(rapid_cache) + .poller() + .until_done() + .await?; + println!("create_test_regional_rapid_bucket: attached rapid-cache-ultra in us-central1-a"); + + Ok((control, created_bucket)) +} + +pub async fn cleanup_regional_rapid_bucket( + control: StorageControl, + bucket_name: String, + project_id: String, +) -> Result<()> { + let mut caches = control + .list_rapid_caches() + .set_parent(&bucket_name) + .by_item(); + while let Some(item) = caches.next().await { + if let Ok(cache) = item { + tracing::info!("disabling rapid cache {}", cache.name); + let _ = control + .disable_rapid_cache() + .set_name(cache.name) + .poller() + .until_done() + .await; + } + } + storage_samples::cleanup_bucket(control, bucket_name, project_id).await } pub async fn objects(builder: StorageBuilder, bucket_name: &str, prefix: &str) -> Result<()> { diff --git a/tests/storage/tests/driver.rs b/tests/storage/tests/driver.rs index af9d0a0919..828bbc89d9 100644 --- a/tests/storage/tests/driver.rs +++ b/tests/storage/tests/driver.rs @@ -266,13 +266,18 @@ mod storage { ) .await .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( + let _ = integration_tests_storage::cleanup_regional_rapid_bucket( control, bucket.name.clone(), bucket.project.clone(), ) .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(|e| { + tracing::error!( + "error cleaning up regional rapid bucket {}: {e:?}", + bucket.name + ) + }) .inspect_err(anydump); result } @@ -292,13 +297,18 @@ mod storage { ) .await .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( + let _ = integration_tests_storage::cleanup_regional_rapid_bucket( control, bucket.name.clone(), bucket.project.clone(), ) .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) + .inspect_err(|e| { + tracing::error!( + "error cleaning up regional rapid bucket {}: {e:?}", + bucket.name + ) + }) .inspect_err(anydump); result } From b8124c9901e07174f851f7ef9fc933eceec0d728 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Tue, 22 Sep 2026 10:23:01 +0000 Subject: [PATCH 7/8] feat(storage): support preprod endpoint for RCU StorageControl in regional rapid tests --- tests/storage/src/lib.rs | 28 +++++++++++++++++++++++++++- 1 file changed, 27 insertions(+), 1 deletion(-) diff --git a/tests/storage/src/lib.rs b/tests/storage/src/lib.rs index 16bf0958b3..a78ec068b1 100644 --- a/tests/storage/src/lib.rs +++ b/tests/storage/src/lib.rs @@ -34,6 +34,7 @@ use google_cloud_storage::model::bucket::iam_config::UniformBucketLevelAccess; use google_cloud_storage::model::bucket::{HierarchicalNamespace, IamConfig}; use google_cloud_storage::model::{Bucket, RapidCache}; use google_cloud_storage::read_object::ReadObjectResponse; +use google_cloud_storage::retry_policy::RetryableErrors; use google_cloud_test_utils::resource_names::random_bucket_id; use google_cloud_test_utils::runtime_config::{project_id, test_service_account}; use google_cloud_wkt::FieldMask; @@ -42,6 +43,31 @@ pub use storage_samples::{ cleanup_stale_buckets, create_test_bucket, create_test_hns_bucket, create_test_rapid_bucket, }; +/// Builds a `StorageControl` client for tests. +/// Defaults to the Preprod endpoint (`https://storage-preprod-test-grpc.googleusercontent.com:443`) +/// unless overridden by `GOOGLE_CLOUD_TEST_STORAGE_CONTROL_ENDPOINT`. +pub async fn build_storage_control_client() -> Result { + let endpoint = + std::env::var("GOOGLE_CLOUD_TEST_STORAGE_CONTROL_ENDPOINT").unwrap_or_else(|_| { + "https://storage-preprod-test-grpc.googleusercontent.com:443".to_string() + }); + tracing::info!("StorageControl endpoint: {endpoint}"); + + let client = StorageControl::builder() + .with_endpoint(&endpoint) + .with_backoff_policy( + ExponentialBackoffBuilder::new() + .with_initial_delay(Duration::from_secs(2)) + .with_maximum_delay(Duration::from_secs(8)) + .build()?, + ) + .with_retry_policy(RetryableErrors.with_attempt_limit(5)) + .build() + .await?; + + Ok(client) +} + pub async fn build_storage_client() -> Result { let mut builder = Storage::builder(); if let Ok(endpoint) = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT") { @@ -62,7 +88,7 @@ pub async fn build_non_colocated_storage_client(off_zone: &str) -> Result Result<(StorageControl, Bucket)> { let project_id = project_id()?; - let control = StorageControl::builder().build().await?; + let control = build_storage_control_client().await?; cleanup_stale_buckets(&control, &project_id).await; let bucket_id = random_bucket_id(); From c70b018f930bc088b1ddd9dce8e980953d7dc7d5 Mon Sep 17 00:00:00 2001 From: Olivia Xiaoni Lai <5503815+xlai20@users.noreply.github.com> Date: Wed, 23 Sep 2026 09:15:26 +0000 Subject: [PATCH 8/8] feat(storage): flatten bidi read conformance tests and encapsulate bucket lifecycles --- tests/storage/src/bidi_read/conformance.rs | 115 ++++++++++ tests/storage/src/lib.rs | 13 ++ tests/storage/tests/driver.rs | 247 +++++++++------------ 3 files changed, 230 insertions(+), 145 deletions(-) diff --git a/tests/storage/src/bidi_read/conformance.rs b/tests/storage/src/bidi_read/conformance.rs index 9554b359a9..ccf8e11052 100644 --- a/tests/storage/src/bidi_read/conformance.rs +++ b/tests/storage/src/bidi_read/conformance.rs @@ -42,6 +42,121 @@ pub async fn run_with_scenario( Ok(()) } +// ----------------------------------------------------------------------------- +// Lifecycle Helpers (manage bucket provisioning -> test execution -> teardown) +// ----------------------------------------------------------------------------- + +async fn with_regional_standard_bucket(hns: bool, f: F) -> anyhow::Result<()> +where + F: FnOnce(Storage, String) -> Fut, + Fut: std::future::Future>, +{ + let (control, bucket) = if hns { + crate::create_test_hns_bucket().await? + } else { + crate::create_test_bucket().await? + }; + let client = crate::build_storage_client().await?; + let result = f(client, bucket.name.clone()).await; + let _ = storage_samples::cleanup_bucket(control, bucket.name, bucket.project).await; + result +} + +async fn with_zonal_rapid_bucket(colocated: bool, f: F) -> anyhow::Result<()> +where + F: FnOnce(Storage, String) -> Fut, + Fut: std::future::Future>, +{ + let (control, bucket) = crate::create_test_rapid_bucket().await?; + let client = if colocated { + crate::build_storage_client().await? + } else { + crate::build_non_colocated_storage_client("us-central1-b").await? + }; + let result = f(client, bucket.name.clone()).await; + let _ = storage_samples::cleanup_bucket(control, bucket.name, bucket.project).await; + result +} + +async fn with_regional_rapid_bucket(hns: bool, f: F) -> anyhow::Result<()> +where + F: FnOnce(Storage, String) -> Fut, + Fut: std::future::Future>, +{ + let (control, bucket) = crate::create_test_regional_rapid_bucket(hns).await?; + let client = crate::build_regional_rapid_storage_client().await?; + let result = f(client, bucket.name.clone()).await; + let _ = crate::cleanup_regional_rapid_bucket(control, bucket.name, bucket.project).await; + result +} + +// ----------------------------------------------------------------------------- +// Self-Contained Test Runners (Invoked by driver.rs) +// ----------------------------------------------------------------------------- + +// Non-bucket-type dependent (Tests 2, 4, 5) +pub async fn run_read_post_stream_close() -> anyhow::Result<()> { + with_regional_standard_bucket(false, |client, bucket| async move { + test_read_post_stream_close(&client, &bucket).await + }) + .await +} + +pub async fn run_non_existent_bucket_read() -> anyhow::Result<()> { + let client = crate::build_storage_client().await?; + test_non_existent_bucket_read(&client).await +} + +pub async fn run_out_of_range() -> anyhow::Result<()> { + with_regional_standard_bucket(false, |client, bucket| async move { + test_out_of_range(&client, &bucket).await + }) + .await +} + +// Bucket-type-dependent (Tests 1 & 3) +pub async fn run_multiple_ranged_read_regional_standard(hns: bool) -> anyhow::Result<()> { + with_regional_standard_bucket(hns, |client, bucket| async move { + test_multiple_ranged_read(&client, &bucket).await + }) + .await +} + +pub async fn run_zero_copy_read_regional_standard(hns: bool) -> anyhow::Result<()> { + with_regional_standard_bucket(hns, |client, bucket| async move { + test_zero_copy_read(&client, &bucket).await + }) + .await +} + +pub async fn run_multiple_ranged_read_zonal_rapid(colocated: bool) -> anyhow::Result<()> { + with_zonal_rapid_bucket(colocated, |client, bucket| async move { + test_multiple_ranged_read(&client, &bucket).await + }) + .await +} + +pub async fn run_zero_copy_read_zonal_rapid(colocated: bool) -> anyhow::Result<()> { + with_zonal_rapid_bucket(colocated, |client, bucket| async move { + test_zero_copy_read(&client, &bucket).await + }) + .await +} + +pub async fn run_multiple_ranged_read_regional_rapid(hns: bool) -> anyhow::Result<()> { + with_regional_rapid_bucket(hns, |client, bucket| async move { + test_multiple_ranged_read(&client, &bucket).await + }) + .await +} + +pub async fn run_zero_copy_read_regional_rapid(hns: bool) -> anyhow::Result<()> { + with_regional_rapid_bucket(hns, |client, bucket| async move { + test_zero_copy_read(&client, &bucket).await + }) + .await +} + /// Test Suite 1 - Test 1: Multiple Ranged Read /// /// Tests reading an object across multiple concurrent range read streams over the diff --git a/tests/storage/src/lib.rs b/tests/storage/src/lib.rs index a78ec068b1..e82c6f1f79 100644 --- a/tests/storage/src/lib.rs +++ b/tests/storage/src/lib.rs @@ -86,6 +86,19 @@ pub async fn build_non_colocated_storage_client(off_zone: &str) -> Result Result { + let endpoint = std::env::var("GOOGLE_CLOUD_TEST_STORAGE_CONTROL_ENDPOINT") + .or_else(|_| std::env::var("GOOGLE_CLOUD_TEST_STORAGE_ENDPOINT")) + .unwrap_or_else(|_| { + "https://storage-preprod-test-grpc.googleusercontent.com:443".to_string() + }); + + Ok(Storage::builder().with_endpoint(endpoint).build().await?) +} + pub async fn create_test_regional_rapid_bucket(hns: bool) -> Result<(StorageControl, Bucket)> { let project_id = project_id()?; let control = build_storage_control_client().await?; diff --git a/tests/storage/tests/driver.rs b/tests/storage/tests/driver.rs index 828bbc89d9..8c68010e4e 100644 --- a/tests/storage/tests/driver.rs +++ b/tests/storage/tests/driver.rs @@ -136,181 +136,138 @@ mod storage { result } - mod regional_standard { + mod conformance { use super::*; + use integration_tests_storage::bidi_read::conformance; + + // ========================================================================= + // Non-bucket-type dependent test cases (Tests 2, 4, 5) + // ========================================================================= #[tokio::test(flavor = "multi_thread")] - async fn hns() -> anyhow::Result<()> { + async fn read_post_stream_close() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_hns_bucket() + conformance::run_read_post_stream_close() .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Standard (HNS)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump) } #[tokio::test(flavor = "multi_thread")] - async fn flat() -> anyhow::Result<()> { + async fn non_existent_bucket_read() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_bucket() + conformance::run_non_existent_bucket_read() .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Standard (Flat)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump) } - } - #[cfg(google_cloud_unstable_storage_bidi)] - mod zonal_rapid { - use super::*; + #[tokio::test(flavor = "multi_thread")] + async fn out_of_range() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_out_of_range().await.inspect_err(anydump) + } + + // ========================================================================= + // Bucket-type-dependent test cases (Tests 1 & 3 permuted) + // Format: [test-case-name]_[bucket_type]_[hns]_[colocated] + // ========================================================================= + // --- 1. Regional Standard --- #[tokio::test(flavor = "multi_thread")] - async fn colocated() -> anyhow::Result<()> { + async fn multiple_ranged_read_regional_standard_hns_colocated() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + conformance::run_multiple_ranged_read_regional_standard(true) .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Co-located, us-central1-a)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump) } #[tokio::test(flavor = "multi_thread")] - async fn non_colocated() -> anyhow::Result<()> { + async fn multiple_ranged_read_regional_standard_flat_colocated() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = integration_tests_storage::create_test_rapid_bucket() + conformance::run_multiple_ranged_read_regional_standard(false) .await - .inspect_err(anydump)?; - let client = - integration_tests_storage::build_non_colocated_storage_client("us-central1-b") - .await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Zonal Rapid (Non Co-located, off-zone endpoint)", - ) - .await - .inspect_err(anydump); - let _ = storage_samples::cleanup_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| tracing::error!("error cleaning up bucket {}: {e:?}", bucket.name)) - .inspect_err(anydump); - result + .inspect_err(anydump) } - } - #[cfg(google_cloud_unstable_storage_bidi)] - mod regional_rapid { - use super::*; + #[tokio::test(flavor = "multi_thread")] + async fn zero_copy_read_regional_standard_hns_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_zero_copy_read_regional_standard(true) + .await + .inspect_err(anydump) + } #[tokio::test(flavor = "multi_thread")] - async fn hns() -> anyhow::Result<()> { + async fn zero_copy_read_regional_standard_flat_colocated() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = - integration_tests_storage::create_test_regional_rapid_bucket(true) - .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Rapid (HNS)", - ) - .await - .inspect_err(anydump); - let _ = integration_tests_storage::cleanup_regional_rapid_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| { - tracing::error!( - "error cleaning up regional rapid bucket {}: {e:?}", - bucket.name - ) - }) - .inspect_err(anydump); - result + conformance::run_zero_copy_read_regional_standard(false) + .await + .inspect_err(anydump) } + // --- 2. Zonal Rapid --- #[tokio::test(flavor = "multi_thread")] - async fn flat() -> anyhow::Result<()> { + async fn multiple_ranged_read_zonal_rapid_hns_colocated() -> anyhow::Result<()> { let _guard = enable_tracing(); - let (control, bucket) = - integration_tests_storage::create_test_regional_rapid_bucket(false) - .await - .inspect_err(anydump)?; - let client = integration_tests_storage::build_storage_client().await?; - let result = integration_tests_storage::bidi_read::conformance::run_with_scenario( - &client, - &bucket.name, - "Regional Rapid (Flat)", - ) - .await - .inspect_err(anydump); - let _ = integration_tests_storage::cleanup_regional_rapid_bucket( - control, - bucket.name.clone(), - bucket.project.clone(), - ) - .await - .inspect_err(|e| { - tracing::error!( - "error cleaning up regional rapid bucket {}: {e:?}", - bucket.name - ) - }) - .inspect_err(anydump); - result + conformance::run_multiple_ranged_read_zonal_rapid(true) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn multiple_ranged_read_zonal_rapid_hns_non_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_multiple_ranged_read_zonal_rapid(false) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn zero_copy_read_zonal_rapid_hns_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_zero_copy_read_zonal_rapid(true) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn zero_copy_read_zonal_rapid_hns_non_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_zero_copy_read_zonal_rapid(false) + .await + .inspect_err(anydump) + } + + // --- 3. Regional Rapid (RCU) --- + #[tokio::test(flavor = "multi_thread")] + async fn multiple_ranged_read_regional_rapid_hns_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_multiple_ranged_read_regional_rapid(true) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn multiple_ranged_read_regional_rapid_flat_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_multiple_ranged_read_regional_rapid(false) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn zero_copy_read_regional_rapid_hns_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_zero_copy_read_regional_rapid(true) + .await + .inspect_err(anydump) + } + + #[tokio::test(flavor = "multi_thread")] + async fn zero_copy_read_regional_rapid_flat_colocated() -> anyhow::Result<()> { + let _guard = enable_tracing(); + conformance::run_zero_copy_read_regional_rapid(false) + .await + .inspect_err(anydump) } } }