Skip to content

Commit f8a0cac

Browse files
authored
feat(storage): add resource span attributes for ACO ( App Centric Observability ) for async client (#16151)
1 parent f07c472 commit f8a0cac

7 files changed

Lines changed: 411 additions & 25 deletions

File tree

google/cloud/storage/internal/async/connection_tracing.cc

Lines changed: 130 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,21 @@
1818
#include "google/cloud/storage/internal/async/reader_connection_tracing.h"
1919
#include "google/cloud/storage/internal/async/rewriter_connection_tracing.h"
2020
#include "google/cloud/storage/internal/async/writer_connection_tracing.h"
21+
#include "google/cloud/storage/internal/bucket_metadata_cache.h"
22+
#include "google/cloud/storage/options.h"
2123
#include "google/cloud/internal/opentelemetry.h"
2224
#include "google/cloud/version.h"
25+
#include "google/storage/v2/storage.pb.h"
26+
#include <algorithm>
27+
#include <chrono>
28+
#include <future>
2329
#include <memory>
30+
#include <mutex>
31+
#include <vector>
2432

2533
namespace google {
2634
namespace cloud {
35+
2736
namespace storage_internal {
2837
GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN
2938

@@ -33,20 +42,18 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
3342
public:
3443
explicit AsyncConnectionTracing(
3544
std::shared_ptr<storage::AsyncConnection> impl)
36-
: impl_(std::move(impl)) {}
45+
: impl_(std::move(impl)),
46+
cache_(std::make_shared<BucketMetadataCache>()) {}
3747

38-
Options options() const override { return impl_->options(); }
48+
~AsyncConnectionTracing() override = default;
3949

40-
future<StatusOr<google::storage::v2::Bucket>> GetBucket(
41-
GetBucketParams p) override {
42-
auto span = internal::MakeSpan("storage::AsyncConnection::GetBucket");
43-
internal::OTelScope scope(span);
44-
return internal::EndSpan(std::move(span), impl_->GetBucket(std::move(p)));
45-
}
50+
Options options() const override { return impl_->options(); }
4651

4752
future<StatusOr<google::storage::v2::Object>> InsertObject(
4853
InsertObjectParams p) override {
4954
auto span = internal::MakeSpan("storage::AsyncConnection::InsertObject");
55+
EnrichSpan(*span, p.options,
56+
p.request.write_object_spec().resource().bucket());
5057
internal::OTelScope scope(span);
5158
return internal::EndSpan(std::move(span),
5259
impl_->InsertObject(std::move(p)));
@@ -55,6 +62,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
5562
future<StatusOr<std::shared_ptr<storage::ObjectDescriptorConnection>>> Open(
5663
OpenParams p) override {
5764
auto span = internal::MakeSpan("storage::AsyncConnection::Open");
65+
EnrichSpan(*span, p.options, p.read_spec.bucket());
5866
internal::OTelScope scope(span);
5967
return impl_->Open(std::move(p))
6068
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
@@ -63,9 +71,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
6371
std::shared_ptr<storage::ObjectDescriptorConnection>> {
6472
auto result = f.get();
6573
internal::DetachOTelContext(oc);
66-
if (!result) {
74+
if (!result)
6775
return internal::EndSpan(*span, std::move(result).status());
68-
}
6976
return MakeTracingObjectDescriptorConnection(std::move(span),
7077
*std::move(result));
7178
});
@@ -74,6 +81,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
7481
future<StatusOr<std::unique_ptr<storage::AsyncReaderConnection>>> ReadObject(
7582
ReadObjectParams p) override {
7683
auto span = internal::MakeSpan("storage::AsyncConnection::ReadObject");
84+
EnrichSpan(*span, p.options, p.request.bucket());
7785
internal::OTelScope scope(span);
7886
auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(),
7987
span = std::move(span)](auto f)
@@ -89,20 +97,18 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
8997
future<StatusOr<storage::ReadPayload>> ReadObjectRange(
9098
ReadObjectParams p) override {
9199
auto span = internal::MakeSpan("storage::AsyncConnection::ReadObjectRange");
100+
EnrichSpan(*span, p.options, p.request.bucket());
92101
internal::OTelScope scope(span);
93-
return impl_->ReadObjectRange(std::move(p))
94-
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
95-
span = std::move(span)](auto f) {
96-
auto result = f.get();
97-
internal::DetachOTelContext(oc);
98-
return internal::EndSpan(*span, std::move(result));
99-
});
102+
return internal::EndSpan(std::move(span),
103+
impl_->ReadObjectRange(std::move(p)));
100104
}
101105

102106
future<StatusOr<std::unique_ptr<storage::AsyncWriterConnection>>>
103107
StartAppendableObjectUpload(AppendableUploadParams p) override {
104108
auto span = internal::MakeSpan(
105109
"storage::AsyncConnection::StartAppendableObjectUpload");
110+
EnrichSpan(*span, p.options,
111+
p.request.write_object_spec().resource().bucket());
106112
internal::OTelScope scope(span);
107113
return impl_->StartAppendableObjectUpload(std::move(p))
108114
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
@@ -119,6 +125,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
119125
ResumeAppendableObjectUpload(AppendableUploadParams p) override {
120126
auto span = internal::MakeSpan(
121127
"storage::AsyncConnection::ResumeAppendableObjectUpload");
128+
EnrichSpan(*span, p.options,
129+
p.request.write_object_spec().resource().bucket());
122130
internal::OTelScope scope(span);
123131
return impl_->ResumeAppendableObjectUpload(std::move(p))
124132
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
@@ -135,6 +143,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
135143
StartUnbufferedUpload(UploadParams p) override {
136144
auto span =
137145
internal::MakeSpan("storage::AsyncConnection::StartUnbufferedUpload");
146+
EnrichSpan(*span, p.options,
147+
p.request.write_object_spec().resource().bucket());
138148
internal::OTelScope scope(span);
139149
return impl_->StartUnbufferedUpload(std::move(p))
140150
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
@@ -151,6 +161,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
151161
StartBufferedUpload(UploadParams p) override {
152162
auto span =
153163
internal::MakeSpan("storage::AsyncConnection::StartBufferedUpload");
164+
EnrichSpan(*span, p.options,
165+
p.request.write_object_spec().resource().bucket());
154166
internal::OTelScope scope(span);
155167
return impl_->StartBufferedUpload(std::move(p))
156168
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
@@ -198,13 +210,15 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
198210
future<StatusOr<google::storage::v2::Object>> ComposeObject(
199211
ComposeObjectParams p) override {
200212
auto span = internal::MakeSpan("storage::AsyncConnection::ComposeObject");
213+
EnrichSpan(*span, p.options, p.request.destination().bucket());
201214
internal::OTelScope scope(span);
202215
return internal::EndSpan(std::move(span),
203216
impl_->ComposeObject(std::move(p)));
204217
}
205218

206219
future<Status> DeleteObject(DeleteObjectParams p) override {
207220
auto span = internal::MakeSpan("storage::AsyncConnection::DeleteObject");
221+
EnrichSpan(*span, p.options, p.request.bucket());
208222
internal::OTelScope scope(span);
209223
return internal::EndSpan(std::move(span),
210224
impl_->DeleteObject(std::move(p)));
@@ -217,8 +231,107 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
217231
impl_->RewriteObject(std::move(p)), enabled);
218232
}
219233

234+
future<StatusOr<google::storage::v2::Bucket>> GetBucket(
235+
GetBucketParams p) override {
236+
auto span = internal::MakeSpan("storage::AsyncConnection::GetBucket");
237+
internal::OTelScope scope(span);
238+
auto const bucket_name = p.request.name();
239+
auto const options = p.options;
240+
return impl_->GetBucket(std::move(p))
241+
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
242+
span = std::move(span), cache = cache_, bucket_name,
243+
options](future<StatusOr<google::storage::v2::Bucket>> f)
244+
-> StatusOr<google::storage::v2::Bucket> {
245+
StatusOr<google::storage::v2::Bucket> result = f.get();
246+
internal::DetachOTelContext(oc);
247+
if (result.ok()) {
248+
EnrichSpan(*span, options, *result, bucket_name, *cache);
249+
} else {
250+
cache->MaybeInvalidate(result, bucket_name);
251+
}
252+
return internal::EndSpan(*span, std::move(result));
253+
});
254+
}
255+
220256
private:
257+
static constexpr char kProjectBucketPrefix[] = "projects/_/buckets/";
258+
static constexpr char kGlobalLocation[] = "global";
259+
260+
BucketMetadataCache& cache() const { return *cache_; }
261+
262+
void MaybeTriggerBackgroundFetch(Options const& options,
263+
std::string const& bucket_name) {
264+
if (!cache().StartFetch(bucket_name)) {
265+
return;
266+
}
267+
268+
auto guard = ScopedFetch(cache_, bucket_name);
269+
google::storage::v2::GetBucketRequest request;
270+
auto const normalized_bucket_name =
271+
BucketMetadataCache::NormalizeBucketName(bucket_name);
272+
request.set_name(std::string(kProjectBucketPrefix) +
273+
normalized_bucket_name);
274+
GetBucketParams params{std::move(request), options};
275+
276+
impl_->GetBucket(std::move(params))
277+
.then([cache = cache_, bucket_name, guard = std::move(guard)](
278+
future<StatusOr<google::storage::v2::Bucket>> f) {
279+
StatusOr<google::storage::v2::Bucket> metadata = f.get();
280+
if (metadata.ok()) {
281+
BucketCacheEntry entry = BucketCacheEntry::FromLocation(
282+
metadata->project() + "/buckets/" +
283+
BucketMetadataCache::NormalizeBucketName(bucket_name),
284+
metadata->location(), metadata->location_type());
285+
cache->Put(bucket_name, std::move(entry));
286+
} else if (metadata.status().code() ==
287+
StatusCode::kPermissionDenied) {
288+
BucketCacheEntry entry{
289+
std::string(kProjectBucketPrefix) +
290+
BucketMetadataCache::NormalizeBucketName(bucket_name),
291+
kGlobalLocation};
292+
cache->Put(bucket_name, std::move(entry));
293+
}
294+
});
295+
}
296+
297+
static void EnrichSpan(opentelemetry::trace::Span& span,
298+
BucketCacheEntry const& entry) {
299+
span.SetAttribute("gcp.resource.destination.id", entry.id);
300+
span.SetAttribute("gcp.resource.destination.location", entry.location);
301+
}
302+
303+
static void EnrichSpan(opentelemetry::trace::Span& span,
304+
Options const& options,
305+
google::storage::v2::Bucket const& bucket,
306+
std::string const& bucket_name,
307+
BucketMetadataCache& cache) {
308+
auto const enabled = options.get<
309+
google::cloud::storage_experimental::OTelSpanEnrichmentOption>();
310+
if (!enabled) return;
311+
auto entry = BucketCacheEntry::FromLocation(
312+
bucket.project() + "/buckets/" +
313+
BucketMetadataCache::NormalizeBucketName(bucket_name),
314+
bucket.location(), bucket.location_type());
315+
EnrichSpan(span, entry);
316+
cache.Put(bucket_name, std::move(entry));
317+
}
318+
319+
void EnrichSpan(opentelemetry::trace::Span& span, Options const& options,
320+
std::string const& bucket_name) {
321+
if (bucket_name.empty()) return;
322+
auto const enabled = options.get<
323+
google::cloud::storage_experimental::OTelSpanEnrichmentOption>();
324+
if (!enabled) return;
325+
auto entry = cache().Get(bucket_name);
326+
if (entry.has_value()) {
327+
EnrichSpan(span, *entry);
328+
} else {
329+
MaybeTriggerBackgroundFetch(options, bucket_name);
330+
}
331+
}
332+
221333
std::shared_ptr<storage::AsyncConnection> impl_;
334+
std::shared_ptr<BucketMetadataCache> cache_;
222335
};
223336

224337
} // namespace

0 commit comments

Comments
 (0)