Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
147 changes: 130 additions & 17 deletions google/cloud/storage/internal/async/connection_tracing.cc
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,21 @@
#include "google/cloud/storage/internal/async/reader_connection_tracing.h"
#include "google/cloud/storage/internal/async/rewriter_connection_tracing.h"
#include "google/cloud/storage/internal/async/writer_connection_tracing.h"
#include "google/cloud/storage/internal/bucket_metadata_cache.h"
#include "google/cloud/storage/options.h"
#include "google/cloud/internal/opentelemetry.h"
#include "google/cloud/version.h"
#include "google/storage/v2/storage.pb.h"
#include <algorithm>
#include <chrono>
#include <future>
#include <memory>
#include <mutex>
#include <vector>

namespace google {
namespace cloud {

namespace storage_internal {
GOOGLE_CLOUD_CPP_INLINE_NAMESPACE_BEGIN

Expand All @@ -33,20 +42,18 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
public:
explicit AsyncConnectionTracing(
std::shared_ptr<storage::AsyncConnection> impl)
: impl_(std::move(impl)) {}
: impl_(std::move(impl)),
cache_(std::make_shared<BucketMetadataCache>()) {}

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

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

future<StatusOr<google::storage::v2::Object>> InsertObject(
InsertObjectParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::InsertObject");
EnrichSpan(*span, p.options,
p.request.write_object_spec().resource().bucket());
internal::OTelScope scope(span);
return internal::EndSpan(std::move(span),
impl_->InsertObject(std::move(p)));
Expand All @@ -55,6 +62,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
future<StatusOr<std::shared_ptr<storage::ObjectDescriptorConnection>>> Open(
OpenParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::Open");
EnrichSpan(*span, p.options, p.read_spec.bucket());
internal::OTelScope scope(span);
return impl_->Open(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
Expand All @@ -63,9 +71,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
std::shared_ptr<storage::ObjectDescriptorConnection>> {
auto result = f.get();
internal::DetachOTelContext(oc);
if (!result) {
if (!result)
return internal::EndSpan(*span, std::move(result).status());
}
return MakeTracingObjectDescriptorConnection(std::move(span),
*std::move(result));
});
Expand All @@ -74,6 +81,7 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
future<StatusOr<std::unique_ptr<storage::AsyncReaderConnection>>> ReadObject(
ReadObjectParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::ReadObject");
EnrichSpan(*span, p.options, p.request.bucket());
internal::OTelScope scope(span);
auto wrap = [oc = opentelemetry::context::RuntimeContext::GetCurrent(),
span = std::move(span)](auto f)
Expand All @@ -89,20 +97,18 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
future<StatusOr<storage::ReadPayload>> ReadObjectRange(
ReadObjectParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::ReadObjectRange");
EnrichSpan(*span, p.options, p.request.bucket());
internal::OTelScope scope(span);
return impl_->ReadObjectRange(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
span = std::move(span)](auto f) {
auto result = f.get();
internal::DetachOTelContext(oc);
return internal::EndSpan(*span, std::move(result));
});
return internal::EndSpan(std::move(span),
impl_->ReadObjectRange(std::move(p)));
}

future<StatusOr<std::unique_ptr<storage::AsyncWriterConnection>>>
StartAppendableObjectUpload(AppendableUploadParams p) override {
auto span = internal::MakeSpan(
"storage::AsyncConnection::StartAppendableObjectUpload");
EnrichSpan(*span, p.options,
p.request.write_object_spec().resource().bucket());
internal::OTelScope scope(span);
return impl_->StartAppendableObjectUpload(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
Expand All @@ -119,6 +125,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
ResumeAppendableObjectUpload(AppendableUploadParams p) override {
auto span = internal::MakeSpan(
"storage::AsyncConnection::ResumeAppendableObjectUpload");
EnrichSpan(*span, p.options,
p.request.write_object_spec().resource().bucket());
internal::OTelScope scope(span);
return impl_->ResumeAppendableObjectUpload(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
Expand All @@ -135,6 +143,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
StartUnbufferedUpload(UploadParams p) override {
auto span =
internal::MakeSpan("storage::AsyncConnection::StartUnbufferedUpload");
EnrichSpan(*span, p.options,
p.request.write_object_spec().resource().bucket());
internal::OTelScope scope(span);
return impl_->StartUnbufferedUpload(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
Expand All @@ -151,6 +161,8 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
StartBufferedUpload(UploadParams p) override {
auto span =
internal::MakeSpan("storage::AsyncConnection::StartBufferedUpload");
EnrichSpan(*span, p.options,
p.request.write_object_spec().resource().bucket());
internal::OTelScope scope(span);
return impl_->StartBufferedUpload(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
Expand Down Expand Up @@ -198,13 +210,15 @@ class AsyncConnectionTracing : public storage::AsyncConnection {
future<StatusOr<google::storage::v2::Object>> ComposeObject(
ComposeObjectParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::ComposeObject");
EnrichSpan(*span, p.options, p.request.destination().bucket());
internal::OTelScope scope(span);
return internal::EndSpan(std::move(span),
impl_->ComposeObject(std::move(p)));
}

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

future<StatusOr<google::storage::v2::Bucket>> GetBucket(
GetBucketParams p) override {
auto span = internal::MakeSpan("storage::AsyncConnection::GetBucket");
internal::OTelScope scope(span);
auto const bucket_name = p.request.name();
auto const options = p.options;
return impl_->GetBucket(std::move(p))
.then([oc = opentelemetry::context::RuntimeContext::GetCurrent(),
span = std::move(span), cache = cache_, bucket_name,
options](future<StatusOr<google::storage::v2::Bucket>> f)
-> StatusOr<google::storage::v2::Bucket> {
StatusOr<google::storage::v2::Bucket> result = f.get();
internal::DetachOTelContext(oc);
if (result.ok()) {
Comment thread
bajajneha27 marked this conversation as resolved.
EnrichSpan(*span, options, *result, bucket_name, *cache);
} else {
cache->MaybeInvalidate(result, bucket_name);
}
return internal::EndSpan(*span, std::move(result));
});
}

private:
static constexpr char kProjectBucketPrefix[] = "projects/_/buckets/";
static constexpr char kGlobalLocation[] = "global";

BucketMetadataCache& cache() const { return *cache_; }

void MaybeTriggerBackgroundFetch(Options const& options,
std::string const& bucket_name) {
if (!cache().StartFetch(bucket_name)) {
return;
}

auto guard = ScopedFetch(cache_, bucket_name);
google::storage::v2::GetBucketRequest request;
auto const normalized_bucket_name =
BucketMetadataCache::NormalizeBucketName(bucket_name);
request.set_name(std::string(kProjectBucketPrefix) +
normalized_bucket_name);
GetBucketParams params{std::move(request), options};

impl_->GetBucket(std::move(params))
.then([cache = cache_, bucket_name, guard = std::move(guard)](
future<StatusOr<google::storage::v2::Bucket>> f) {
StatusOr<google::storage::v2::Bucket> metadata = f.get();
if (metadata.ok()) {
BucketCacheEntry entry = BucketCacheEntry::FromLocation(
metadata->project() + "/buckets/" +
BucketMetadataCache::NormalizeBucketName(bucket_name),
metadata->location(), metadata->location_type());
cache->Put(bucket_name, std::move(entry));
} else if (metadata.status().code() ==
StatusCode::kPermissionDenied) {
BucketCacheEntry entry{
std::string(kProjectBucketPrefix) +
BucketMetadataCache::NormalizeBucketName(bucket_name),
kGlobalLocation};
cache->Put(bucket_name, std::move(entry));
}
});
}

static void EnrichSpan(opentelemetry::trace::Span& span,
BucketCacheEntry const& entry) {
span.SetAttribute("gcp.resource.destination.id", entry.id);
span.SetAttribute("gcp.resource.destination.location", entry.location);
}

static void EnrichSpan(opentelemetry::trace::Span& span,
Options const& options,
google::storage::v2::Bucket const& bucket,
std::string const& bucket_name,
BucketMetadataCache& cache) {
auto const enabled = options.get<
google::cloud::storage_experimental::OTelSpanEnrichmentOption>();
if (!enabled) return;
auto entry = BucketCacheEntry::FromLocation(
bucket.project() + "/buckets/" +
BucketMetadataCache::NormalizeBucketName(bucket_name),
bucket.location(), bucket.location_type());
EnrichSpan(span, entry);
cache.Put(bucket_name, std::move(entry));
}

void EnrichSpan(opentelemetry::trace::Span& span, Options const& options,
std::string const& bucket_name) {
if (bucket_name.empty()) return;
auto const enabled = options.get<
google::cloud::storage_experimental::OTelSpanEnrichmentOption>();
if (!enabled) return;
auto entry = cache().Get(bucket_name);
if (entry.has_value()) {
EnrichSpan(span, *entry);
} else {
MaybeTriggerBackgroundFetch(options, bucket_name);
}
}

std::shared_ptr<storage::AsyncConnection> impl_;
std::shared_ptr<BucketMetadataCache> cache_;
};

} // namespace
Expand Down
Loading
Loading