From afc726143044c0796dbbda23d9a0b0312c3f10b1 Mon Sep 17 00:00:00 2001 From: NickGoog Date: Thu, 24 Sep 2026 14:22:47 +0000 Subject: [PATCH 1/4] docs(storage): add zonal bucket pre-warmed writer pool sample Adds OptimizeWriteLatencyPool sample (region tag: storage_optimize_write_latency_pool) to storage_async_samples.cc demonstrating a pre-warmed pool of AsyncWriter instances with unfinalized objects and Flush() to avoid object creation and finalization metadata overhead on the critical write path. Verified with both mock unit tests and live integration testing against a Rapid (zonal) bucket in us-central1-a: Running live C++ test against bucket=, prefix=live_cpp_pool_1790087459 C++ read back: "0123456789", pool size after refill: 3 [ OK ] WriterPoolCppTest.PreWarmedPoolWithUnfinalizedObjects (1464 ms) [ PASSED ] 1 test. --- .../storage/examples/storage_async_samples.cc | 74 +++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index 97a9bc454adb7..dd1d253affb7a 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -26,6 +26,7 @@ #include #include #include +#include #include #include #include @@ -903,6 +904,71 @@ void FinalizeAppendableObjectUpload(google::cloud::storage::AsyncClient& client, std::cout << "Finalized object: " << object.DebugString() << "\n"; } +void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, + std::vector const& argv) { + //! [optimize-write-latency-pool] + // [START storage_optimize_write_latency_pool] + namespace gcs = google::cloud::storage; + auto coro = [](gcs::AsyncClient& client, std::string bucket_name, + std::string key_prefix, + int pool_size) -> google::cloud::future { + std::string const next_object_name = + key_prefix + "_" + std::to_string(pool_size); + + // 1. Init pool: Sized to ensure pre-warmed writers are always available. + std::deque> pool; + for (int i = 0; i < pool_size; ++i) { + auto [writer, token] = (co_await client.StartAppendableObjectUpload( + gcs::BucketName(bucket_name), + key_prefix + "_" + std::to_string(i))) + .value(); + pool.emplace_back(std::move(writer), std::move(token)); + } + + // 2. Write: Pop a pre-warmed writer and commit with Flush() (~1-2 ms) + // instead of Finalize(). + auto [writer, token] = std::move(pool.front()); + pool.pop_front(); + token = (co_await writer.Write(std::move(token), + gcs::WritePayload("0123456789"))) + .value(); + (void)co_await writer.Flush(); + + // 3. Pool maintenance (run asynchronously off the critical write path): + // Close the used writer without finalizing and refill the pool. + auto maintain_pool = [](gcs::AsyncClient client, std::string bucket_name, + std::string next_object_name, + gcs::AsyncWriter writer) + -> google::cloud::future> { + (void)co_await writer.Close(); + auto [new_writer, new_token] = + (co_await client.StartAppendableObjectUpload( + gcs::BucketName(bucket_name), next_object_name)) + .value(); + co_return {std::move(new_writer), std::move(new_token)}; + }; + auto maintenance_future = + maintain_pool(client, bucket_name, next_object_name, std::move(writer)); + + // 4. Read: Unfinalized objects are readable after Flush(). + gcs::ObjectDescriptor descriptor = + (co_await client.Open(gcs::BucketName(bucket_name), key_prefix + "_0")) + .value(); + auto [reader, read_token] = descriptor.Read(0, 10); + auto [payload, next_token] = + (co_await reader.Read(std::move(read_token))).value(); + + auto [new_writer, new_token] = co_await std::move(maintenance_future); + pool.emplace_back(std::move(new_writer), std::move(new_token)); + for (auto& [rem_writer, rem_token] : pool) { + (void)co_await rem_writer.Close(); + } + }; + // [END storage_optimize_write_latency_pool] + //! [optimize-write-latency-pool] + coro(client, argv.at(0), argv.at(1), 3).get(); +} + void ReadAppendableObjectTail(google::cloud::storage::AsyncClient& client, std::vector const& argv) { //! [read-appendable-object-tail] @@ -1492,6 +1558,13 @@ void AutoRun(std::vector const& argv) { scheduled_for_delete.push_back(std::move(object_name)); object_name = examples::MakeRandomObjectName(generator, "object-"); + std::cout << "Running OptimizeWriteLatencyPool() example" << std::endl; + OptimizeWriteLatencyPool(client, {bucket_name, object_name}); + for (int i = 0; i != 4; ++i) { + scheduled_for_delete.push_back(object_name + "_" + std::to_string(i)); + } + object_name = examples::MakeRandomObjectName(generator, "object-"); + std::cout << "Running ReadAppendableObjectTail() example" << std::endl; // Create a dummy object for the tail example to read. In a real // application another process would be writing to this object. @@ -1657,6 +1730,7 @@ int main(int argc, char* argv[]) try { PauseAndResumeAppendableUpload), make_entry("finalize-appendable-object-upload", {}, FinalizeAppendableObjectUpload), + make_entry("optimize-write-latency-pool", {}, OptimizeWriteLatencyPool), make_entry("rewrite-object", {""}, RewriteObject), make_entry("resume-rewrite-object", {""}, ResumeRewrite), From 4e9da63ea2bf0483f7947728ce38a0f37db61f40 Mon Sep 17 00:00:00 2001 From: NickGoog Date: Wed, 30 Sep 2026 15:49:07 +0000 Subject: [PATCH 2/4] fix(storage): address review feedback in writer pool sample --- .../storage/examples/storage_async_samples.cc | 31 ++++++++++++++----- 1 file changed, 24 insertions(+), 7 deletions(-) diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index dd1d253affb7a..ca496eb07c332 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -925,14 +925,15 @@ void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, pool.emplace_back(std::move(writer), std::move(token)); } - // 2. Write: Pop a pre-warmed writer and commit with Flush() (~1-2 ms) - // instead of Finalize(). + // 2. Write: Pop a pre-warmed writer and commit with Flush() instead of + // Finalize(). auto [writer, token] = std::move(pool.front()); pool.pop_front(); token = (co_await writer.Write(std::move(token), gcs::WritePayload("0123456789"))) .value(); - (void)co_await writer.Flush(); + auto flush_status = co_await writer.Flush(); + if (!flush_status.ok()) throw std::runtime_error(flush_status.message()); // 3. Pool maintenance (run asynchronously off the critical write path): // Close the used writer without finalizing and refill the pool. @@ -940,7 +941,8 @@ void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, std::string next_object_name, gcs::AsyncWriter writer) -> google::cloud::future> { - (void)co_await writer.Close(); + auto close_status = co_await writer.Close(); + if (!close_status.ok()) throw std::runtime_error(close_status.message()); auto [new_writer, new_token] = (co_await client.StartAppendableObjectUpload( gcs::BucketName(bucket_name), next_object_name)) @@ -955,13 +957,22 @@ void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, (co_await client.Open(gcs::BucketName(bucket_name), key_prefix + "_0")) .value(); auto [reader, read_token] = descriptor.Read(0, 10); - auto [payload, next_token] = - (co_await reader.Read(std::move(read_token))).value(); + std::string contents; + while (read_token.valid()) { + auto [payload, t] = (co_await reader.Read(std::move(read_token))).value(); + read_token = std::move(t); + for (auto const& buffer : payload.contents()) { + contents.append(buffer.begin(), buffer.end()); + } + } + std::cout << "Read unfinalized object " << key_prefix << "_0: " << contents + << "\n"; auto [new_writer, new_token] = co_await std::move(maintenance_future); pool.emplace_back(std::move(new_writer), std::move(new_token)); for (auto& [rem_writer, rem_token] : pool) { - (void)co_await rem_writer.Close(); + auto close_status = co_await rem_writer.Close(); + if (!close_status.ok()) throw std::runtime_error(close_status.message()); } }; // [END storage_optimize_write_latency_pool] @@ -1211,6 +1222,12 @@ void FinalizeAppendableObjectUpload(google::cloud::storage::AsyncClient&, "coroutines\n"; } +void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient&, + std::vector const&) { + std::cerr << "AsyncClient::OptimizeWriteLatencyPool() example requires " + "coroutines\n"; +} + void ReadAppendableObjectTail(google::cloud::storage::AsyncClient&, std::vector const&) { std::cerr << "AsyncClient::ReadAppendableObjectTail() example requires " From 4a7e02eae409f33865f890544f6698851bc22be7 Mon Sep 17 00:00:00 2001 From: NickGoog Date: Wed, 30 Sep 2026 16:52:45 +0000 Subject: [PATCH 3/4] docs(storage): describe Flush() as faster instead of citing latency --- google/cloud/storage/examples/storage_async_samples.cc | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index ca496eb07c332..65d613fb24335 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -925,8 +925,8 @@ void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, pool.emplace_back(std::move(writer), std::move(token)); } - // 2. Write: Pop a pre-warmed writer and commit with Flush() instead of - // Finalize(). + // 2. Write: Pop a pre-warmed writer and commit with the faster Flush() + // instead of Finalize(). auto [writer, token] = std::move(pool.front()); pool.pop_front(); token = (co_await writer.Write(std::move(token), From 7c925c78ecfb553ba0c4202642150b577228f896 Mon Sep 17 00:00:00 2001 From: NickGoog Date: Thu, 1 Oct 2026 15:29:44 +0000 Subject: [PATCH 4/4] docs(storage): use make_bucket_entry and move args in writer pool sample --- google/cloud/storage/examples/storage_async_samples.cc | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/google/cloud/storage/examples/storage_async_samples.cc b/google/cloud/storage/examples/storage_async_samples.cc index 494df30154041..cdc252b5049b6 100644 --- a/google/cloud/storage/examples/storage_async_samples.cc +++ b/google/cloud/storage/examples/storage_async_samples.cc @@ -945,7 +945,8 @@ void OptimizeWriteLatencyPool(google::cloud::storage::AsyncClient& client, if (!close_status.ok()) throw std::runtime_error(close_status.message()); auto [new_writer, new_token] = (co_await client.StartAppendableObjectUpload( - gcs::BucketName(bucket_name), next_object_name)) + gcs::BucketName(std::move(bucket_name)), + std::move(next_object_name))) .value(); co_return {std::move(new_writer), std::move(new_token)}; }; @@ -1747,7 +1748,8 @@ int main(int argc, char* argv[]) try { PauseAndResumeAppendableUpload), make_entry("finalize-appendable-object-upload", {}, FinalizeAppendableObjectUpload), - make_entry("optimize-write-latency-pool", {}, OptimizeWriteLatencyPool), + make_bucket_entry("optimize-write-latency-pool", {""}, + OptimizeWriteLatencyPool), make_entry("rewrite-object", {""}, RewriteObject), make_entry("resume-rewrite-object", {""}, ResumeRewrite),