rls: update rls cache metrics to use async gauge framework (#8808)

This PR changes the way rls cache metrics ("grpc.lb.rls.cache_entries" &
"grpc.lb.rls.cache_size") to use async gauge framework instead of
synchronoush capturing.
This also does cleanup to remove the NoopMetricsRecorder struct along
with addressing some old nits.

RELEASE NOTES:
* rls: Change rls cache metrics ("grpc.lb.rls.cache_entries" &
"grpc.lb.rls.cache_size") to be recorded asynchronously (i.e the metric
will be recorded once per collection cycle, rather than every time its
value changes) going forward.

---------

Signed-off-by: Tom Wieczorek <twieczorek@mirantis.com>
Signed-off-by: iamrajiv <rajivperfect007@gmail.com>
Co-authored-by: eshitachandwani <59800922+eshitachandwani@users.noreply.github.com>
Co-authored-by: Arjan Singh Bal <46515553+arjan-bal@users.noreply.github.com>
Co-authored-by: Tom Wieczorek <twz123@users.noreply.github.com>
Co-authored-by: Easwar Swaminathan <easwars@google.com>
Co-authored-by: Sankalp Tripathi <sankalpt92@gmail.com>
Co-authored-by: Joy Bestourous <32601978+joybestourous@users.noreply.github.com>
Co-authored-by: Pranjali-2501 <87357388+Pranjali-2501@users.noreply.github.com>
Co-authored-by: Vanja Pejovic <vvaffle@gmail.com>
Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
Co-authored-by: Yuki Ito <mrno110y@gmail.com>
Co-authored-by: yy <yhymmt37@gmail.com>
Co-authored-by: Evan Jones <evan.jones@datadoghq.com>
Co-authored-by: Antoine Tollenaere <atollena@gmail.com>
Co-authored-by: Rajiv Singh <rajivperfect007@gmail.com>
Co-authored-by: Mikhail Mazurskiy <126021+ash2k@users.noreply.github.com>
Co-authored-by: Jille Timmermans <jille@quis.cx>
diff --git a/balancer/rls/balancer.go b/balancer/rls/balancer.go
index a45e830..abcfcb3 100644
--- a/balancer/rls/balancer.go
+++ b/balancer/rls/balancer.go
@@ -79,14 +79,14 @@
 	dataCachePurgeHook   = func() {}
 	resetBackoffHook     = func() {}
 
-	cacheEntriesMetric = estats.RegisterInt64Gauge(estats.MetricDescriptor{
+	cacheEntriesMetric = estats.RegisterInt64AsyncGauge(estats.MetricDescriptor{
 		Name:        "grpc.lb.rls.cache_entries",
 		Description: "EXPERIMENTAL. Number of entries in the RLS cache.",
 		Unit:        "{entry}",
 		Labels:      []string{"grpc.target", "grpc.lb.rls.server_target", "grpc.lb.rls.instance_uuid"},
 		Default:     false,
 	})
-	cacheSizeMetric = estats.RegisterInt64Gauge(estats.MetricDescriptor{
+	cacheSizeMetric = estats.RegisterInt64AsyncGauge(estats.MetricDescriptor{
 		Name:        "grpc.lb.rls.cache_size",
 		Description: "EXPERIMENTAL. The current size of the RLS cache.",
 		Unit:        "By",
@@ -140,7 +140,9 @@
 		updateCh:           buffer.NewUnbounded(),
 	}
 	lb.logger = internalgrpclog.NewPrefixLogger(logger, fmt.Sprintf("[rls-experimental-lb %p] ", lb))
-	lb.dataCache = newDataCache(maxCacheSize, lb.logger, cc.MetricsRecorder(), opts.Target.String())
+	lb.dataCache = newDataCache(maxCacheSize, lb.logger, opts.Target.String())
+	metricsRecorder := cc.MetricsRecorder()
+	lb.unregisterMetricHandler = metricsRecorder.RegisterAsyncReporter(lb, cacheEntriesMetric, cacheSizeMetric)
 	lb.bg = balancergroup.New(balancergroup.Options{
 		CC:                      cc,
 		BuildOpts:               opts,
@@ -162,6 +164,9 @@
 	dataCachePurgeHook func()
 	logger             *internalgrpclog.PrefixLogger
 
+	// unregisterMetricHandler is the function to deregister the async metric reporter.
+	unregisterMetricHandler func()
+
 	// If both cacheMu and stateMu need to be acquired, the former must be
 	// acquired first to prevent a deadlock. This order restriction is due to the
 	// fact that in places where we need to acquire both the locks, we always
@@ -488,6 +493,7 @@
 	if b.ctrlCh != nil {
 		b.ctrlCh.close()
 	}
+	b.unregisterMetricHandler()
 	b.bg.Close()
 	b.stateMu.Unlock()
 
@@ -702,3 +708,23 @@
 	}
 	b.stateMu.Unlock()
 }
+
+// Report reports the metrics data to the provided recorder.
+func (b *rlsBalancer) Report(r estats.AsyncMetricsRecorder) error {
+	b.cacheMu.Lock()
+	currentSize := b.dataCache.currentSize
+	entriesLen := int64(len(b.dataCache.entries))
+	rlsServerTarget := b.dataCache.rlsServerTarget
+	grpcTarget := b.dataCache.grpcTarget
+	uuid := b.dataCache.uuid
+	shutdown := b.dataCache.shutdown.HasFired()
+	b.cacheMu.Unlock()
+
+	if shutdown {
+		return nil
+	}
+
+	cacheSizeMetric.Record(r, currentSize, grpcTarget, rlsServerTarget, uuid)
+	cacheEntriesMetric.Record(r, entriesLen, grpcTarget, rlsServerTarget, uuid)
+	return nil
+}
diff --git a/balancer/rls/cache.go b/balancer/rls/cache.go
index 7fe796c..2f48d85 100644
--- a/balancer/rls/cache.go
+++ b/balancer/rls/cache.go
@@ -23,7 +23,6 @@
 	"time"
 
 	"github.com/google/uuid"
-	estats "google.golang.org/grpc/experimental/stats"
 	"google.golang.org/grpc/internal/backoff"
 	internalgrpclog "google.golang.org/grpc/internal/grpclog"
 	"google.golang.org/grpc/internal/grpcsync"
@@ -174,21 +173,19 @@
 	rlsServerTarget string
 
 	// Read only after initialization.
-	grpcTarget      string
-	uuid            string
-	metricsRecorder estats.MetricsRecorder
+	grpcTarget string
+	uuid       string
 }
 
-func newDataCache(size int64, logger *internalgrpclog.PrefixLogger, metricsRecorder estats.MetricsRecorder, grpcTarget string) *dataCache {
+func newDataCache(size int64, logger *internalgrpclog.PrefixLogger, grpcTarget string) *dataCache {
 	return &dataCache{
-		maxSize:         size,
-		keys:            newLRU(),
-		entries:         make(map[cacheKey]*cacheEntry),
-		logger:          logger,
-		shutdown:        grpcsync.NewEvent(),
-		grpcTarget:      grpcTarget,
-		uuid:            uuid.New().String(),
-		metricsRecorder: metricsRecorder,
+		maxSize:    size,
+		keys:       newLRU(),
+		entries:    make(map[cacheKey]*cacheEntry),
+		logger:     logger,
+		shutdown:   grpcsync.NewEvent(),
+		grpcTarget: grpcTarget,
+		uuid:       uuid.New().String(),
 	}
 }
 
@@ -327,8 +324,7 @@
 	if dc.currentSize > dc.maxSize {
 		backoffCancelled = dc.resize(dc.maxSize)
 	}
-	cacheSizeMetric.Record(dc.metricsRecorder, dc.currentSize, dc.grpcTarget, dc.rlsServerTarget, dc.uuid)
-	cacheEntriesMetric.Record(dc.metricsRecorder, int64(len(dc.entries)), dc.grpcTarget, dc.rlsServerTarget, dc.uuid)
+
 	return backoffCancelled, true
 }
 
@@ -338,7 +334,7 @@
 	dc.currentSize -= entry.size
 	entry.size = newSize
 	dc.currentSize += entry.size
-	cacheSizeMetric.Record(dc.metricsRecorder, dc.currentSize, dc.grpcTarget, dc.rlsServerTarget, dc.uuid)
+
 }
 
 func (dc *dataCache) getEntry(key cacheKey) *cacheEntry {
@@ -371,8 +367,7 @@
 	delete(dc.entries, key)
 	dc.currentSize -= entry.size
 	dc.keys.removeEntry(key)
-	cacheSizeMetric.Record(dc.metricsRecorder, dc.currentSize, dc.grpcTarget, dc.rlsServerTarget, dc.uuid)
-	cacheEntriesMetric.Record(dc.metricsRecorder, int64(len(dc.entries)), dc.grpcTarget, dc.rlsServerTarget, dc.uuid)
+
 }
 
 func (dc *dataCache) stop() {
diff --git a/balancer/rls/cache_test.go b/balancer/rls/cache_test.go
index 52239de..5a35665 100644
--- a/balancer/rls/cache_test.go
+++ b/balancer/rls/cache_test.go
@@ -25,7 +25,6 @@
 	"github.com/google/go-cmp/cmp"
 	"github.com/google/go-cmp/cmp/cmpopts"
 	"google.golang.org/grpc/internal/backoff"
-	"google.golang.org/grpc/internal/testutils/stats"
 )
 
 var (
@@ -120,7 +119,7 @@
 
 func (s) TestDataCache_BasicOperations(t *testing.T) {
 	initCacheEntries()
-	dc := newDataCache(5, nil, &stats.NoopMetricsRecorder{}, "")
+	dc := newDataCache(5, nil, "")
 	for i, k := range cacheKeys {
 		dc.addEntry(k, cacheEntries[i])
 	}
@@ -134,7 +133,7 @@
 
 func (s) TestDataCache_AddForcesResize(t *testing.T) {
 	initCacheEntries()
-	dc := newDataCache(1, nil, &stats.NoopMetricsRecorder{}, "")
+	dc := newDataCache(1, nil, "")
 
 	// The first entry in cacheEntries has a minimum expiry time in the future.
 	// This entry would stop the resize operation since we do not evict entries
@@ -163,7 +162,7 @@
 
 func (s) TestDataCache_Resize(t *testing.T) {
 	initCacheEntries()
-	dc := newDataCache(5, nil, &stats.NoopMetricsRecorder{}, "")
+	dc := newDataCache(5, nil, "")
 	for i, k := range cacheKeys {
 		dc.addEntry(k, cacheEntries[i])
 	}
@@ -194,7 +193,7 @@
 
 func (s) TestDataCache_EvictExpiredEntries(t *testing.T) {
 	initCacheEntries()
-	dc := newDataCache(5, nil, &stats.NoopMetricsRecorder{}, "")
+	dc := newDataCache(5, nil, "")
 	for i, k := range cacheKeys {
 		dc.addEntry(k, cacheEntries[i])
 	}
@@ -221,7 +220,7 @@
 	}
 
 	initCacheEntries()
-	dc := newDataCache(5, nil, &stats.NoopMetricsRecorder{}, "")
+	dc := newDataCache(5, nil, "")
 	for i, k := range cacheKeys {
 		dc.addEntry(k, cacheEntries[i])
 	}
@@ -242,61 +241,3 @@
 		t.Fatalf("unexpected diff in backoffState for cache entry after dataCache.resetBackoffState(): %s", diff)
 	}
 }
-
-func (s) TestDataCache_Metrics(t *testing.T) {
-	cacheEntriesMetricsTests := []*cacheEntry{
-		{size: 1},
-		{size: 2},
-		{size: 3},
-		{size: 4},
-		{size: 5},
-	}
-	tmr := stats.NewTestMetricsRecorder()
-	dc := newDataCache(50, nil, tmr, "")
-
-	dc.updateRLSServerTarget("rls-server-target")
-	for i, k := range cacheKeys {
-		dc.addEntry(k, cacheEntriesMetricsTests[i])
-	}
-
-	const cacheEntriesKey = "grpc.lb.rls.cache_entries"
-	const cacheSizeKey = "grpc.lb.rls.cache_size"
-	// 5 total entries which add up to 15 size, so should record that.
-	if got, _ := tmr.Metric(cacheEntriesKey); got != 5 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheEntriesKey, got, 5)
-	}
-	if got, _ := tmr.Metric(cacheSizeKey); got != 15 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheSizeKey, got, 15)
-	}
-
-	// Resize down the cache to 2 entries (deterministic as based of LRU).
-	dc.resize(9)
-	if got, _ := tmr.Metric(cacheEntriesKey); got != 2 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheEntriesKey, got, 2)
-	}
-	if got, _ := tmr.Metric(cacheSizeKey); got != 9 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheSizeKey, got, 9)
-	}
-
-	// Update an entry to have size 6. This should reflect in the size metrics,
-	// which will increase by 1 to 11, while the number of cache entries should
-	// stay same. This write is deterministic and writes to the last one.
-	dc.updateEntrySize(cacheEntriesMetricsTests[4], 6)
-
-	if got, _ := tmr.Metric(cacheEntriesKey); got != 2 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheEntriesKey, got, 2)
-	}
-	if got, _ := tmr.Metric(cacheSizeKey); got != 10 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheSizeKey, got, 10)
-	}
-
-	// Delete this scaled up cache key. This should scale down the cache to 1
-	// entries, and remove 6 size so cache size should be 4.
-	dc.deleteAndCleanup(cacheKeys[4], cacheEntriesMetricsTests[4])
-	if got, _ := tmr.Metric(cacheEntriesKey); got != 1 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheEntriesKey, got, 1)
-	}
-	if got, _ := tmr.Metric(cacheSizeKey); got != 4 {
-		t.Fatalf("Unexpected data for metric %v, got: %v, want: %v", cacheSizeKey, got, 4)
-	}
-}
diff --git a/internal/leakcheck/leakcheck_test.go b/internal/leakcheck/leakcheck_test.go
index 35d0a04..c61cc43 100644
--- a/internal/leakcheck/leakcheck_test.go
+++ b/internal/leakcheck/leakcheck_test.go
@@ -134,28 +134,21 @@
 }
 
 func TestLeakChecker_DetectsLeak(t *testing.T) {
-	// 1. Setup the tracker (swaps the delegate in internal).
 	TrackAsyncReporters()
-
 	// Safety defer: ensure we restore the default delegate even if the test crashes
 	// before CheckAsyncReporters is called.
 	defer func() {
 		internal.AsyncReporterCleanupDelegate = func(f func()) func() { return f }
 	}()
 
-	// 2. Simulate a registration using the swapped delegate.
 	// We utilize the internal delegate directly to simulate stats.RegisterAsyncReporter behavior.
 	noOpCleanup := func() {}
 	wrappedCleanup := internal.AsyncReporterCleanupDelegate(noOpCleanup)
 
-	// 3. Create a leak: We discard 'wrappedCleanup' without calling it.
+	// Create a leak: We discard 'wrappedCleanup' without calling it.
 	_ = wrappedCleanup
-
-	// 4. Check for leaks.
 	tl := &testLogger{}
 	CheckAsyncReporters(tl)
-
-	// 5. Assertions.
 	if tl.errorCount == 0 {
 		t.Error("Expected leak checker to report a leak, but it succeeded silently.")
 	}
@@ -165,24 +158,15 @@
 }
 
 func TestLeakChecker_PassesOnCleanup(t *testing.T) {
-	// 1. Setup.
 	TrackAsyncReporters()
 	defer func() {
 		internal.AsyncReporterCleanupDelegate = func(f func()) func() { return f }
 	}()
-
-	// 2. Simulate registration.
 	noOpCleanup := func() {}
 	wrappedCleanup := internal.AsyncReporterCleanupDelegate(noOpCleanup)
-
-	// 3. Behave correctly: Call the cleanup.
 	wrappedCleanup()
-
-	// 4. Check for leaks.
 	tl := &testLogger{}
 	CheckAsyncReporters(tl)
-
-	// 5. Assertions.
 	if tl.errorCount > 0 {
 		t.Errorf("Expected no leaks, but got errors: %v", tl.errors)
 	}
diff --git a/internal/testutils/stats/test_metrics_recorder.go b/internal/testutils/stats/test_metrics_recorder.go
index 5481bce..1c91495 100644
--- a/internal/testutils/stats/test_metrics_recorder.go
+++ b/internal/testutils/stats/test_metrics_recorder.go
@@ -277,6 +277,24 @@
 	r.data[handle.Name] = float64(incr)
 }
 
+// To implement a estats.AsyncMetricsRecorder, which allows it to be used in async metrics:
+
+// RecordInt64AsyncGauge sends the metrics data to the intGaugeCh channel and updates
+// the internal data map with the recorded value.
+func (r *TestMetricsRecorder) RecordInt64AsyncGauge(handle *estats.Int64AsyncGaugeHandle, incr int64, labels ...string) {
+	r.intGaugeCh.ReceiveOrFail()
+	r.intGaugeCh.Send(MetricsData{
+		Handle:    handle.Descriptor(),
+		IntIncr:   incr,
+		LabelKeys: append(handle.Labels, handle.OptionalLabels...),
+		LabelVals: labels,
+	})
+
+	r.mu.Lock()
+	defer r.mu.Unlock()
+	r.data[handle.Name] = float64(incr)
+}
+
 // To implement a stats.Handler, which allows it to be set as a dial option:
 
 // TagRPC is TestMetricsRecorder's implementation of TagRPC.
@@ -294,9 +312,3 @@
 
 // HandleConn is TestMetricsRecorder's implementation of HandleConn.
 func (r *TestMetricsRecorder) HandleConn(context.Context, stats.ConnStats) {}
-
-// NoopMetricsRecorder is a noop MetricsRecorder to be used in tests to prevent
-// nil panics.
-type NoopMetricsRecorder struct {
-	estats.UnimplementedMetricsRecorder
-}
diff --git a/internal/xds/xdsclient/client_refcounted_test.go b/internal/xds/xdsclient/client_refcounted_test.go
index ba2331c..fe44e6e 100644
--- a/internal/xds/xdsclient/client_refcounted_test.go
+++ b/internal/xds/xdsclient/client_refcounted_test.go
@@ -24,8 +24,8 @@
 	"testing"
 
 	"github.com/google/uuid"
+	estats "google.golang.org/grpc/experimental/stats"
 	"google.golang.org/grpc/internal/testutils"
-	"google.golang.org/grpc/internal/testutils/stats"
 	"google.golang.org/grpc/internal/testutils/xds/e2e"
 	"google.golang.org/grpc/internal/xds/bootstrap"
 )
@@ -61,7 +61,7 @@
 	defer func() { xdsClientImplCloseHook = origClientImplCloseHook }()
 
 	// The first call to New() should create a new client.
-	_, closeFunc, err := pool.NewClient(t.Name(), &stats.NoopMetricsRecorder{})
+	_, closeFunc, err := pool.NewClient(t.Name(), &estats.UnimplementedMetricsRecorder{})
 	if err != nil {
 		t.Fatalf("Failed to create xDS client: %v", err)
 	}
@@ -77,7 +77,7 @@
 	closeFuncs := make([]func(), count)
 	for i := 0; i < count; i++ {
 		func() {
-			_, closeFuncs[i], err = pool.NewClient(t.Name(), &stats.NoopMetricsRecorder{})
+			_, closeFuncs[i], err = pool.NewClient(t.Name(), &estats.UnimplementedMetricsRecorder{})
 			if err != nil {
 				t.Fatalf("%d-th call to New() failed with error: %v", i, err)
 			}
@@ -115,7 +115,7 @@
 
 	// Calling New() again, after the previous Client was actually closed,
 	// should create a new one.
-	_, closeFunc, err = pool.NewClient(t.Name(), &stats.NoopMetricsRecorder{})
+	_, closeFunc, err = pool.NewClient(t.Name(), &estats.UnimplementedMetricsRecorder{})
 	if err != nil {
 		t.Fatalf("Failed to create xDS client: %v", err)
 	}
@@ -157,7 +157,7 @@
 
 	// Create two xDS clients.
 	client1Name := t.Name() + "-1"
-	_, closeFunc1, err := pool.NewClient(client1Name, &stats.NoopMetricsRecorder{})
+	_, closeFunc1, err := pool.NewClient(client1Name, &estats.UnimplementedMetricsRecorder{})
 	if err != nil {
 		t.Fatalf("Failed to create xDS client: %v", err)
 	}
@@ -172,7 +172,7 @@
 	}
 
 	client2Name := t.Name() + "-2"
-	_, closeFunc2, err := pool.NewClient(client2Name, &stats.NoopMetricsRecorder{})
+	_, closeFunc2, err := pool.NewClient(client2Name, &estats.UnimplementedMetricsRecorder{})
 	if err != nil {
 		t.Fatalf("Failed to create xDS client: %v", err)
 	}
@@ -194,7 +194,7 @@
 		defer wg.Done()
 		for i := 0; i < count; i++ {
 			var err error
-			_, closeFuncs1[i], err = pool.NewClient(client1Name, &stats.NoopMetricsRecorder{})
+			_, closeFuncs1[i], err = pool.NewClient(client1Name, &estats.UnimplementedMetricsRecorder{})
 			if err != nil {
 				t.Errorf("%d-th call to New() failed with error: %v", i, err)
 			}
@@ -204,7 +204,7 @@
 		defer wg.Done()
 		for i := 0; i < count; i++ {
 			var err error
-			_, closeFuncs2[i], err = pool.NewClient(client2Name, &stats.NoopMetricsRecorder{})
+			_, closeFuncs2[i], err = pool.NewClient(client2Name, &estats.UnimplementedMetricsRecorder{})
 			if err != nil {
 				t.Errorf("%d-th call to New() failed with error: %v", i, err)
 			}
diff --git a/internal/xds/xdsclient/pool/pool_ext_test.go b/internal/xds/xdsclient/pool/pool_ext_test.go
index 2150a27..83f8a99 100644
--- a/internal/xds/xdsclient/pool/pool_ext_test.go
+++ b/internal/xds/xdsclient/pool/pool_ext_test.go
@@ -31,7 +31,6 @@
 	"google.golang.org/grpc/internal/grpctest"
 	"google.golang.org/grpc/internal/stubserver"
 	"google.golang.org/grpc/internal/testutils"
-	"google.golang.org/grpc/internal/testutils/stats"
 	"google.golang.org/grpc/internal/testutils/xds/e2e"
 	"google.golang.org/grpc/internal/xds/bootstrap"
 	"google.golang.org/grpc/internal/xds/xdsclient"
@@ -41,6 +40,7 @@
 	v3endpointpb "github.com/envoyproxy/go-control-plane/envoy/config/endpoint/v3"
 	v3listenerpb "github.com/envoyproxy/go-control-plane/envoy/config/listener/v3"
 	v3routepb "github.com/envoyproxy/go-control-plane/envoy/config/route/v3"
+	estats "google.golang.org/grpc/experimental/stats"
 	testgrpc "google.golang.org/grpc/interop/grpc_testing"
 	testpb "google.golang.org/grpc/interop/grpc_testing"
 )
@@ -65,7 +65,7 @@
 // Then it sets the env var XDSBootstrapFileName and retry creating a client
 // in DefaultPool. This should succeed.
 func (s) TestDefaultPool_LazyLoadBootstrapConfig(t *testing.T) {
-	_, closeFunc, err := xdsclient.DefaultPool.NewClient(t.Name(), &stats.NoopMetricsRecorder{})
+	_, closeFunc, err := xdsclient.DefaultPool.NewClient(t.Name(), &estats.UnimplementedMetricsRecorder{})
 	if err == nil {
 		t.Fatalf("xdsclient.DefaultPool.NewClient() succeeded without setting bootstrap config env vars, want failure")
 	}
@@ -93,7 +93,7 @@
 	// state to make it re-read the env vars during next client creation.
 	xdsclient.DefaultPool.UnsetBootstrapConfigForTesting()
 
-	_, closeFunc, err = xdsclient.DefaultPool.NewClient(t.Name(), &stats.NoopMetricsRecorder{})
+	_, closeFunc, err = xdsclient.DefaultPool.NewClient(t.Name(), &estats.UnimplementedMetricsRecorder{})
 	if err != nil {
 		t.Fatalf("Failed to create xDS client: %v", err)
 	}