1414
1515#include " google/cloud/storage/internal/async/connection_tracing.h"
1616#include " google/cloud/storage/async/writer_connection.h"
17+ #include " google/cloud/storage/internal/async/default_options.h"
1718#include " google/cloud/storage/internal/async/object_descriptor_connection_tracing.h"
1819#include " google/cloud/storage/internal/async/reader_connection_tracing.h"
1920#include " google/cloud/storage/internal/async/rewriter_connection_tracing.h"
2021#include " google/cloud/storage/internal/async/writer_connection_tracing.h"
22+ #include " google/cloud/storage/internal/bucket_metadata_cache.h"
23+ #include " google/cloud/storage/internal/connection_factory.h"
2124#include " google/cloud/internal/opentelemetry.h"
2225#include " google/cloud/version.h"
26+ #include < algorithm>
27+ #include < chrono>
28+ #include < future>
2329#include < memory>
30+ #include < mutex>
31+ #include < vector>
2432
2533namespace google {
2634namespace cloud {
35+
2736namespace storage_internal {
2837GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN
2938
@@ -35,6 +44,12 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
3544 std::shared_ptr<storage::AsyncConnection> impl)
3645 : impl_(std::move(impl)) {}
3746
47+ ~AsyncConnectionTracing () override {
48+ for (auto & f : bg_tasks_) {
49+ if (f.valid ()) f.wait ();
50+ }
51+ }
52+
3853 Options options () const override { return impl_->options (); }
3954
4055 future<StatusOr<google::storage::v2::Bucket>> GetBucket (
@@ -47,6 +62,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
4762 future<StatusOr<google::storage::v2::Object>> InsertObject (
4863 InsertObjectParams p) override {
4964 auto span = internal::MakeSpan (" storage::AsyncConnection::InsertObject" );
65+ EnrichSpan (*span, p.options ,
66+ p.request .write_object_spec ().resource ().bucket ());
5067 internal::OTelScope scope (span);
5168 return internal::EndSpan (std::move (span),
5269 impl_->InsertObject (std::move (p)));
@@ -55,6 +72,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
5572 future<StatusOr<std::shared_ptr<storage::ObjectDescriptorConnection>>> Open (
5673 OpenParams p) override {
5774 auto span = internal::MakeSpan (" storage::AsyncConnection::Open" );
75+ EnrichSpan (*span, p.options , p.read_spec .bucket ());
5876 internal::OTelScope scope (span);
5977 return impl_->Open (std::move (p))
6078 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -74,6 +92,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
7492 future<StatusOr<std::unique_ptr<storage::AsyncReaderConnection>>> ReadObject (
7593 ReadObjectParams p) override {
7694 auto span = internal::MakeSpan (" storage::AsyncConnection::ReadObject" );
95+ EnrichSpan (*span, p.options , p.request .bucket ());
7796 internal::OTelScope scope (span);
7897 auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent (),
7998 span = std::move (span)](auto f)
@@ -89,6 +108,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
89108 future<StatusOr<storage::ReadPayload>> ReadObjectRange (
90109 ReadObjectParams p) override {
91110 auto span = internal::MakeSpan (" storage::AsyncConnection::ReadObjectRange" );
111+ EnrichSpan (*span, p.options , p.request .bucket ());
92112 internal::OTelScope scope (span);
93113 return impl_->ReadObjectRange (std::move (p))
94114 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -103,6 +123,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
103123 StartAppendableObjectUpload (AppendableUploadParams p) override {
104124 auto span = internal::MakeSpan (
105125 " storage::AsyncConnection::StartAppendableObjectUpload" );
126+ EnrichSpan (*span, p.options ,
127+ p.request .write_object_spec ().resource ().bucket ());
106128 internal::OTelScope scope (span);
107129 return impl_->StartAppendableObjectUpload (std::move (p))
108130 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -119,6 +141,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
119141 ResumeAppendableObjectUpload (AppendableUploadParams p) override {
120142 auto span = internal::MakeSpan (
121143 " storage::AsyncConnection::ResumeAppendableObjectUpload" );
144+ EnrichSpan (*span, p.options ,
145+ p.request .write_object_spec ().resource ().bucket ());
122146 internal::OTelScope scope (span);
123147 return impl_->ResumeAppendableObjectUpload (std::move (p))
124148 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -135,6 +159,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
135159 StartUnbufferedUpload (UploadParams p) override {
136160 auto span =
137161 internal::MakeSpan (" storage::AsyncConnection::StartUnbufferedUpload" );
162+ EnrichSpan (*span, p.options ,
163+ p.request .write_object_spec ().resource ().bucket ());
138164 internal::OTelScope scope (span);
139165 return impl_->StartUnbufferedUpload (std::move (p))
140166 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -151,6 +177,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
151177 StartBufferedUpload (UploadParams p) override {
152178 auto span =
153179 internal::MakeSpan (" storage::AsyncConnection::StartBufferedUpload" );
180+ EnrichSpan (*span, p.options ,
181+ p.request .write_object_spec ().resource ().bucket ());
154182 internal::OTelScope scope (span);
155183 return impl_->StartBufferedUpload (std::move (p))
156184 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -198,13 +226,15 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
198226 future<StatusOr<google::storage::v2::Object>> ComposeObject (
199227 ComposeObjectParams p) override {
200228 auto span = internal::MakeSpan (" storage::AsyncConnection::ComposeObject" );
229+ EnrichSpan (*span, p.options , p.request .destination ().bucket ());
201230 internal::OTelScope scope (span);
202231 return internal::EndSpan (std::move (span),
203232 impl_->ComposeObject (std::move (p)));
204233 }
205234
206235 future<Status> DeleteObject (DeleteObjectParams p) override {
207236 auto span = internal::MakeSpan (" storage::AsyncConnection::DeleteObject" );
237+ EnrichSpan (*span, p.options , p.request .bucket ());
208238 internal::OTelScope scope (span);
209239 return internal::EndSpan (std::move (span),
210240 impl_->DeleteObject (std::move (p)));
@@ -218,7 +248,61 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
218248 }
219249
220250 private:
251+ void CleanupCompletedTasks () {
252+ std::unique_lock<std::mutex> lk (mu_);
253+ bg_tasks_.erase (
254+ std::remove_if (bg_tasks_.begin (), bg_tasks_.end (),
255+ [](std::future<void > const & f) {
256+ return f.wait_for (std::chrono::seconds (0 )) ==
257+ std::future_status::ready;
258+ }),
259+ bg_tasks_.end ());
260+ }
261+
262+ void MaybeTriggerBackgroundFetch (Options const & options,
263+ std::string const & bucket_name) {
264+ CleanupCompletedTasks ();
265+
266+ if (!BucketMetadataCache::Singleton ().StartFetch (bucket_name)) {
267+ return ;
268+ }
269+
270+ auto f = std::async (std::launch::async, [bucket_name, options]() {
271+ google::cloud::internal::OptionsSpan span (options);
272+ auto conn = MakeStorageConnection (options);
273+ auto const normalized =
274+ BucketMetadataCache::NormalizeBucketName (bucket_name);
275+ storage::internal::GetBucketMetadataRequest request (normalized);
276+ auto metadata = conn->GetBucketMetadata (request);
277+ if (metadata.ok ()) {
278+ BucketMetadataCache::Singleton ().Put (
279+ bucket_name, BucketCacheEntry::FromMetadata (*metadata));
280+ } else if (metadata.status ().code () == StatusCode::kPermissionDenied ) {
281+ BucketMetadataCache::Singleton ().Put (
282+ bucket_name, {" projects/_/buckets/" + normalized, " global" });
283+ }
284+ BucketMetadataCache::Singleton ().EndFetch (bucket_name);
285+ });
286+
287+ std::unique_lock<std::mutex> lk (mu_);
288+ bg_tasks_.push_back (std::move (f));
289+ }
290+
291+ void EnrichSpan (opentelemetry::trace::Span& span, Options const & options,
292+ std::string const & bucket_name) {
293+ if (bucket_name.empty ()) return ;
294+ auto entry = BucketMetadataCache::Singleton ().Get (bucket_name);
295+ if (entry.has_value ()) {
296+ span.SetAttribute (" gcp.resource.destination.id" , entry->id );
297+ span.SetAttribute (" gcp.resource.destination.location" , entry->location );
298+ } else {
299+ MaybeTriggerBackgroundFetch (options, bucket_name);
300+ }
301+ }
302+
221303 std::shared_ptr<storage::AsyncConnection> impl_;
304+ std::vector<std::future<void >> bg_tasks_;
305+ std::mutex mu_;
222306};
223307
224308} // namespace
0 commit comments