From 072276bacc934185074a5ceb65f269b0d8a2cbca Mon Sep 17 00:00:00 2001 From: Markus Wieland Date: Wed, 12 Aug 2026 13:47:52 +0200 Subject: [PATCH] feat: add monitor for caching client Signed-off-by: Markus Wieland --- cmd/manager/main.go | 4 +- pkg/clientcache/client.go | 30 +++-- pkg/clientcache/client_test.go | 14 +-- pkg/clientcache/monitor.go | 103 +++++++++++++++++ pkg/clientcache/monitor_test.go | 190 ++++++++++++++++++++++++++++++++ 5 files changed, 326 insertions(+), 15 deletions(-) create mode 100644 pkg/clientcache/monitor.go create mode 100644 pkg/clientcache/monitor_test.go diff --git a/cmd/manager/main.go b/cmd/manager/main.go index ab888d37b..16708b713 100644 --- a/cmd/manager/main.go +++ b/cmd/manager/main.go @@ -403,7 +403,8 @@ func main() { // observed. *multicluster.Client satisfies clientcache.Client, providing // both the inner client.Client and informer access for eviction. clientCacheConfig := conf.GetConfigOrDie[clientcache.RootConfig]() - cachingClient, err := clientcache.New(multiclusterClient, scheme, clientCacheConfig.ClientCache) + clientCacheMonitor := clientcache.NewMonitor("cortex_") + cachingClient, err := clientcache.New(multiclusterClient, scheme, clientCacheConfig.ClientCache, clientCacheMonitor) if err != nil { setupLog.Error(err, "unable to create client cache") os.Exit(1) @@ -419,6 +420,7 @@ func main() { metrics.Registry = monitoring.WrapRegistry(metrics.Registry, metricsConfig) metrics.Registry.MustRegister(&logMetricsMonitor) metrics.Registry.MustRegister(multiclusterMonitor) + metrics.Registry.MustRegister(clientCacheMonitor) // TODO: Remove me after scheduling pipeline steps don't require DB connections anymore. metrics.Registry.MustRegister(&db.Monitor) diff --git a/pkg/clientcache/client.go b/pkg/clientcache/client.go index c80fde968..33261e17c 100644 --- a/pkg/clientcache/client.go +++ b/pkg/clientcache/client.go @@ -152,10 +152,11 @@ const defaultTTL = 2 * time.Minute type CachingClient struct { client.Client // inner client, used for delegation - inner Client - scheme *runtime.Scheme - ttl time.Duration - gvks map[schema.GroupVersionKind]bool + inner Client + scheme *runtime.Scheme + ttl time.Duration + gvks map[schema.GroupVersionKind]bool + monitor Monitor mu sync.RWMutex byGVK map[schema.GroupVersionKind]map[objectKey]*entry @@ -170,8 +171,8 @@ type CachingClient struct { // New builds a CachingClient wrapping inner. informers supplies the informers // used for eviction, scheme resolves object GVKs, and conf lists the GVKs to // overlay and the TTL. GVK strings are formatted as "//" -// and are resolved against scheme. -func New(inner Client, scheme *runtime.Scheme, conf Config) (*CachingClient, error) { +// and are resolved against scheme. A nil mon disables metric recording. +func New(inner Client, scheme *runtime.Scheme, conf Config, mon Monitor) (*CachingClient, error) { gvks, err := resolveGVKs(scheme, conf.GVKs) if err != nil { return nil, err @@ -186,6 +187,7 @@ func New(inner Client, scheme *runtime.Scheme, conf Config) (*CachingClient, err scheme: scheme, ttl: ttl, gvks: gvks, + monitor: mon, byGVK: make(map[schema.GroupVersionKind]map[objectKey]*entry), indexers: make(map[schema.GroupVersionKind]map[string]client.IndexerFunc), writeLocks: newKeyedMutex(), @@ -256,6 +258,7 @@ func (c *CachingClient) upsert(gvk schema.GroupVersionKind, obj client.Object) { deleted: false, expiresAt: time.Now().Add(c.ttl), } + c.recordSizeLocked(gvk) } // tombstone marks the object as deleted in the overlay so it is filtered out @@ -271,6 +274,7 @@ func (c *CachingClient) tombstone(gvk schema.GroupVersionKind, obj client.Object deleted: true, expiresAt: time.Now().Add(c.ttl), } + c.recordSizeLocked(gvk) } // evictIfSeen removes the overlay entry for obj if the informer-observed object @@ -297,6 +301,7 @@ func (c *CachingClient) evictIfSeen(gvk schema.GroupVersionKind, obj client.Obje return } delete(entries, key) + c.recordSizeLocked(gvk) } // getEntry returns the overlay entry for the key, if present. @@ -315,12 +320,13 @@ func (c *CachingClient) getEntry(gvk schema.GroupVersionKind, key objectKey) (*e func (c *CachingClient) cleanupExpired(now time.Time) { c.mu.Lock() defer c.mu.Unlock() - for _, entries := range c.byGVK { + for gvk, entries := range c.byGVK { for key, e := range entries { if now.After(e.expiresAt) { delete(entries, key) } } + c.recordSizeLocked(gvk) } } @@ -342,6 +348,15 @@ func (c *CachingClient) ensureGVK(gvk schema.GroupVersionKind) { } } +// recordSizeLocked reports the current overlay size (live + tombstones) for the +// GVK to the monitor, if one is configured. Callers must hold c.mu. +func (c *CachingClient) recordSizeLocked(gvk schema.GroupVersionKind) { + if c.monitor == nil { + return + } + c.monitor.observe(gvk, len(c.byGVK[gvk])) +} + // overlayList merges the overlay entries for the GVK into the informer result, // deduplicating by objectKey (overlay wins), dropping tombstones, and filtering // overlay-only entries against the list options' label and field selectors. @@ -533,6 +548,7 @@ func (c *CachingClient) DeleteAllOf(ctx context.Context, obj client.Object, opts expiresAt: time.Now().Add(c.ttl), } } + c.recordSizeLocked(gvk) return nil } diff --git a/pkg/clientcache/client_test.go b/pkg/clientcache/client_test.go index 75aa3db9e..cfec38cb2 100644 --- a/pkg/clientcache/client_test.go +++ b/pkg/clientcache/client_test.go @@ -224,7 +224,7 @@ func reservationConfig() Config { func newCaching(t *testing.T, inner Client) *CachingClient { t.Helper() - c, err := New(inner, testScheme(t), reservationConfig()) + c, err := New(inner, testScheme(t), reservationConfig(), nil) if err != nil { t.Fatalf("New: %v", err) } @@ -261,7 +261,7 @@ func waitFor(t *testing.T, cond func() bool) { func TestNewUnknownGVKError(t *testing.T) { _, err := New(newTestClient(t), testScheme(t), Config{ GVKs: []string{"cortex.cloud/v1alpha1/DoesNotExist"}, - }) + }, nil) if err == nil { t.Fatalf("expected error for unknown GVK, got nil") } @@ -270,7 +270,7 @@ func TestNewUnknownGVKError(t *testing.T) { func TestNewDefaultTTL(t *testing.T) { c, err := New(newTestClient(t), testScheme(t), Config{ GVKs: []string{"cortex.cloud/v1alpha1/Reservation"}, - }) + }, nil) if err != nil { t.Fatalf("New: %v", err) } @@ -283,7 +283,7 @@ func TestNewExplicitTTL(t *testing.T) { c, err := New(newTestClient(t), testScheme(t), Config{ GVKs: []string{"cortex.cloud/v1alpha1/Reservation"}, TTL: metav1.Duration{Duration: 90 * time.Second}, - }) + }, nil) if err != nil { t.Fatalf("New: %v", err) } @@ -548,7 +548,7 @@ func TestFieldMatching(t *testing.T) { func TestNonCachedGVKPassthrough(t *testing.T) { inner := newTestClient(t) - c, err := New(inner, testScheme(t), Config{}) + c, err := New(inner, testScheme(t), Config{}, nil) if err != nil { t.Fatalf("New: %v", err) } @@ -754,7 +754,7 @@ func TestGetNotFoundWithNoOverlay(t *testing.T) { func TestGetNonCachedPropagatesError(t *testing.T) { sentinel := errors.New("get boom") - c, err := New(&errClient{Client: newTestClient(t), getErr: sentinel}, testScheme(t), Config{}) + c, err := New(&errClient{Client: newTestClient(t), getErr: sentinel}, testScheme(t), Config{}, nil) if err != nil { t.Fatalf("New: %v", err) } @@ -802,7 +802,7 @@ func TestStatusCreateDelegates(t *testing.T) { func TestStatusUpdateNonCachedNoOverlay(t *testing.T) { r := newReservation("res-sn", "az-1", "") inner := newTestClient(t, r) - c, err := New(inner, testScheme(t), Config{}) + c, err := New(inner, testScheme(t), Config{}, nil) if err != nil { t.Fatalf("New: %v", err) } diff --git a/pkg/clientcache/monitor.go b/pkg/clientcache/monitor.go new file mode 100644 index 000000000..37e3b4cc1 --- /dev/null +++ b/pkg/clientcache/monitor.go @@ -0,0 +1,103 @@ +// Copyright SAP SE +// SPDX-License-Identifier: Apache-2.0 + +package clientcache + +import ( + "sync" + + "github.com/prometheus/client_golang/prometheus" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// Monitor is the metrics sink for the CachingClient. It is optional on the +// client: a nil Monitor causes recording to be skipped entirely. It embeds +// prometheus.Collector so a concrete implementation can be registered with a +// Prometheus registry. +type Monitor interface { + prometheus.Collector + + // observe records the current overlay size (live + tombstones) for a GVK. + // It is called by the client while holding c.mu so size is consistent. + observe(gvk schema.GroupVersionKind, size int) +} + +// monitor is the default Prometheus-backed Monitor implementation. +// +// Overlay entries are normally evicted within milliseconds once the informer +// observes the change, so a plain gauge scraped by Prometheus would almost +// always read 0 and give no signal about transient spikes. To make abnormal +// growth visible even when the overlay is empty at scrape time, the monitor +// tracks a per-GVK high-watermark (the max size seen since the last scrape) +// which is reset on each Collect, plus a current-size gauge to confirm the +// overlay drains correctly. +type monitor struct { + mu sync.Mutex + // maxSinceScrape is the high-watermark of overlay size per GVK since the + // last scrape. It is reset to a fresh map on each Collect. + maxSinceScrape map[schema.GroupVersionKind]int + // current is the last observed overlay size per GVK. It is a snapshot, not + // reset on scrape. + current map[schema.GroupVersionKind]int + + maxDesc *prometheus.Desc + currentDesc *prometheus.Desc +} + +// NewMonitor creates a new Prometheus-backed CachingClient monitor. The prefix +// is prepended to every metric name (e.g. pass "cortex_" to produce +// "cortex_clientcache_overlay_entries_max"). +func NewMonitor(prefix string) Monitor { + return &monitor{ + maxSinceScrape: make(map[schema.GroupVersionKind]int), + current: make(map[schema.GroupVersionKind]int), + maxDesc: prometheus.NewDesc( + prefix+"clientcache_overlay_entries_max", + "Maximum overlay entries (live + tombstones) per GVK since the last scrape", + []string{"gvk"}, nil, + ), + currentDesc: prometheus.NewDesc( + prefix+"clientcache_overlay_entries", + "Current overlay entries (live + tombstones) per GVK at scrape time", + []string{"gvk"}, nil, + ), + } +} + +// observe records the current overlay size for a GVK, updating the +// high-watermark if it grew. +func (m *monitor) observe(gvk schema.GroupVersionKind, size int) { + m.mu.Lock() + defer m.mu.Unlock() + m.current[gvk] = size + if size > m.maxSinceScrape[gvk] { + m.maxSinceScrape[gvk] = size + } +} + +// Describe implements prometheus.Collector. +func (m *monitor) Describe(ch chan<- *prometheus.Desc) { + ch <- m.maxDesc + ch <- m.currentDesc +} + +// Collect implements prometheus.Collector. It emits the current size and the +// high-watermark per GVK, then resets the high-watermark so the next scrape +// captures a fresh maximum. +func (m *monitor) Collect(ch chan<- prometheus.Metric) { + m.mu.Lock() + defer m.mu.Unlock() + for gvk, size := range m.current { + ch <- prometheus.MustNewConstMetric( + m.currentDesc, prometheus.GaugeValue, float64(size), gvk.String(), + ) + } + for gvk, max := range m.maxSinceScrape { + ch <- prometheus.MustNewConstMetric( + m.maxDesc, prometheus.GaugeValue, float64(max), gvk.String(), + ) + } + // Reset the high-watermark after emitting so the next scrape window starts + // fresh. current is intentionally left as-is (it is a snapshot). + m.maxSinceScrape = make(map[schema.GroupVersionKind]int) +} diff --git a/pkg/clientcache/monitor_test.go b/pkg/clientcache/monitor_test.go new file mode 100644 index 000000000..146305acb --- /dev/null +++ b/pkg/clientcache/monitor_test.go @@ -0,0 +1,190 @@ +// Copyright SAP SE +// SPDX-License-Identifier: Apache-2.0 + +package clientcache + +import ( + "context" + "testing" + + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// scrape performs a single registry Gather (which triggers one Collect, and +// thus one high-watermark reset) and returns the gauge values keyed by +// "|". +func scrape(t *testing.T, reg *prometheus.Registry) map[string]float64 { + t.Helper() + families, err := reg.Gather() + if err != nil { + t.Fatalf("gather: %v", err) + } + out := make(map[string]float64) + for _, f := range families { + for _, m := range f.GetMetric() { + out[f.GetName()+"|"+labelValue(m, "gvk")] = m.GetGauge().GetValue() + } + } + return out +} + +func labelValue(m *dto.Metric, name string) string { + for _, l := range m.GetLabel() { + if l.GetName() == name { + return l.GetValue() + } + } + return "" +} + +func TestMonitor_Registration(t *testing.T) { + m := NewMonitor("cortex_") + reg := prometheus.NewRegistry() + if err := reg.Register(m); err != nil { + t.Fatalf("register: %v", err) + } + + gvk := schema.GroupVersionKind{Group: "cortex.cloud", Version: "v1alpha1", Kind: "Reservation"} + m.observe(gvk, 2) + + families, err := reg.Gather() + if err != nil { + t.Fatalf("gather: %v", err) + } + var maxFound, curFound bool + for _, f := range families { + switch f.GetName() { + case "cortex_clientcache_overlay_entries_max": + maxFound = true + case "cortex_clientcache_overlay_entries": + curFound = true + } + } + if !maxFound { + t.Error("expected cortex_clientcache_overlay_entries_max to be registered") + } + if !curFound { + t.Error("expected cortex_clientcache_overlay_entries to be registered") + } +} + +func TestMonitor_HighWatermarkAndReset(t *testing.T) { + m := NewMonitor("cortex_") + reg := prometheus.NewRegistry() + if err := reg.Register(m); err != nil { + t.Fatalf("register: %v", err) + } + + gvk := schema.GroupVersionKind{Group: "cortex.cloud", Version: "v1alpha1", Kind: "Reservation"} + otherGVK := schema.GroupVersionKind{Group: "kvm.cloud.sap", Version: "v1", Kind: "Hypervisor"} + + // Observe a spike (max 3) that drains back to 2. + m.observe(gvk, 1) + m.observe(gvk, 3) + m.observe(gvk, 2) + // A second GVK stays independent. + m.observe(otherGVK, 5) + + // First scrape: max reflects the spike, current the last observed size. + const maxName = "cortex_clientcache_overlay_entries_max" + const curName = "cortex_clientcache_overlay_entries" + s := scrape(t, reg) + if got := s[maxName+"|"+gvk.String()]; got != 3 { + t.Errorf("first scrape max for %s: got %v, want 3", gvk, got) + } + if got := s[curName+"|"+gvk.String()]; got != 2 { + t.Errorf("first scrape current for %s: got %v, want 2", gvk, got) + } + if got := s[maxName+"|"+otherGVK.String()]; got != 5 { + t.Errorf("first scrape max for %s: got %v, want 5", otherGVK, got) + } + + // Second scrape without further observation: the high-watermark was reset + // by the previous Collect, so no _max sample remains for the gvk. current + // is a snapshot and still reports the last size. + s = scrape(t, reg) + if _, ok := s[maxName+"|"+gvk.String()]; ok { + t.Errorf("second scrape max for %s: expected reset (absent), got a sample", gvk) + } + if got := s[curName+"|"+gvk.String()]; got != 2 { + t.Errorf("second scrape current for %s: got %v, want 2", gvk, got) + } + + // Observing again after reset with size 0 does not create a _max sample + // (observe only records a watermark when size grows above the current max, + // and 0 never exceeds the zero-valued default). The current gauge reflects + // the new size of 0. + m.observe(gvk, 0) + s = scrape(t, reg) + if _, ok := s[maxName+"|"+gvk.String()]; ok { + t.Errorf("post-reset max for %s: expected absent for size 0, got a sample", gvk) + } + if got, ok := s[curName+"|"+gvk.String()]; !ok || got != 0 { + t.Errorf("post-reset current for %s: got %v (ok=%v), want 0", gvk, got, ok) + } + + // Observing a non-zero size after reset re-establishes the watermark. + m.observe(gvk, 4) + s = scrape(t, reg) + if got := s[maxName+"|"+gvk.String()]; got != 4 { + t.Errorf("post-reset max for %s: got %v, want 4", gvk, got) + } +} + +// TestMonitor_ClientRecords exercises the client-level recording path with a +// real monitor: creating and deleting a cached object moves the overlay size +// and the high-watermark tracks the peak. +func TestMonitor_ClientRecords(t *testing.T) { + inner := newTestClient(t) + m := NewMonitor("cortex_") + reg := prometheus.NewRegistry() + if err := reg.Register(m); err != nil { + t.Fatalf("register: %v", err) + } + c, err := New(inner, testScheme(t), reservationConfig(), m) + if err != nil { + t.Fatalf("New: %v", err) + } + + ctx := context.Background() + if err := c.Create(ctx, newReservation("a", "az1", "")); err != nil { + t.Fatalf("Create a: %v", err) + } + if err := c.Create(ctx, newReservation("b", "az1", "")); err != nil { + t.Fatalf("Create b: %v", err) + } + // Delete tombstones stay in the overlay, so the size does not shrink here. + if err := c.Delete(ctx, newReservation("a", "az1", "")); err != nil { + t.Fatalf("Delete a: %v", err) + } + + gvk := reservationGVK().String() + s := scrape(t, reg) + if got := s["cortex_clientcache_overlay_entries_max|"+gvk]; got != 2 { + t.Errorf("max: got %v, want 2", got) + } + if got := s["cortex_clientcache_overlay_entries|"+gvk]; got != 2 { + t.Errorf("current: got %v, want 2", got) + } +} + +// TestMonitor_NilSafe verifies a client built with a nil monitor does not panic +// on the recording paths. +func TestMonitor_NilSafe(t *testing.T) { + inner := newTestClient(t) + c, err := New(inner, testScheme(t), reservationConfig(), nil) + if err != nil { + t.Fatalf("New: %v", err) + } + ctx := context.Background() + res := newReservation("a", "az1", "") + if err := c.Create(ctx, res); err != nil { + t.Fatalf("Create: %v", err) + } + if err := c.Delete(ctx, res); err != nil { + t.Fatalf("Delete: %v", err) + } + c.cleanupExpired(res.CreationTimestamp.Time) +}