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,11 +44,19 @@ 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::Object>> InsertObject (
4156 InsertObjectParams p) override {
4257 auto span = internal::MakeSpan (" storage::AsyncConnection::InsertObject" );
58+ EnrichSpan (*span, p.options ,
59+ p.request .write_object_spec ().resource ().bucket ());
4360 internal::OTelScope scope (span);
4461 return internal::EndSpan (std::move (span),
4562 impl_->InsertObject (std::move (p)));
@@ -48,6 +65,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
4865 future<StatusOr<std::shared_ptr<storage::ObjectDescriptorConnection>>> Open (
4966 OpenParams p) override {
5067 auto span = internal::MakeSpan (" storage::AsyncConnection::Open" );
68+ EnrichSpan (*span, p.options , p.read_spec .bucket ());
5169 internal::OTelScope scope (span);
5270 return impl_->Open (std::move (p))
5371 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -67,6 +85,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
6785 future<StatusOr<std::unique_ptr<storage::AsyncReaderConnection>>> ReadObject (
6886 ReadObjectParams p) override {
6987 auto span = internal::MakeSpan (" storage::AsyncConnection::ReadObject" );
88+ EnrichSpan (*span, p.options , p.request .bucket ());
7089 internal::OTelScope scope (span);
7190 auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent (),
7291 span = std::move (span)](auto f)
@@ -82,6 +101,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
82101 future<StatusOr<storage::ReadPayload>> ReadObjectRange (
83102 ReadObjectParams p) override {
84103 auto span = internal::MakeSpan (" storage::AsyncConnection::ReadObjectRange" );
104+ EnrichSpan (*span, p.options , p.request .bucket ());
85105 internal::OTelScope scope (span);
86106 return impl_->ReadObjectRange (std::move (p))
87107 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -96,6 +116,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
96116 StartAppendableObjectUpload (AppendableUploadParams p) override {
97117 auto span = internal::MakeSpan (
98118 " storage::AsyncConnection::StartAppendableObjectUpload" );
119+ EnrichSpan (*span, p.options ,
120+ p.request .write_object_spec ().resource ().bucket ());
99121 internal::OTelScope scope (span);
100122 return impl_->StartAppendableObjectUpload (std::move (p))
101123 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -112,6 +134,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
112134 ResumeAppendableObjectUpload (AppendableUploadParams p) override {
113135 auto span = internal::MakeSpan (
114136 " storage::AsyncConnection::ResumeAppendableObjectUpload" );
137+ EnrichSpan (*span, p.options ,
138+ p.request .write_object_spec ().resource ().bucket ());
115139 internal::OTelScope scope (span);
116140 return impl_->ResumeAppendableObjectUpload (std::move (p))
117141 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -128,6 +152,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
128152 StartUnbufferedUpload (UploadParams p) override {
129153 auto span =
130154 internal::MakeSpan (" storage::AsyncConnection::StartUnbufferedUpload" );
155+ EnrichSpan (*span, p.options ,
156+ p.request .write_object_spec ().resource ().bucket ());
131157 internal::OTelScope scope (span);
132158 return impl_->StartUnbufferedUpload (std::move (p))
133159 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -144,6 +170,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
144170 StartBufferedUpload (UploadParams p) override {
145171 auto span =
146172 internal::MakeSpan (" storage::AsyncConnection::StartBufferedUpload" );
173+ EnrichSpan (*span, p.options ,
174+ p.request .write_object_spec ().resource ().bucket ());
147175 internal::OTelScope scope (span);
148176 return impl_->StartBufferedUpload (std::move (p))
149177 .then ([oc = opentelemetry::context::RuntimeContext::GetCurrent (),
@@ -191,13 +219,15 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
191219 future<StatusOr<google::storage::v2::Object>> ComposeObject (
192220 ComposeObjectParams p) override {
193221 auto span = internal::MakeSpan (" storage::AsyncConnection::ComposeObject" );
222+ EnrichSpan (*span, p.options , p.request .destination ().bucket ());
194223 internal::OTelScope scope (span);
195224 return internal::EndSpan (std::move (span),
196225 impl_->ComposeObject (std::move (p)));
197226 }
198227
199228 future<Status> DeleteObject (DeleteObjectParams p) override {
200229 auto span = internal::MakeSpan (" storage::AsyncConnection::DeleteObject" );
230+ EnrichSpan (*span, p.options , p.request .bucket ());
201231 internal::OTelScope scope (span);
202232 return internal::EndSpan (std::move (span),
203233 impl_->DeleteObject (std::move (p)));
@@ -211,7 +241,61 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
211241 }
212242
213243 private:
244+ void CleanupCompletedTasks () {
245+ std::unique_lock<std::mutex> lk (mu_);
246+ bg_tasks_.erase (
247+ std::remove_if (bg_tasks_.begin (), bg_tasks_.end (),
248+ [](std::future<void > const & f) {
249+ return f.wait_for (std::chrono::seconds (0 )) ==
250+ std::future_status::ready;
251+ }),
252+ bg_tasks_.end ());
253+ }
254+
255+ void MaybeTriggerBackgroundFetch (Options const & options,
256+ std::string const & bucket_name) {
257+ CleanupCompletedTasks ();
258+
259+ if (!BucketMetadataCache::Singleton ().StartFetch (bucket_name)) {
260+ return ;
261+ }
262+
263+ auto f = std::async (std::launch::async, [bucket_name, options]() {
264+ google::cloud::internal::OptionsSpan span (options);
265+ auto conn = MakeStorageConnection (options);
266+ auto const normalized =
267+ BucketMetadataCache::NormalizeBucketName (bucket_name);
268+ storage::internal::GetBucketMetadataRequest request (normalized);
269+ auto metadata = conn->GetBucketMetadata (request);
270+ if (metadata.ok ()) {
271+ BucketMetadataCache::Singleton ().Put (
272+ bucket_name, BucketCacheEntry::FromMetadata (*metadata));
273+ } else if (metadata.status ().code () == StatusCode::kPermissionDenied ) {
274+ BucketMetadataCache::Singleton ().Put (
275+ bucket_name, {" projects/_/buckets/" + normalized, " global" });
276+ }
277+ BucketMetadataCache::Singleton ().EndFetch (bucket_name);
278+ });
279+
280+ std::unique_lock<std::mutex> lk (mu_);
281+ bg_tasks_.push_back (std::move (f));
282+ }
283+
284+ void EnrichSpan (opentelemetry::trace::Span& span, Options const & options,
285+ std::string const & bucket_name) {
286+ if (bucket_name.empty ()) return ;
287+ auto entry = BucketMetadataCache::Singleton ().Get (bucket_name);
288+ if (entry.has_value ()) {
289+ span.SetAttribute (" gcp.resource.destination.id" , entry->id );
290+ span.SetAttribute (" gcp.resource.destination.location" , entry->location );
291+ } else {
292+ MaybeTriggerBackgroundFetch (options, bucket_name);
293+ }
294+ }
295+
214296 std::shared_ptr<storage::AsyncConnection> impl_;
297+ std::vector<std::future<void >> bg_tasks_;
298+ std::mutex mu_;
215299};
216300
217301} // namespace
0 commit comments