From 77ad5a786e6eb2d2c5303ff4e5aafe8189013789 Mon Sep 17 00:00:00 2001 From: Yijie Ma Date: Thu, 11 Jan 2024 13:35:10 -0800 Subject: [PATCH] [CSM O11Y] CSM Service Label Plumbing from LB Policies to CallAttemptTracer (#35210) Closes #35210 COPYBARA_INTEGRATE_REVIEW=https://github.com/grpc/grpc/pull/35210 from yijiem:csm-service-label 6a6a7d177478a31248be8f4e551c7c991157c4be PiperOrigin-RevId: 597641393 --- CMakeLists.txt | 2 + build_autogenerated.yaml | 6 +- .../filters/client_channel/client_channel.cc | 7 + .../client_channel/client_channel_internal.h | 2 + .../lb_policy/xds/xds_cluster_impl.cc | 10 + src/core/ext/xds/xds_cluster.cc | 32 +++ src/core/ext/xds/xds_cluster.h | 2 + src/core/lib/channel/call_tracer.cc | 7 + src/core/lib/channel/call_tracer.h | 10 + src/cpp/ext/csm/metadata_exchange.cc | 42 ++++ src/cpp/ext/csm/metadata_exchange.h | 14 ++ .../filters/census/open_census_call_tracer.h | 3 + src/cpp/ext/otel/key_value_iterable.h | 40 +++- src/cpp/ext/otel/otel_call_tracer.h | 7 + src/cpp/ext/otel/otel_client_filter.cc | 20 +- src/cpp/ext/otel/otel_plugin.h | 26 ++- src/cpp/ext/otel/otel_server_call_tracer.cc | 13 +- src/proto/grpc/testing/xds/v3/base.proto | 40 ++++ src/proto/grpc/testing/xds/v3/cluster.proto | 7 + .../grpc_observability/client_call_tracer.h | 4 + test/core/channel/BUILD | 1 + test/core/channel/call_tracer_test.cc | 102 +--------- .../lb_policy/lb_policy_test_lib.h | 4 + test/core/end2end/tests/http2_stats.cc | 5 + test/core/util/BUILD | 10 + test/core/util/fake_stats_plugin.cc | 82 ++++++++ test/core/util/fake_stats_plugin.h | 192 ++++++++++++++++++ .../xds/xds_cluster_resource_type_test.cc | 90 ++++++++ test/cpp/end2end/xds/BUILD | 1 + .../end2end/xds/xds_cluster_end2end_test.cc | 39 ++++ test/cpp/ext/csm/metadata_exchange_test.cc | 39 +++- test/cpp/ext/otel/otel_plugin_test.cc | 28 +++ test/cpp/ext/otel/otel_test_library.cc | 59 +++++- test/cpp/ext/otel/otel_test_library.h | 2 + 34 files changed, 827 insertions(+), 121 deletions(-) create mode 100644 test/core/util/fake_stats_plugin.cc create mode 100644 test/core/util/fake_stats_plugin.h diff --git a/CMakeLists.txt b/CMakeLists.txt index eb2eac86917..92234187bb7 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -7809,6 +7809,7 @@ if(gRPC_BUILD_TESTS) add_executable(call_tracer_test test/core/channel/call_tracer_test.cc + test/core/util/fake_stats_plugin.cc ) target_compile_features(call_tracer_test PUBLIC cxx_std_14) target_include_directories(call_tracer_test @@ -27043,6 +27044,7 @@ if(_gRPC_PLATFORM_LINUX OR _gRPC_PLATFORM_MAC OR _gRPC_PLATFORM_POSIX) ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/xds/v3/string.grpc.pb.cc ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/xds/v3/string.pb.h ${_gRPC_PROTO_GENS_DIR}/src/proto/grpc/testing/xds/v3/string.grpc.pb.h + test/core/util/fake_stats_plugin.cc test/cpp/end2end/connection_attempt_injector.cc test/cpp/end2end/test_service_impl.cc test/cpp/end2end/xds/xds_cluster_end2end_test.cc diff --git a/build_autogenerated.yaml b/build_autogenerated.yaml index 3a956abcf0c..33c211723ad 100644 --- a/build_autogenerated.yaml +++ b/build_autogenerated.yaml @@ -6418,9 +6418,11 @@ targets: gtest: true build: test language: c++ - headers: [] + headers: + - test/core/util/fake_stats_plugin.h src: - test/core/channel/call_tracer_test.cc + - test/core/util/fake_stats_plugin.cc deps: - gtest - grpc_test_util @@ -18016,6 +18018,7 @@ targets: run: false language: c++ headers: + - test/core/util/fake_stats_plugin.h - test/core/util/scoped_env_var.h - test/cpp/end2end/connection_attempt_injector.h - test/cpp/end2end/counted_service.h @@ -18056,6 +18059,7 @@ targets: - src/proto/grpc/testing/xds/v3/route.proto - src/proto/grpc/testing/xds/v3/router.proto - src/proto/grpc/testing/xds/v3/string.proto + - test/core/util/fake_stats_plugin.cc - test/cpp/end2end/connection_attempt_injector.cc - test/cpp/end2end/test_service_impl.cc - test/cpp/end2end/xds/xds_cluster_end2end_test.cc diff --git a/src/core/ext/filters/client_channel/client_channel.cc b/src/core/ext/filters/client_channel/client_channel.cc index fc4ce34926e..e7583e10ed3 100644 --- a/src/core/ext/filters/client_channel/client_channel.cc +++ b/src/core/ext/filters/client_channel/client_channel.cc @@ -2603,6 +2603,8 @@ class ClientChannel::LoadBalancedCall::LbCallState ServiceConfigCallData::CallAttributeInterface* GetCallAttribute( UniqueTypeName type) const override; + ClientCallTracer::CallAttemptTracer* GetCallAttemptTracer() const override; + private: LoadBalancedCall* lb_call_; }; @@ -2694,6 +2696,11 @@ ClientChannel::LoadBalancedCall::LbCallState::GetCallAttribute( return service_config_call_data->GetCallAttribute(type); } +ClientCallTracer::CallAttemptTracer* +ClientChannel::LoadBalancedCall::LbCallState::GetCallAttemptTracer() const { + return lb_call_->call_attempt_tracer(); +} + // // ClientChannel::LoadBalancedCall::BackendMetricAccessor // diff --git a/src/core/ext/filters/client_channel/client_channel_internal.h b/src/core/ext/filters/client_channel/client_channel_internal.h index c559eeec7f8..24cded22538 100644 --- a/src/core/ext/filters/client_channel/client_channel_internal.h +++ b/src/core/ext/filters/client_channel/client_channel_internal.h @@ -25,6 +25,7 @@ #include +#include "src/core/lib/channel/call_tracer.h" #include "src/core/lib/channel/context.h" #include "src/core/lib/gprpp/unique_type_name.h" #include "src/core/lib/load_balancing/lb_policy.h" @@ -49,6 +50,7 @@ class ClientChannelLbCallState : public LoadBalancingPolicy::CallState { public: virtual ServiceConfigCallData::CallAttributeInterface* GetCallAttribute( UniqueTypeName type) const = 0; + virtual ClientCallTracer::CallAttemptTracer* GetCallAttemptTracer() const = 0; }; // Internal type for ServiceConfigCallData. Handles call commits. diff --git a/src/core/ext/filters/client_channel/lb_policy/xds/xds_cluster_impl.cc b/src/core/ext/filters/client_channel/lb_policy/xds/xds_cluster_impl.cc index e2ff8824663..c495d83cd54 100644 --- a/src/core/ext/filters/client_channel/lb_policy/xds/xds_cluster_impl.cc +++ b/src/core/ext/filters/client_channel/lb_policy/xds/xds_cluster_impl.cc @@ -37,6 +37,7 @@ #include #include +#include "src/core/ext/filters/client_channel/client_channel_internal.h" #include "src/core/ext/filters/client_channel/lb_policy/backend_metric_data.h" #include "src/core/ext/filters/client_channel/lb_policy/child_policy_handler.h" #include "src/core/ext/filters/client_channel/lb_policy/xds/xds_channel_args.h" @@ -76,6 +77,8 @@ TraceFlag grpc_xds_cluster_impl_lb_trace(false, "xds_cluster_impl_lb"); namespace { +using OptionalLabelComponent = + ClientCallTracer::CallAttemptTracer::OptionalLabelComponent; using XdsConfig = XdsDependencyManager::XdsConfig; // @@ -215,6 +218,7 @@ class XdsClusterImplLb : public LoadBalancingPolicy { RefCountedPtr call_counter_; uint32_t max_concurrent_requests_; + std::shared_ptr> service_labels_; RefCountedPtr drop_config_; RefCountedPtr drop_stats_; RefCountedPtr picker_; @@ -358,6 +362,7 @@ XdsClusterImplLb::Picker::Picker(XdsClusterImplLb* xds_cluster_impl_lb, : call_counter_(xds_cluster_impl_lb->call_counter_), max_concurrent_requests_( xds_cluster_impl_lb->cluster_resource_->max_concurrent_requests), + service_labels_(xds_cluster_impl_lb->cluster_resource_->telemetry_labels), drop_config_(xds_cluster_impl_lb->drop_config_), drop_stats_(xds_cluster_impl_lb->drop_stats_), picker_(std::move(picker)) { @@ -369,6 +374,11 @@ XdsClusterImplLb::Picker::Picker(XdsClusterImplLb* xds_cluster_impl_lb, LoadBalancingPolicy::PickResult XdsClusterImplLb::Picker::Pick( LoadBalancingPolicy::PickArgs args) { + auto* call_state = static_cast(args.call_state); + if (call_state->GetCallAttemptTracer() != nullptr) { + call_state->GetCallAttemptTracer()->AddOptionalLabels( + OptionalLabelComponent::kXdsServiceLabels, service_labels_); + } // Handle EDS drops. const std::string* drop_category; if (drop_config_ != nullptr && drop_config_->ShouldDrop(&drop_category)) { diff --git a/src/core/ext/xds/xds_cluster.cc b/src/core/ext/xds/xds_cluster.cc index 7edd3e8be3c..5c2cca5933e 100644 --- a/src/core/ext/xds/xds_cluster.cc +++ b/src/core/ext/xds/xds_cluster.cc @@ -48,6 +48,7 @@ #include "envoy/extensions/upstreams/http/v3/http_protocol_options.upb.h" #include "google/protobuf/any.upb.h" #include "google/protobuf/duration.upb.h" +#include "google/protobuf/struct.upb.h" #include "google/protobuf/wrappers.upb.h" #include "upb/base/string_view.h" #include "upb/text/encode.h" @@ -703,6 +704,37 @@ absl::StatusOr> CdsResourceParse( cds_update->override_host_statuses.Add( XdsHealthStatus(XdsHealthStatus::kHealthy)); } + // Record telemetry labels (if any). + const envoy_config_core_v3_Metadata* metadata = + envoy_config_cluster_v3_Cluster_metadata(cluster); + if (metadata != nullptr) { + google_protobuf_Struct* telemetry_labels_struct; + if (envoy_config_core_v3_Metadata_filter_metadata_get( + metadata, + StdStringToUpbString( + absl::string_view("com.google.csm.telemetry_labels")), + &telemetry_labels_struct)) { + auto telemetry_labels = + std::make_shared>(); + size_t iter = kUpb_Map_Begin; + const google_protobuf_Struct_FieldsEntry* fields_entry; + while ((fields_entry = google_protobuf_Struct_fields_next( + telemetry_labels_struct, &iter)) != nullptr) { + // Adds any entry whose value is a string to telemetry_labels. + const google_protobuf_Value* value = + google_protobuf_Struct_FieldsEntry_value(fields_entry); + if (google_protobuf_Value_has_string_value(value)) { + telemetry_labels->emplace( + UpbStringToStdString( + google_protobuf_Struct_FieldsEntry_key(fields_entry)), + UpbStringToStdString(google_protobuf_Value_string_value(value))); + } + } + if (!telemetry_labels->empty()) { + cds_update->telemetry_labels = std::move(telemetry_labels); + } + } + } // Return result. if (!errors.ok()) { return errors.status(absl::StatusCode::kInvalidArgument, diff --git a/src/core/ext/xds/xds_cluster.h b/src/core/ext/xds/xds_cluster.h index f92cba623ee..405e6e9c25a 100644 --- a/src/core/ext/xds/xds_cluster.h +++ b/src/core/ext/xds/xds_cluster.h @@ -101,6 +101,8 @@ struct XdsClusterResource : public XdsResourceType::ResourceData { XdsHealthStatusSet override_host_statuses; + std::shared_ptr> telemetry_labels; + bool operator==(const XdsClusterResource& other) const { return type == other.type && lb_policy_config == other.lb_policy_config && lrs_load_reporting_server == other.lrs_load_reporting_server && diff --git a/src/core/lib/channel/call_tracer.cc b/src/core/lib/channel/call_tracer.cc index 74f58ab234b..d98a944b599 100644 --- a/src/core/lib/channel/call_tracer.cc +++ b/src/core/lib/channel/call_tracer.cc @@ -146,6 +146,13 @@ class DelegatingClientCallTracer : public ClientCallTracer { std::shared_ptr StartNewTcpTrace() override { return nullptr; } + void AddOptionalLabels( + OptionalLabelComponent component, + std::shared_ptr> labels) override { + for (auto* tracer : tracers_) { + tracer->AddOptionalLabels(component, labels); + } + } std::string TraceId() override { return tracers_[0]->TraceId(); } std::string SpanId() override { return tracers_[0]->SpanId(); } bool IsSampled() override { return tracers_[0]->IsSampled(); } diff --git a/src/core/lib/channel/call_tracer.h b/src/core/lib/channel/call_tracer.h index 3f2b10ead22..34c8342dc35 100644 --- a/src/core/lib/channel/call_tracer.h +++ b/src/core/lib/channel/call_tracer.h @@ -128,6 +128,11 @@ class ClientCallTracer : public CallTracerAnnotationInterface { // as transparent retry attempts.) class CallAttemptTracer : public CallTracerInterface { public: + enum class OptionalLabelComponent : std::uint8_t { + kXdsServiceLabels = 0, + kSize = 1, // keep last + }; + ~CallAttemptTracer() override {} // TODO(yashykt): The following two methods `RecordReceivedTrailingMetadata` // and `RecordEnd` should be moved into CallTracerInterface. @@ -140,6 +145,11 @@ class ClientCallTracer : public CallTracerAnnotationInterface { // Should be the last API call to the object. Once invoked, the tracer // library is free to destroy the object. virtual void RecordEnd(const gpr_timespec& latency) = 0; + + // Adds optional labels to be reported by the underlying tracer in a call. + virtual void AddOptionalLabels( + OptionalLabelComponent component, + std::shared_ptr> labels) = 0; }; ~ClientCallTracer() override {} diff --git a/src/cpp/ext/csm/metadata_exchange.cc b/src/cpp/ext/csm/metadata_exchange.cc index e0a34827658..459a03a15a5 100644 --- a/src/cpp/ext/csm/metadata_exchange.cc +++ b/src/cpp/ext/csm/metadata_exchange.cc @@ -42,6 +42,7 @@ #include +#include "src/core/lib/channel/call_tracer.h" #include "src/core/lib/gprpp/env.h" #include "src/core/lib/iomgr/error.h" #include "src/core/lib/iomgr/load_file.h" @@ -49,10 +50,14 @@ #include "src/core/lib/json/json_object_loader.h" #include "src/core/lib/json/json_reader.h" #include "src/core/lib/slice/slice_internal.h" +#include "src/cpp/ext/otel/key_value_iterable.h" namespace grpc { namespace internal { +using OptionalLabelComponent = + grpc_core::ClientCallTracer::CallAttemptTracer::OptionalLabelComponent; + namespace { // The keys that will be used in the Metadata Exchange between local and remote. @@ -427,5 +432,42 @@ void ServiceMeshLabelsInjector::AddLabels( serialized_labels_to_send_.Ref()); } +bool ServiceMeshLabelsInjector::AddOptionalLabels( + absl::Span>> + optional_labels_span, + opentelemetry::nostd::function_ref< + bool(opentelemetry::nostd::string_view, + opentelemetry::common::AttributeValue)> + callback) const { + // According to the CSM Observability Metric spec, if the control plane fails + // to provide these labels, the client will set their values to "unknown". + // These default values are set below. + absl::string_view service_name = "unknown"; + absl::string_view service_namespace = "unknown"; + // Performs JSON label name format to CSM Observability Metric spec format + // conversion. + if (optional_labels_span.size() > + static_cast(OptionalLabelComponent::kXdsServiceLabels)) { + const auto& optional_labels = optional_labels_span[static_cast( + OptionalLabelComponent::kXdsServiceLabels)]; + if (optional_labels != nullptr) { + auto it = optional_labels->find("service_name"); + if (it != optional_labels->end()) service_name = it->second; + it = optional_labels->find("service_namespace"); + if (it != optional_labels->end()) service_namespace = it->second; + } + } + return callback("csm.service_name", + AbslStrViewToOpenTelemetryStrView(service_name)) && + callback("csm.service_namespace_name", + AbslStrViewToOpenTelemetryStrView(service_namespace)); +} + +size_t ServiceMeshLabelsInjector::GetOptionalLabelsSize( + absl::Span>>) + const { + return 2; +} + } // namespace internal } // namespace grpc diff --git a/src/cpp/ext/csm/metadata_exchange.h b/src/cpp/ext/csm/metadata_exchange.h index c2c24ff2843..40f9ff5d51f 100644 --- a/src/cpp/ext/csm/metadata_exchange.h +++ b/src/cpp/ext/csm/metadata_exchange.h @@ -50,6 +50,20 @@ class ServiceMeshLabelsInjector : public LabelsInjector { void AddLabels(grpc_metadata_batch* outgoing_initial_metadata, LabelsIterable* labels_from_incoming_metadata) const override; + // Add optional labels to the traced calls. + bool AddOptionalLabels( + absl::Span>> + optional_labels_span, + opentelemetry::nostd::function_ref< + bool(opentelemetry::nostd::string_view, + opentelemetry::common::AttributeValue)> + callback) const override; + + // Gets the size of the actual optional labels. + size_t GetOptionalLabelsSize( + absl::Span>> + optional_labels_span) const override; + private: std::vector> local_labels_; grpc_core::Slice serialized_labels_to_send_; diff --git a/src/cpp/ext/filters/census/open_census_call_tracer.h b/src/cpp/ext/filters/census/open_census_call_tracer.h index 7b8fc4f61b5..235a19ff6a5 100644 --- a/src/cpp/ext/filters/census/open_census_call_tracer.h +++ b/src/cpp/ext/filters/census/open_census_call_tracer.h @@ -103,6 +103,9 @@ class OpenCensusCallTracer : public grpc_core::ClientCallTracer { void RecordAnnotation(absl::string_view annotation) override; void RecordAnnotation(const Annotation& annotation) override; std::shared_ptr StartNewTcpTrace() override; + void AddOptionalLabels( + OptionalLabelComponent, + std::shared_ptr>) override {} experimental::CensusContext* context() { return &context_; } diff --git a/src/cpp/ext/otel/key_value_iterable.h b/src/cpp/ext/otel/key_value_iterable.h index 89d0e6a2630..fe43b0136c9 100644 --- a/src/cpp/ext/otel/key_value_iterable.h +++ b/src/cpp/ext/otel/key_value_iterable.h @@ -53,11 +53,16 @@ class KeyValueIterable : public opentelemetry::common::KeyValueIterable { const std::vector>& injected_labels_from_plugin_options, absl::Span> - additional_labels) + additional_labels, + const ActivePluginOptionsView* active_plugin_options_view, + absl::Span>> + optional_labels_span) : injected_labels_iterable_(injected_labels_iterable), injected_labels_from_plugin_options_( injected_labels_from_plugin_options), - additional_labels_(additional_labels) {} + additional_labels_(additional_labels), + active_plugin_options_view_(active_plugin_options_view), + optional_labels_(optional_labels_span) {} bool ForEachKeyValue(opentelemetry::nostd::function_ref< bool(opentelemetry::nostd::string_view, @@ -72,6 +77,21 @@ class KeyValueIterable : public opentelemetry::common::KeyValueIterable { } } } + if (OpenTelemetryPluginState().labels_injector != nullptr && + !OpenTelemetryPluginState().labels_injector->AddOptionalLabels( + optional_labels_, callback)) { + return false; + } + if (active_plugin_options_view_ != nullptr && + !active_plugin_options_view_->ForEach( + [callback, this]( + const InternalOpenTelemetryPluginOption& plugin_option, + size_t /*index*/) { + return plugin_option.labels_injector()->AddOptionalLabels( + optional_labels_, callback); + })) { + return false; + } for (const auto& plugin_option_injected_iterable : injected_labels_from_plugin_options_) { if (plugin_option_injected_iterable != nullptr) { @@ -104,6 +124,19 @@ class KeyValueIterable : public opentelemetry::common::KeyValueIterable { } } size += additional_labels_.size(); + if (OpenTelemetryPluginState().labels_injector != nullptr) { + size += OpenTelemetryPluginState().labels_injector->GetOptionalLabelsSize( + optional_labels_); + } + if (active_plugin_options_view_ != nullptr) { + active_plugin_options_view_->ForEach( + [&size, this](const InternalOpenTelemetryPluginOption& plugin_option, + size_t /*index*/) { + size += plugin_option.labels_injector()->GetOptionalLabelsSize( + optional_labels_); + return true; + }); + } return size; } @@ -113,6 +146,9 @@ class KeyValueIterable : public opentelemetry::common::KeyValueIterable { injected_labels_from_plugin_options_; absl::Span> additional_labels_; + const ActivePluginOptionsView* active_plugin_options_view_; + absl::Span>> + optional_labels_; }; } // namespace internal diff --git a/src/cpp/ext/otel/otel_call_tracer.h b/src/cpp/ext/otel/otel_call_tracer.h index 8b90c92835f..cda80c39f74 100644 --- a/src/cpp/ext/otel/otel_call_tracer.h +++ b/src/cpp/ext/otel/otel_call_tracer.h @@ -91,6 +91,9 @@ class OpenTelemetryCallTracer : public grpc_core::ClientCallTracer { void RecordAnnotation(absl::string_view /*annotation*/) override; void RecordAnnotation(const Annotation& /*annotation*/) override; std::shared_ptr StartNewTcpTrace() override; + void AddOptionalLabels(OptionalLabelComponent component, + std::shared_ptr> + optional_labels) override; private: const OpenTelemetryCallTracer* parent_; @@ -98,6 +101,10 @@ class OpenTelemetryCallTracer : public grpc_core::ClientCallTracer { // Start time (for measuring latency). absl::Time start_time_; std::unique_ptr injected_labels_; + // The indices of the array correspond to the OptionalLabelComponent enum. + std::array>, + static_cast(OptionalLabelComponent::kSize)> + optional_labels_array_; std::vector> injected_labels_from_plugin_options_; }; diff --git a/src/cpp/ext/otel/otel_client_filter.cc b/src/cpp/ext/otel/otel_client_filter.cc index 4587b3836d5..0abd897bad3 100644 --- a/src/cpp/ext/otel/otel_client_filter.cc +++ b/src/cpp/ext/otel/otel_client_filter.cc @@ -132,7 +132,9 @@ OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer:: // avoid recording a subset of injected labels here. OpenTelemetryPluginState().client.attempt.started->Add( 1, KeyValueIterable(/*injected_labels_iterable=*/nullptr, {}, - additional_labels)); + additional_labels, + /*active_plugin_options_view=*/nullptr, + /*optional_labels_span=*/{})); } } @@ -150,6 +152,7 @@ void OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer:: injected_labels_from_plugin_options_.push_back( labels_injector->GetLabels(recv_initial_metadata)); } + return true; }); } @@ -166,6 +169,7 @@ void OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer:: if (labels_injector != nullptr) { labels_injector->AddLabels(send_initial_metadata, nullptr); } + return true; }); } @@ -206,9 +210,10 @@ void OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer:: {OpenTelemetryStatusKey(), grpc_status_code_to_string( static_cast(status.code()))}}}; - KeyValueIterable labels(injected_labels_.get(), - injected_labels_from_plugin_options_, - additional_labels); + KeyValueIterable labels( + injected_labels_.get(), injected_labels_from_plugin_options_, + additional_labels, &parent_->parent_->active_plugin_options_view(), + optional_labels_array_); if (OpenTelemetryPluginState().client.attempt.duration != nullptr) { OpenTelemetryPluginState().client.attempt.duration->Record( absl::ToDoubleSeconds(absl::Now() - start_time_), labels, @@ -262,6 +267,13 @@ OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer::StartNewTcpTrace() { return nullptr; } +void OpenTelemetryCallTracer::OpenTelemetryCallAttemptTracer::AddOptionalLabels( + OptionalLabelComponent component, + std::shared_ptr> optional_labels) { + optional_labels_array_[static_cast(component)] = + std::move(optional_labels); +} + // // OpenTelemetryCallTracer // diff --git a/src/cpp/ext/otel/otel_plugin.h b/src/cpp/ext/otel/otel_plugin.h index 76c4d0b62f4..f2c87da3f6a 100644 --- a/src/cpp/ext/otel/otel_plugin.h +++ b/src/cpp/ext/otel/otel_plugin.h @@ -79,6 +79,23 @@ class LabelsInjector { virtual void AddLabels( grpc_metadata_batch* outgoing_initial_metadata, LabelsIterable* labels_from_incoming_metadata) const = 0; + + // Adds optional labels to the traced calls. Each entry in the span + // corresponds to the CallAttemptTracer::OptionalLabelComponent enum. Returns + // false when callback returns false. + virtual bool AddOptionalLabels( + absl::Span>> + optional_labels_span, + opentelemetry::nostd::function_ref< + bool(opentelemetry::nostd::string_view, + opentelemetry::common::AttributeValue)> + callback) const = 0; + + // Gets the actual size of the optional labels that the Plugin is going to + // produce through the AddOptionalLabels method. + virtual size_t GetOptionalLabelsSize( + absl::Span>> + optional_labels_span) const = 0; }; class InternalOpenTelemetryPluginOption @@ -223,14 +240,17 @@ class ActivePluginOptionsView { }); } - void ForEach( - absl::FunctionRef + bool ForEach( + absl::FunctionRef func) const { for (size_t i = 0; i < OpenTelemetryPluginState().plugin_options.size(); ++i) { const auto& plugin_option = OpenTelemetryPluginState().plugin_options[i]; - if (active_mask_[i]) func(*plugin_option, i); + if (active_mask_[i] && !func(*plugin_option, i)) { + return false; + } } + return true; } private: diff --git a/src/cpp/ext/otel/otel_server_call_tracer.cc b/src/cpp/ext/otel/otel_server_call_tracer.cc index 7fa40c6ccc9..fd9734e502c 100644 --- a/src/cpp/ext/otel/otel_server_call_tracer.cc +++ b/src/cpp/ext/otel/otel_server_call_tracer.cc @@ -95,6 +95,7 @@ class OpenTelemetryServerCallTracer : public grpc_core::ServerCallTracer { send_initial_metadata, injected_labels_from_plugin_options_[index].get()); } + return true; }); } @@ -186,6 +187,7 @@ void OpenTelemetryServerCallTracer::RecordReceivedInitialMetadata( injected_labels_from_plugin_options_[index] = labels_injector->GetLabels(recv_initial_metadata); } + return true; }); registered_method_ = recv_initial_metadata->get(grpc_core::GrpcRegisteredMethod()) @@ -197,7 +199,8 @@ void OpenTelemetryServerCallTracer::RecordReceivedInitialMetadata( // avoid recording a subset of injected labels here. OpenTelemetryPluginState().server.call.started->Add( 1, KeyValueIterable(/*injected_labels_iterable=*/nullptr, {}, - additional_labels)); + additional_labels, + /*active_plugin_options_view=*/nullptr, {})); } } @@ -215,9 +218,11 @@ void OpenTelemetryServerCallTracer::RecordEnd( {{OpenTelemetryMethodKey(), MethodForStats()}, {OpenTelemetryStatusKey(), grpc_status_code_to_string(final_info->final_status)}}}; - KeyValueIterable labels(injected_labels_.get(), - injected_labels_from_plugin_options_, - additional_labels); + // Currently we do not have any optional labels on the server side. + KeyValueIterable labels( + injected_labels_.get(), injected_labels_from_plugin_options_, + additional_labels, + /*active_plugin_options_view=*/nullptr, /*optional_labels_span=*/{}); if (OpenTelemetryPluginState().server.call.duration != nullptr) { OpenTelemetryPluginState().server.call.duration->Record( absl::ToDoubleSeconds(elapsed_time_), labels, diff --git a/src/proto/grpc/testing/xds/v3/base.proto b/src/proto/grpc/testing/xds/v3/base.proto index 33719f687c5..fcf78419f58 100644 --- a/src/proto/grpc/testing/xds/v3/base.proto +++ b/src/proto/grpc/testing/xds/v3/base.proto @@ -129,3 +129,43 @@ message TransportSocket { google.protobuf.Any typed_config = 3; } } + +// Metadata provides additional inputs to filters based on matched listeners, +// filter chains, routes and endpoints. It is structured as a map, usually from +// filter name (in reverse DNS format) to metadata specific to the filter. Metadata +// key-values for a filter are merged as connection and request handling occurs, +// with later values for the same key overriding earlier values. +// +// An example use of metadata is providing additional values to +// http_connection_manager in the envoy.http_connection_manager.access_log +// namespace. +// +// Another example use of metadata is to per service config info in cluster metadata, which may get +// consumed by multiple filters. +// +// For load balancing, Metadata provides a means to subset cluster endpoints. +// Endpoints have a Metadata object associated and routes contain a Metadata +// object to match against. There are some well defined metadata used today for +// this purpose: +// +// * ``{"envoy.lb": {"canary": }}`` This indicates the canary status of an +// endpoint and is also used during header processing +// (x-envoy-upstream-canary) and for stats purposes. +// [#next-major-version: move to type/metadata/v2] +message Metadata { + // Key is the reverse DNS filter name, e.g. com.acme.widget. The ``envoy.*`` + // namespace is reserved for Envoy's built-in filters. + // If both ``filter_metadata`` and + // :ref:`typed_filter_metadata ` + // fields are present in the metadata with same keys, + // only ``typed_filter_metadata`` field will be parsed. + map filter_metadata = 1; + + // Key is the reverse DNS filter name, e.g. com.acme.widget. The ``envoy.*`` + // namespace is reserved for Envoy's built-in filters. + // The value is encoded as google.protobuf.Any. + // If both :ref:`filter_metadata ` + // and ``typed_filter_metadata`` fields are present in the metadata with same keys, + // only ``typed_filter_metadata`` field will be parsed. + map typed_filter_metadata = 2; +} diff --git a/src/proto/grpc/testing/xds/v3/cluster.proto b/src/proto/grpc/testing/xds/v3/cluster.proto index 75c01303d6a..a7c438399a4 100644 --- a/src/proto/grpc/testing/xds/v3/cluster.proto +++ b/src/proto/grpc/testing/xds/v3/cluster.proto @@ -252,6 +252,13 @@ message Cluster { // from the LRS stream here.] core.v3.ConfigSource lrs_server = 42; + // The Metadata field can be used to provide additional information about the + // cluster. It can be used for stats, logging, and varying filter behavior. + // Fields should use reverse DNS notation to denote which entity within Envoy + // will need the information. For instance, if the metadata is intended for + // the Router filter, the filter name should be specified as ``envoy.filters.http.router``. + core.v3.Metadata metadata = 25; + core.v3.TypedExtensionConfig upstream_config = 48; } diff --git a/src/python/grpcio_observability/grpc_observability/client_call_tracer.h b/src/python/grpcio_observability/grpc_observability/client_call_tracer.h index f51edc50620..d49ca42bd20 100644 --- a/src/python/grpcio_observability/grpc_observability/client_call_tracer.h +++ b/src/python/grpcio_observability/grpc_observability/client_call_tracer.h @@ -73,6 +73,10 @@ class PythonOpenCensusCallTracer : public grpc_core::ClientCallTracer { void RecordAnnotation(absl::string_view annotation) override; void RecordAnnotation(const Annotation& annotation) override; std::shared_ptr StartNewTcpTrace() override; + void AddOptionalLabels( + OptionalLabelComponent /*component*/, + std::shared_ptr> /*labels*/) + override {} private: // Maximum size of trace context is sent on the wire. diff --git a/test/core/channel/BUILD b/test/core/channel/BUILD index f374c745b5d..d872f7a5263 100644 --- a/test/core/channel/BUILD +++ b/test/core/channel/BUILD @@ -27,6 +27,7 @@ grpc_cc_test( uses_polling = False, deps = [ "//:grpc", + "//test/core/util:fake_stats_plugin", "//test/core/util:grpc_test_util", ], ) diff --git a/test/core/channel/call_tracer_test.cc b/test/core/channel/call_tracer_test.cc index b0dbcd5704d..d83093cb128 100644 --- a/test/core/channel/call_tracer_test.cc +++ b/test/core/channel/call_tracer_test.cc @@ -18,7 +18,6 @@ #include "src/core/lib/channel/call_tracer.h" -#include #include #include "gtest/gtest.h" @@ -26,115 +25,16 @@ #include #include -#include "src/core/lib/channel/tcp_tracer.h" #include "src/core/lib/gprpp/ref_counted_ptr.h" #include "src/core/lib/promise/context.h" #include "src/core/lib/resource_quota/memory_quota.h" #include "src/core/lib/resource_quota/resource_quota.h" +#include "test/core/util/fake_stats_plugin.h" #include "test/core/util/test_config.h" namespace grpc_core { namespace { -class FakeClientCallTracer : public ClientCallTracer { - public: - class FakeClientCallAttemptTracer - : public ClientCallTracer::CallAttemptTracer { - public: - explicit FakeClientCallAttemptTracer( - std::vector* annotation_logger) - : annotation_logger_(annotation_logger) {} - ~FakeClientCallAttemptTracer() override {} - void RecordSendInitialMetadata( - grpc_metadata_batch* /*send_initial_metadata*/) override {} - void RecordSendTrailingMetadata( - grpc_metadata_batch* /*send_trailing_metadata*/) override {} - void RecordSendMessage(const SliceBuffer& /*send_message*/) override {} - void RecordSendCompressedMessage( - const SliceBuffer& /*send_compressed_message*/) override {} - void RecordReceivedInitialMetadata( - grpc_metadata_batch* /*recv_initial_metadata*/) override {} - void RecordReceivedMessage(const SliceBuffer& /*recv_message*/) override {} - void RecordReceivedDecompressedMessage( - const SliceBuffer& /*recv_decompressed_message*/) override {} - void RecordCancel(grpc_error_handle /*cancel_error*/) override {} - void RecordReceivedTrailingMetadata( - absl::Status /*status*/, - grpc_metadata_batch* /*recv_trailing_metadata*/, - const grpc_transport_stream_stats* /*transport_stream_stats*/) - override {} - void RecordEnd(const gpr_timespec& /*latency*/) override { delete this; } - void RecordAnnotation(absl::string_view annotation) override { - annotation_logger_->push_back(std::string(annotation)); - } - void RecordAnnotation(const Annotation& /*annotation*/) override {} - std::shared_ptr StartNewTcpTrace() override { - return nullptr; - } - std::string TraceId() override { return ""; } - std::string SpanId() override { return ""; } - bool IsSampled() override { return false; } - - private: - std::vector* annotation_logger_; - }; - - explicit FakeClientCallTracer(std::vector* annotation_logger) - : annotation_logger_(annotation_logger) {} - ~FakeClientCallTracer() override {} - CallAttemptTracer* StartNewAttempt(bool /*is_transparent_retry*/) override { - return GetContext()->ManagedNew( - annotation_logger_); - } - - void RecordAnnotation(absl::string_view annotation) override { - annotation_logger_->push_back(std::string(annotation)); - } - void RecordAnnotation(const Annotation& /*annotation*/) override {} - std::string TraceId() override { return ""; } - std::string SpanId() override { return ""; } - bool IsSampled() override { return false; } - - private: - std::vector* annotation_logger_; -}; - -class FakeServerCallTracer : public ServerCallTracer { - public: - explicit FakeServerCallTracer(std::vector* annotation_logger) - : annotation_logger_(annotation_logger) {} - ~FakeServerCallTracer() override {} - void RecordSendInitialMetadata( - grpc_metadata_batch* /*send_initial_metadata*/) override {} - void RecordSendTrailingMetadata( - grpc_metadata_batch* /*send_trailing_metadata*/) override {} - void RecordSendMessage(const SliceBuffer& /*send_message*/) override {} - void RecordSendCompressedMessage( - const SliceBuffer& /*send_compressed_message*/) override {} - void RecordReceivedInitialMetadata( - grpc_metadata_batch* /*recv_initial_metadata*/) override {} - void RecordReceivedMessage(const SliceBuffer& /*recv_message*/) override {} - void RecordReceivedDecompressedMessage( - const SliceBuffer& /*recv_decompressed_message*/) override {} - void RecordCancel(grpc_error_handle /*cancel_error*/) override {} - void RecordReceivedTrailingMetadata( - grpc_metadata_batch* /*recv_trailing_metadata*/) override {} - void RecordEnd(const grpc_call_final_info* /*final_info*/) override {} - void RecordAnnotation(absl::string_view annotation) override { - annotation_logger_->push_back(std::string(annotation)); - } - void RecordAnnotation(const Annotation& /*annotation*/) override {} - std::shared_ptr StartNewTcpTrace() override { - return nullptr; - } - std::string TraceId() override { return ""; } - std::string SpanId() override { return ""; } - bool IsSampled() override { return false; } - - private: - std::vector* annotation_logger_; -}; - class CallTracerTest : public ::testing::Test { protected: void SetUp() override { diff --git a/test/core/client_channel/lb_policy/lb_policy_test_lib.h b/test/core/client_channel/lb_policy/lb_policy_test_lib.h index 4be86bcad9b..48f97ca9524 100644 --- a/test/core/client_channel/lb_policy/lb_policy_test_lib.h +++ b/test/core/client_channel/lb_policy/lb_policy_test_lib.h @@ -671,6 +671,10 @@ class LoadBalancingPolicyTest : public ::testing::Test { return nullptr; } + ClientCallTracer::CallAttemptTracer* GetCallAttemptTracer() const override { + return nullptr; + } + std::vector allocations_; std::map attributes_; diff --git a/test/core/end2end/tests/http2_stats.cc b/test/core/end2end/tests/http2_stats.cc index 14f8dbeb648..aac381e2644 100644 --- a/test/core/end2end/tests/http2_stats.cc +++ b/test/core/end2end/tests/http2_stats.cc @@ -97,6 +97,11 @@ class FakeCallTracer : public ClientCallTracer { void RecordAnnotation(absl::string_view /*annotation*/) override {} void RecordAnnotation(const Annotation& /*annotation*/) override {} + void AddOptionalLabels( + OptionalLabelComponent /*component*/, + std::shared_ptr> /*labels*/) + override {} + static grpc_transport_stream_stats transport_stream_stats() { MutexLock lock(g_mu); return transport_stream_stats_; diff --git a/test/core/util/BUILD b/test/core/util/BUILD index d6e0bad20c9..b519c306ca0 100644 --- a/test/core/util/BUILD +++ b/test/core/util/BUILD @@ -491,3 +491,13 @@ grpc_cc_library( "//src/core:resource_quota", ], ) + +grpc_cc_library( + name = "fake_stats_plugin", + srcs = ["fake_stats_plugin.cc"], + hdrs = ["fake_stats_plugin.h"], + deps = [ + "//:grpc", + "//src/core:examine_stack", + ], +) diff --git a/test/core/util/fake_stats_plugin.cc b/test/core/util/fake_stats_plugin.cc new file mode 100644 index 00000000000..f8819461da4 --- /dev/null +++ b/test/core/util/fake_stats_plugin.cc @@ -0,0 +1,82 @@ +// Copyright 2023 The gRPC Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#include "test/core/util/fake_stats_plugin.h" + +#include "src/core/lib/config/core_configuration.h" + +namespace grpc_core { + +class FakeStatsClientFilter : public ChannelFilter { + public: + static const grpc_channel_filter kFilter; + + static absl::StatusOr Create( + const ChannelArgs& /*args*/, ChannelFilter::Args /*filter_args*/); + + ArenaPromise MakeCallPromise( + CallArgs call_args, NextPromiseFactory next_promise_factory) override; + + private: + explicit FakeStatsClientFilter( + FakeClientCallTracerFactory* fake_client_call_tracer_factory); + FakeClientCallTracerFactory* const fake_client_call_tracer_factory_; +}; + +const grpc_channel_filter FakeStatsClientFilter::kFilter = + MakePromiseBasedFilter( + "fake_stats_client"); + +absl::StatusOr FakeStatsClientFilter::Create( + const ChannelArgs& args, ChannelFilter::Args /*filter_args*/) { + auto* fake_client_call_tracer_factory = + args.GetPointer( + GRPC_ARG_INJECT_FAKE_CLIENT_CALL_TRACER_FACTORY); + GPR_ASSERT(fake_client_call_tracer_factory != nullptr); + return FakeStatsClientFilter(fake_client_call_tracer_factory); +} + +ArenaPromise FakeStatsClientFilter::MakeCallPromise( + CallArgs call_args, NextPromiseFactory next_promise_factory) { + FakeClientCallTracer* client_call_tracer = + fake_client_call_tracer_factory_->CreateFakeClientCallTracer(); + if (client_call_tracer != nullptr) { + auto* call_context = GetContext(); + call_context[GRPC_CONTEXT_CALL_TRACER_ANNOTATION_INTERFACE].value = + client_call_tracer; + call_context[GRPC_CONTEXT_CALL_TRACER_ANNOTATION_INTERFACE].destroy = + nullptr; + } + return next_promise_factory(std::move(call_args)); +} + +FakeStatsClientFilter::FakeStatsClientFilter( + FakeClientCallTracerFactory* fake_client_call_tracer_factory) + : fake_client_call_tracer_factory_(fake_client_call_tracer_factory) {} + +void RegisterFakeStatsPlugin() { + CoreConfiguration::RegisterBuilder( + [](CoreConfiguration::Builder* builder) mutable { + builder->channel_init() + ->RegisterFilter(GRPC_CLIENT_CHANNEL, + &FakeStatsClientFilter::kFilter) + .If([](const ChannelArgs& args) { + return args.GetPointer( + GRPC_ARG_INJECT_FAKE_CLIENT_CALL_TRACER_FACTORY) != + nullptr; + }); + }); +} + +} // namespace grpc_core diff --git a/test/core/util/fake_stats_plugin.h b/test/core/util/fake_stats_plugin.h new file mode 100644 index 00000000000..57a659ba4bd --- /dev/null +++ b/test/core/util/fake_stats_plugin.h @@ -0,0 +1,192 @@ +// Copyright 2023 The gRPC Authors. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#ifndef GRPC_TEST_CORE_UTIL_FAKE_STATS_PLUGIN_H +#define GRPC_TEST_CORE_UTIL_FAKE_STATS_PLUGIN_H + +#include +#include +#include + +#include "src/core/lib/channel/call_tracer.h" +#include "src/core/lib/channel/promise_based_filter.h" +#include "src/core/lib/channel/tcp_tracer.h" + +namespace grpc_core { + +// Registers a FakeStatsClientFilter as a client channel filter if there is a +// FakeClientCallTracerFactory in the channel args. This filter will use the +// FakeClientCallTracerFactory to create and inject a FakeClientCallTracer into +// the call context. +// Example usage: +// RegisterFakeStatsPlugin(); // before grpc_init() +// +// // Creates a FakeClientCallTracerFactory and adds it into the channel args. +// FakeClientCallTracerFactory fake_client_call_tracer_factory; +// ChannelArguments channel_args; +// channel_args.SetPointer(GRPC_ARG_INJECT_FAKE_CLIENT_CALL_TRACER_FACTORY, +// &fake_client_call_tracer_factory); +// +// // After the system under test has been executed (e.g. an RPC has been +// // sent), use the FakeClientCallTracerFactory to verify certain +// // expectations. +// EXPECT_THAT(fake_client_call_tracer_factory.GetLastFakeClientCallTracer() +// ->GetLastCallAttemptTracer() +// ->GetOptionalLabels(), +// VerifyCsmServiceLabels()); +void RegisterFakeStatsPlugin(); + +class FakeClientCallTracer : public ClientCallTracer { + public: + class FakeClientCallAttemptTracer + : public ClientCallTracer::CallAttemptTracer { + public: + explicit FakeClientCallAttemptTracer( + std::vector* annotation_logger) + : annotation_logger_(annotation_logger) {} + ~FakeClientCallAttemptTracer() override {} + void RecordSendInitialMetadata( + grpc_metadata_batch* /*send_initial_metadata*/) override {} + void RecordSendTrailingMetadata( + grpc_metadata_batch* /*send_trailing_metadata*/) override {} + void RecordSendMessage(const SliceBuffer& /*send_message*/) override {} + void RecordSendCompressedMessage( + const SliceBuffer& /*send_compressed_message*/) override {} + void RecordReceivedInitialMetadata( + grpc_metadata_batch* /*recv_initial_metadata*/) override {} + void RecordReceivedMessage(const SliceBuffer& /*recv_message*/) override {} + void RecordReceivedDecompressedMessage( + const SliceBuffer& /*recv_decompressed_message*/) override {} + void RecordCancel(grpc_error_handle /*cancel_error*/) override {} + void RecordReceivedTrailingMetadata( + absl::Status /*status*/, + grpc_metadata_batch* /*recv_trailing_metadata*/, + const grpc_transport_stream_stats* /*transport_stream_stats*/) + override {} + void RecordEnd(const gpr_timespec& /*latency*/) override {} + void RecordAnnotation(absl::string_view annotation) override { + annotation_logger_->push_back(std::string(annotation)); + } + void RecordAnnotation(const Annotation& /*annotation*/) override {} + std::shared_ptr StartNewTcpTrace() override { + return nullptr; + } + void AddOptionalLabels( + OptionalLabelComponent component, + std::shared_ptr> labels) override { + optional_labels_.emplace(component, std::move(labels)); + } + std::string TraceId() override { return ""; } + std::string SpanId() override { return ""; } + bool IsSampled() override { return false; } + + const std::map>>& + GetOptionalLabels() const { + return optional_labels_; + } + + private: + std::vector* annotation_logger_; + std::map>> + optional_labels_; + }; + + explicit FakeClientCallTracer(std::vector* annotation_logger) + : annotation_logger_(annotation_logger) {} + ~FakeClientCallTracer() override {} + CallAttemptTracer* StartNewAttempt(bool /*is_transparent_retry*/) override { + call_attempt_tracers_.emplace_back( + new FakeClientCallAttemptTracer(annotation_logger_)); + return call_attempt_tracers_.back().get(); + } + + void RecordAnnotation(absl::string_view annotation) override { + annotation_logger_->push_back(std::string(annotation)); + } + void RecordAnnotation(const Annotation& /*annotation*/) override {} + std::string TraceId() override { return ""; } + std::string SpanId() override { return ""; } + bool IsSampled() override { return false; } + + FakeClientCallAttemptTracer* GetLastCallAttemptTracer() const { + return call_attempt_tracers_.back().get(); + } + + private: + std::vector* annotation_logger_; + std::vector> + call_attempt_tracers_; +}; + +#define GRPC_ARG_INJECT_FAKE_CLIENT_CALL_TRACER_FACTORY \ + "grpc.testing.inject_fake_client_call_tracer_factory" + +class FakeClientCallTracerFactory { + public: + FakeClientCallTracer* CreateFakeClientCallTracer() { + fake_client_call_tracers_.emplace_back( + new FakeClientCallTracer(&annotation_logger_)); + return fake_client_call_tracers_.back().get(); + } + + FakeClientCallTracer* GetLastFakeClientCallTracer() { + return fake_client_call_tracers_.back().get(); + } + + private: + std::vector annotation_logger_; + std::vector> fake_client_call_tracers_; +}; + +class FakeServerCallTracer : public ServerCallTracer { + public: + explicit FakeServerCallTracer(std::vector* annotation_logger) + : annotation_logger_(annotation_logger) {} + ~FakeServerCallTracer() override {} + void RecordSendInitialMetadata( + grpc_metadata_batch* /*send_initial_metadata*/) override {} + void RecordSendTrailingMetadata( + grpc_metadata_batch* /*send_trailing_metadata*/) override {} + void RecordSendMessage(const SliceBuffer& /*send_message*/) override {} + void RecordSendCompressedMessage( + const SliceBuffer& /*send_compressed_message*/) override {} + void RecordReceivedInitialMetadata( + grpc_metadata_batch* /*recv_initial_metadata*/) override {} + void RecordReceivedMessage(const SliceBuffer& /*recv_message*/) override {} + void RecordReceivedDecompressedMessage( + const SliceBuffer& /*recv_decompressed_message*/) override {} + void RecordCancel(grpc_error_handle /*cancel_error*/) override {} + void RecordReceivedTrailingMetadata( + grpc_metadata_batch* /*recv_trailing_metadata*/) override {} + void RecordEnd(const grpc_call_final_info* /*final_info*/) override {} + void RecordAnnotation(absl::string_view annotation) override { + annotation_logger_->push_back(std::string(annotation)); + } + void RecordAnnotation(const Annotation& /*annotation*/) override {} + std::shared_ptr StartNewTcpTrace() override { + return nullptr; + } + std::string TraceId() override { return ""; } + std::string SpanId() override { return ""; } + bool IsSampled() override { return false; } + + private: + std::vector* annotation_logger_; +}; + +} // namespace grpc_core + +#endif // GRPC_TEST_CORE_UTIL_FAKE_STATS_PLUGIN_H diff --git a/test/core/xds/xds_cluster_resource_type_test.cc b/test/core/xds/xds_cluster_resource_type_test.cc index 2feb6199fc2..ed8744fa467 100644 --- a/test/core/xds/xds_cluster_resource_type_test.cc +++ b/test/core/xds/xds_cluster_resource_type_test.cc @@ -20,6 +20,7 @@ #include #include +#include #include #include "absl/status/status.h" @@ -1613,6 +1614,95 @@ TEST_F(HostOverrideStatusTest, CanExplicitlySetToEmpty) { EXPECT_EQ(resource.override_host_statuses.ToString(), "{}"); } +using TelemetryLabelTest = XdsClusterTest; + +TEST_F(TelemetryLabelTest, ValidServiceLabelsConfig) { + Cluster cluster; + cluster.set_type(cluster.EDS); + cluster.mutable_eds_cluster_config()->mutable_eds_config()->mutable_self(); + auto& filter_map = *cluster.mutable_metadata()->mutable_filter_metadata(); + auto& label_map = + *filter_map["com.google.csm.telemetry_labels"].mutable_fields(); + *label_map["service_name"].mutable_string_value() = "abc"; + *label_map["service_namespace"].mutable_string_value() = "xyz"; + std::string serialized_resource; + ASSERT_TRUE(cluster.SerializeToString(&serialized_resource)); + auto* resource_type = XdsClusterResourceType::Get(); + auto decode_result = + resource_type->Decode(decode_context_, serialized_resource); + ASSERT_TRUE(decode_result.resource.ok()) << decode_result.resource.status(); + auto& resource = + static_cast(**decode_result.resource); + EXPECT_THAT(*resource.telemetry_labels, + ::testing::UnorderedElementsAre( + ::testing::Pair("service_name", "abc"), + ::testing::Pair("service_namespace", "xyz"))); +} + +TEST_F(TelemetryLabelTest, MissingMetadataField) { + Cluster cluster; + cluster.set_type(cluster.EDS); + cluster.mutable_eds_cluster_config()->mutable_eds_config()->mutable_self(); + std::string serialized_resource; + ASSERT_TRUE(cluster.SerializeToString(&serialized_resource)); + auto* resource_type = XdsClusterResourceType::Get(); + auto decode_result = + resource_type->Decode(decode_context_, serialized_resource); + ASSERT_TRUE(decode_result.resource.ok()) << decode_result.resource.status(); + auto& resource = + static_cast(**decode_result.resource); + EXPECT_EQ(resource.telemetry_labels, nullptr); +} + +TEST_F(TelemetryLabelTest, MissingCsmFilterMetadataField) { + Cluster cluster; + cluster.set_type(cluster.EDS); + cluster.mutable_eds_cluster_config()->mutable_eds_config()->mutable_self(); + auto& filter_map = *cluster.mutable_metadata()->mutable_filter_metadata(); + auto& label_map = *filter_map["some_key"].mutable_fields(); + *label_map["some_value"].mutable_string_value() = "abc"; + std::string serialized_resource; + ASSERT_TRUE(cluster.SerializeToString(&serialized_resource)); + auto* resource_type = XdsClusterResourceType::Get(); + auto decode_result = + resource_type->Decode(decode_context_, serialized_resource); + ASSERT_TRUE(decode_result.resource.ok()) << decode_result.resource.status(); + auto& resource = + static_cast(**decode_result.resource); + EXPECT_EQ(resource.telemetry_labels, nullptr); +} + +TEST_F(TelemetryLabelTest, IgnoreNonStringEntries) { + Cluster cluster; + cluster.set_type(cluster.EDS); + cluster.mutable_eds_cluster_config()->mutable_eds_config()->mutable_self(); + auto& filter_map = *cluster.mutable_metadata()->mutable_filter_metadata(); + auto& label_map = + *filter_map["com.google.csm.telemetry_labels"].mutable_fields(); + label_map["bool_value"].set_bool_value(true); + label_map["number_value"].set_number_value(3.14); + *label_map["string_value"].mutable_string_value() = "abc"; + label_map["null_value"].set_null_value(::google::protobuf::NULL_VALUE); + auto& list_value_values = + *label_map["list_value"].mutable_list_value()->mutable_values(); + *list_value_values.Add()->mutable_string_value() = "efg"; + list_value_values.Add()->set_number_value(3.14); + auto& struct_value_fields = + *label_map["struct_value"].mutable_struct_value()->mutable_fields(); + struct_value_fields["bool_value"].set_bool_value(false); + std::string serialized_resource; + ASSERT_TRUE(cluster.SerializeToString(&serialized_resource)); + auto* resource_type = XdsClusterResourceType::Get(); + auto decode_result = + resource_type->Decode(decode_context_, serialized_resource); + ASSERT_TRUE(decode_result.resource.ok()) << decode_result.resource.status(); + auto& resource = + static_cast(**decode_result.resource); + EXPECT_THAT( + *resource.telemetry_labels, + ::testing::UnorderedElementsAre(::testing::Pair("string_value", "abc"))); +} + } // namespace } // namespace testing } // namespace grpc_core diff --git a/test/cpp/end2end/xds/BUILD b/test/cpp/end2end/xds/BUILD index 21a034e74c4..8fad62d50bf 100644 --- a/test/cpp/end2end/xds/BUILD +++ b/test/cpp/end2end/xds/BUILD @@ -169,6 +169,7 @@ grpc_cc_test( "//:gpr", "//:grpc", "//:grpc++", + "//test/core/util:fake_stats_plugin", "//test/core/util:grpc_test_util", "//test/core/util:scoped_env_var", "//test/cpp/end2end:connection_attempt_injector", diff --git a/test/cpp/end2end/xds/xds_cluster_end2end_test.cc b/test/cpp/end2end/xds/xds_cluster_end2end_test.cc index d4c0ad4a823..7662f064ed8 100644 --- a/test/cpp/end2end/xds/xds_cluster_end2end_test.cc +++ b/test/cpp/end2end/xds/xds_cluster_end2end_test.cc @@ -25,9 +25,12 @@ #include "src/core/ext/filters/client_channel/backup_poller.h" #include "src/core/lib/address_utils/sockaddr_utils.h" +#include "src/core/lib/channel/call_tracer.h" #include "src/core/lib/config/config_vars.h" #include "src/core/lib/experiments/experiments.h" +#include "src/core/lib/surface/call.h" #include "src/proto/grpc/testing/xds/v3/orca_load_report.pb.h" +#include "test/core/util/fake_stats_plugin.h" #include "test/core/util/scoped_env_var.h" #include "test/cpp/end2end/connection_attempt_injector.h" #include "test/cpp/end2end/xds/xds_end2end_test_lib.h" @@ -42,6 +45,8 @@ using ::envoy::config::core::v3::HealthStatus; using ::envoy::type::v3::FractionalPercent; using ClientStats = LrsServiceImpl::ClientStats; +using OptionalLabelComponent = + grpc_core::ClientCallTracer::CallAttemptTracer::OptionalLabelComponent; constexpr char kLbDropType[] = "lb"; constexpr char kThrottleDropType[] = "throttle"; @@ -304,6 +309,39 @@ TEST_P(CdsTest, ClusterChangeAfterAdsCallFails) { WaitForBackend(DEBUG_LOCATION, 1); } +TEST_P(CdsTest, VerifyCsmServiceLabelsParsing) { + // Injects a fake client call tracer factory. Try keep this at top. + grpc_core::FakeClientCallTracerFactory fake_client_call_tracer_factory; + CreateAndStartBackends(1); + // Populates EDS resources. + EdsResourceArgs args({{"locality0", CreateEndpointsForBackends()}}); + balancer_->ads_service()->SetEdsResource(BuildEdsResource(args)); + // Populates service labels to CDS resources. + auto cluster = default_cluster_; + auto& filter_map = *cluster.mutable_metadata()->mutable_filter_metadata(); + auto& label_map = + *filter_map["com.google.csm.telemetry_labels"].mutable_fields(); + *label_map["service_name"].mutable_string_value() = "myservice"; + *label_map["service_namespace"].mutable_string_value() = "mynamespace"; + balancer_->ads_service()->SetCdsResource(cluster); + ChannelArguments channel_args; + channel_args.SetPointer(GRPC_ARG_INJECT_FAKE_CLIENT_CALL_TRACER_FACTORY, + &fake_client_call_tracer_factory); + ResetStub(/*failover_timeout_ms=*/0, &channel_args); + // Sends an RPC and verifies that the service labels are recorded in the fake + // client call tracer. + CheckRpcSendOk(DEBUG_LOCATION); + EXPECT_THAT(fake_client_call_tracer_factory.GetLastFakeClientCallTracer() + ->GetLastCallAttemptTracer() + ->GetOptionalLabels(), + ::testing::ElementsAre(::testing::Pair( + OptionalLabelComponent::kXdsServiceLabels, + ::testing::Pointee(::testing::ElementsAre( + ::testing::Pair("service_name", "myservice"), + ::testing::Pair("service_namespace", "mynamespace")))))); + balancer_->Shutdown(); +} + // // CDS deletion tests // @@ -1874,6 +1912,7 @@ int main(int argc, char** argv) { // Workaround Apple CFStream bug grpc_core::SetEnv("grpc_cfstream", "0"); #endif + grpc_core::RegisterFakeStatsPlugin(); grpc_init(); grpc::testing::ConnectionAttemptInjector::Init(); const auto result = RUN_ALL_TESTS(); diff --git a/test/cpp/ext/csm/metadata_exchange_test.cc b/test/cpp/ext/csm/metadata_exchange_test.cc index d5411dcafe9..cdef39dd921 100644 --- a/test/cpp/ext/csm/metadata_exchange_test.cc +++ b/test/cpp/ext/csm/metadata_exchange_test.cc @@ -113,7 +113,8 @@ class MetadataExchangeTest public ::testing::WithParamInterface { protected: void Init(const absl::flat_hash_set& metric_names, - bool enable_client_side_injector = true) { + bool enable_client_side_injector = true, + const std::map& labels_to_inject = {}) { const char* kBootstrap = "{\"node\": {\"id\": " "\"projects/1234567890/networks/mesh:mesh-id/nodes/" @@ -137,7 +138,7 @@ class MetadataExchangeTest /*labels_injector=*/ std::make_unique( GetParam().GetTestResource().GetAttributes()), - /*test_no_meter_provider=*/false, + /*test_no_meter_provider=*/false, labels_to_inject, /*target_selector=*/ [enable_client_side_injector](absl::string_view /*target*/) { return enable_client_side_injector; @@ -156,11 +157,19 @@ class MetadataExchangeTest void VerifyServiceMeshAttributes( const std::map& - attributes) { + attributes, + bool verify_client_only_attributes = true) { EXPECT_EQ( absl::get(attributes.at("csm.workload_canonical_service")), "canonical_service"); EXPECT_EQ(absl::get(attributes.at("csm.mesh_id")), "mesh-id"); + if (verify_client_only_attributes) { + EXPECT_EQ(absl::get(attributes.at("csm.service_name")), + "unknown"); + EXPECT_EQ( + absl::get(attributes.at("csm.service_namespace_name")), + "unknown"); + } switch (GetParam().type()) { case TestScenario::ResourceType::kGke: EXPECT_EQ( @@ -295,7 +304,8 @@ TEST_P(MetadataExchangeTest, ServerCallDuration) { const auto& attributes = data[kMetricName][0].attributes.GetAttributes(); EXPECT_EQ(absl::get(attributes.at("grpc.method")), kMethodName); EXPECT_EQ(absl::get(attributes.at("grpc.status")), "OK"); - VerifyServiceMeshAttributes(attributes); + VerifyServiceMeshAttributes(attributes, + /*verify_client_only_attributes=*/false); } // Test that the server records unknown when the client does not send metadata @@ -328,6 +338,27 @@ TEST_P(MetadataExchangeTest, ClientDoesNotSendMetadata) { "unknown"); } +TEST_P(MetadataExchangeTest, VerifyCsmServiceLabels) { + Init(/*metric_names=*/{grpc::experimental::OpenTelemetryPluginBuilder:: + kClientAttemptDurationInstrumentName}, + /*enable_client_side_injector=*/true, + // Injects CSM service labels to be recorded in the call. + {{"service_name", "myservice"}, {"service_namespace", "mynamespace"}}); + SendRPC(); + const char* kMetricName = "grpc.client.attempt.duration"; + auto data = ReadCurrentMetricsData( + [&](const absl::flat_hash_map< + std::string, + std::vector>& + data) { return !data.contains(kMetricName); }); + ASSERT_EQ(data[kMetricName].size(), 1); + const auto& attributes = data[kMetricName][0].attributes.GetAttributes(); + EXPECT_EQ(absl::get(attributes.at("csm.service_name")), + "myservice"); + EXPECT_EQ(absl::get(attributes.at("csm.service_namespace_name")), + "mynamespace"); +} + INSTANTIATE_TEST_SUITE_P( MetadataExchange, MetadataExchangeTest, ::testing::Values( diff --git a/test/cpp/ext/otel/otel_plugin_test.cc b/test/cpp/ext/otel/otel_plugin_test.cc index 124004ab37a..541e49d4177 100644 --- a/test/cpp/ext/otel/otel_plugin_test.cc +++ b/test/cpp/ext/otel/otel_plugin_test.cc @@ -309,6 +309,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, TargetSelectorReturnsTrue) { opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/ [](absl::string_view /*target*/) { return true; }); SendRPC(); @@ -345,6 +346,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, TargetSelectorReturnsFalse) { opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/ [](absl::string_view /*target*/) { return false; }); SendRPC(); @@ -364,6 +366,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, TargetAttributeFilterReturnsTrue) { opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/[](absl::string_view /*target*/) { return true; @@ -402,6 +405,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, TargetAttributeFilterReturnsFalse) { opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ [server_address = canonical_server_address_]( @@ -469,6 +473,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -508,6 +513,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -575,6 +581,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -613,6 +620,7 @@ TEST_F(OpenTelemetryPluginEnd2EndTest, /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -685,6 +693,22 @@ class CustomLabelInjector : public grpc::internal::LabelsInjector { grpc::internal::LabelsIterable* /*labels_from_incoming_metadata*/) const override {} + bool AddOptionalLabels( + absl::Span>> + /*optional_labels_span*/, + opentelemetry::nostd::function_ref< + bool(opentelemetry::nostd::string_view, + opentelemetry::common::AttributeValue)> + /*callback*/) const override { + return true; + } + + size_t GetOptionalLabelsSize( + absl::Span>> + /*optional_labels_span*/) const override { + return 0; + } + private: std::pair label_; }; @@ -731,6 +755,7 @@ TEST_F(OpenTelemetryPluginOptionEnd2EndTest, Basic) { /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -772,6 +797,7 @@ TEST_F(OpenTelemetryPluginOptionEnd2EndTest, ClientOnlyPluginOption) { /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -814,6 +840,7 @@ TEST_F(OpenTelemetryPluginOptionEnd2EndTest, ServerOnlyPluginOption) { /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), @@ -870,6 +897,7 @@ TEST_F(OpenTelemetryPluginOptionEnd2EndTest, /*resource=*/opentelemetry::sdk::resource::Resource::Create({}), /*labels_injector=*/nullptr, /*test_no_meter_provider=*/false, + /*labels_to_inject=*/{}, /*target_selector=*/absl::AnyInvocable(), /*target_attribute_filter=*/ absl::AnyInvocable(), diff --git a/test/cpp/ext/otel/otel_test_library.cc b/test/cpp/ext/otel/otel_test_library.cc index 53aa350435f..6f8e55a96d1 100644 --- a/test/cpp/ext/otel/otel_test_library.cc +++ b/test/cpp/ext/otel/otel_test_library.cc @@ -29,6 +29,7 @@ #include #include "src/core/lib/channel/call_tracer.h" +#include "src/core/lib/channel/promise_based_filter.h" #include "src/core/lib/config/core_configuration.h" #include "src/core/lib/gprpp/notification.h" #include "test/core/util/test_config.h" @@ -38,11 +39,55 @@ namespace grpc { namespace testing { +#define GRPC_ARG_LABELS_TO_INJECT "grpc.testing.labels_to_inject" + +// A subchannel filter that adds the service labels for test to the +// CallAttemptTracer in a call. +class AddServiceLabelsFilter : public grpc_core::ChannelFilter { + public: + static const grpc_channel_filter kFilter; + + static absl::StatusOr Create( + const grpc_core::ChannelArgs& args, ChannelFilter::Args /*filter_args*/) { + return AddServiceLabelsFilter( + args.GetPointer>( + GRPC_ARG_LABELS_TO_INJECT)); + } + + grpc_core::ArenaPromise MakeCallPromise( + grpc_core::CallArgs call_args, + grpc_core::NextPromiseFactory next_promise_factory) override { + using CallAttemptTracer = grpc_core::ClientCallTracer::CallAttemptTracer; + auto* call_context = grpc_core::GetContext(); + auto* call_tracer = static_cast( + call_context[GRPC_CONTEXT_CALL_TRACER].value); + EXPECT_NE(call_tracer, nullptr); + call_tracer->AddOptionalLabels( + CallAttemptTracer::OptionalLabelComponent::kXdsServiceLabels, + std::make_shared>( + *labels_to_inject_)); + return next_promise_factory(std::move(call_args)); + } + + private: + explicit AddServiceLabelsFilter( + const std::map* labels_to_inject) + : labels_to_inject_(labels_to_inject) {} + + const std::map* labels_to_inject_; +}; + +const grpc_channel_filter AddServiceLabelsFilter::kFilter = + grpc_core::MakePromiseBasedFilter( + "add_service_labels_filter"); + void OpenTelemetryPluginEnd2EndTest::Init( const absl::flat_hash_set& metric_names, opentelemetry::sdk::resource::Resource resource, std::unique_ptr labels_injector, bool test_no_meter_provider, + const std::map& labels_to_inject, absl::AnyInvocable target_selector, absl::AnyInvocable @@ -83,6 +128,16 @@ void OpenTelemetryPluginEnd2EndTest::Init( ot_builder.AddPluginOption(std::move(option)); } ot_builder.BuildAndRegisterGlobal(); + ChannelArguments channel_args; + if (!labels_to_inject.empty()) { + labels_to_inject_ = labels_to_inject; + grpc_core::CoreConfiguration::RegisterBuilder( + [](grpc_core::CoreConfiguration::Builder* builder) mutable { + builder->channel_init()->RegisterFilter( + GRPC_CLIENT_SUBCHANNEL, &AddServiceLabelsFilter::kFilter); + }); + channel_args.SetPointer(GRPC_ARG_LABELS_TO_INJECT, &labels_to_inject_); + } grpc_init(); grpc::ServerBuilder builder; int port; @@ -96,8 +151,8 @@ void OpenTelemetryPluginEnd2EndTest::Init( server_address_ = absl::StrCat("localhost:", port); canonical_server_address_ = absl::StrCat("dns:///", server_address_); - auto channel = - grpc::CreateChannel(server_address_, grpc::InsecureChannelCredentials()); + auto channel = grpc::CreateCustomChannel( + server_address_, grpc::InsecureChannelCredentials(), channel_args); stub_ = EchoTestService::NewStub(channel); generic_stub_ = std::make_unique(std::move(channel)); } diff --git a/test/cpp/ext/otel/otel_test_library.h b/test/cpp/ext/otel/otel_test_library.h index 629d8c12762..0f6edfd1571 100644 --- a/test/cpp/ext/otel/otel_test_library.h +++ b/test/cpp/ext/otel/otel_test_library.h @@ -63,6 +63,7 @@ class OpenTelemetryPluginEnd2EndTest : public ::testing::Test { opentelemetry::sdk::resource::Resource::Create({}), std::unique_ptr labels_injector = nullptr, bool test_no_meter_provider = false, + const std::map& labels_to_inject = {}, absl::AnyInvocable target_selector = absl::AnyInvocable(), absl::AnyInvocable @@ -94,6 +95,7 @@ class OpenTelemetryPluginEnd2EndTest : public ::testing::Test { const absl::string_view kMethodName = "grpc.testing.EchoTestService/Echo"; const absl::string_view kGenericMethodName = "foo/bar"; + std::map labels_to_inject_; std::shared_ptr reader_; std::string server_address_; std::string canonical_server_address_;