diff --git a/internal/xds/balancer/clusterimpl/clusterimpl.go b/internal/xds/balancer/clusterimpl/clusterimpl.go index 72d143a3586b..cdc7bf2eb579 100644 --- a/internal/xds/balancer/clusterimpl/clusterimpl.go +++ b/internal/xds/balancer/clusterimpl/clusterimpl.go @@ -163,6 +163,7 @@ func (b *clusterImplBalancer) newPickerLocked() *picker { counter: b.requestCounter, countMax: b.requestCountMax, telemetryLabels: b.telemetryLabels, + clusterName: b.clusterName, } } @@ -414,6 +415,7 @@ func (b *clusterImplBalancer) getClusterName() string { // SubConn to the wrapper for this purpose. type scWrapper struct { balancer.SubConn + // locality needs to be atomic because it can be updated while being read by // the picker. locality atomic.Pointer[clients.Locality] diff --git a/internal/xds/balancer/clusterimpl/picker.go b/internal/xds/balancer/clusterimpl/picker.go index fab09fa0978b..0c03326178c3 100644 --- a/internal/xds/balancer/clusterimpl/picker.go +++ b/internal/xds/balancer/clusterimpl/picker.go @@ -20,6 +20,7 @@ package clusterimpl import ( "context" + "maps" v3orcapb "github.com/cncf/xds/go/xds/data/orca/v3" "google.golang.org/grpc/balancer" @@ -87,6 +88,7 @@ type picker struct { counter *xdsclient.ClusterRequestsCounter countMax uint32 telemetryLabels map[string]string + clusterName string } func telemetryLabels(ctx context.Context) map[string]string { @@ -103,10 +105,9 @@ func telemetryLabels(ctx context.Context) map[string]string { func (d *picker) Pick(info balancer.PickInfo) (balancer.PickResult, error) { // Unconditionally set labels if present, even dropped or queued RPC's can // use these labels. - if labels := telemetryLabels(info.Ctx); labels != nil { - for key, value := range d.telemetryLabels { - labels[key] = value - } + labels := telemetryLabels(info.Ctx) + if labels != nil { + maps.Copy(labels, d.telemetryLabels) } // Don't drop unless the inner picker is READY. Similar to @@ -154,8 +155,9 @@ func (d *picker) Pick(info balancer.PickInfo) (balancer.PickResult, error) { return pr, err } - if labels := telemetryLabels(info.Ctx); labels != nil { + if labels != nil { labels["grpc.lb.locality"] = xdsinternal.LocalityString(lID) + labels["grpc.lb.backend_service"] = d.clusterName } if d.loadStore != nil { diff --git a/stats/opentelemetry/client_metrics.go b/stats/opentelemetry/client_metrics.go index d2d80494d4ea..c3b77da5bf76 100644 --- a/stats/opentelemetry/client_metrics.go +++ b/stats/opentelemetry/client_metrics.go @@ -174,7 +174,8 @@ func (h *clientMetricsHandler) TagRPC(ctx context.Context, info *stats.RPCTagInf // executes on the callpath that this OpenTelemetry component // currently supports. TelemetryLabels: map[string]string{ - "grpc.lb.locality": "", + "grpc.lb.locality": "", + "grpc.lb.backend_service": "", }, } ctx = istats.SetLabels(ctx, labels) diff --git a/test/xds/xds_telemetry_labels_test.go b/test/xds/xds_telemetry_labels_test.go index 3ed589e5e5b6..36187419b850 100644 --- a/test/xds/xds_telemetry_labels_test.go +++ b/test/xds/xds_telemetry_labels_test.go @@ -44,7 +44,8 @@ const serviceNamespaceKey = "service_namespace" const serviceNamespaceKeyCSM = "csm.service_namespace_name" const serviceNameValue = "grpc-service" const serviceNamespaceValue = "grpc-service-namespace" - +const backendServiceKey = "grpc.lb.backend_service" +const backendServiceValue = "cluster-my-service-client-side-xds" const localityKey = "grpc.lb.locality" const localityValue = `{region="region-1", zone="zone-1", sub_zone="subzone-1"}` @@ -124,23 +125,23 @@ func (fsh *fakeStatsHandler) TagRPC(ctx context.Context, _ *stats.RPCTagInfo) co func (fsh *fakeStatsHandler) HandleRPC(_ context.Context, rs stats.RPCStats) { switch rs.(type) { - // stats.Begin won't get Telemetry Labels because happens after picker - // picks. - - // These three stats callouts trigger all metrics for OpenTelemetry that - // aren't started. All of these should have access to the desired telemetry + // stats.Begin is called before the picker runs, so it won't have telemetry // labels. + // The following three stats callouts trigger OpenTelemetry metrics and are + // guaranteed to run after the picker has selected a subchannel. Therefore, + // they should have access to the desired telemetry labels. case *stats.OutPayload, *stats.InPayload, *stats.End: want := map[string]string{ serviceNameKeyCSM: serviceNameValue, serviceNamespaceKeyCSM: serviceNamespaceValue, localityKey: localityValue, + backendServiceKey: backendServiceValue, } if diff := cmp.Diff(fsh.labels.TelemetryLabels, want); diff != "" { fsh.t.Fatalf("fsh.labels.TelemetryLabels (-got +want): %v", diff) } default: // Nothing to assert for the other stats.Handler callouts. - } + } }