From b4405a9903fa71e6a2bd19bee36700aa306f8bb3 Mon Sep 17 00:00:00 2001 From: YuMin Kim Date: Thu, 25 Jun 2026 12:50:16 +0900 Subject: [PATCH 1/2] Add lifecycle endpoint to clear metrics Signed-off-by: YuMin Kim --- CHANGELOG.md | 4 +++ README.md | 4 +-- main.go | 15 ++++++++ main_test.go | 50 ++++++++++++++++++++++++++ pkg/exporter/exporter.go | 12 +++++++ pkg/exporter/exporter_test.go | 68 +++++++++++++++++++++++++++++++++++ pkg/registry/registry.go | 25 +++++++++++-- 7 files changed, 173 insertions(+), 5 deletions(-) create mode 100644 main_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index f4c40995..796f8e9d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,7 @@ +## unreleased + +* [FEATURE] Add a lifecycle `/-/clear` endpoint to clear registered StatsD time series without restarting the exporter ([#714](https://github.com/prometheus/statsd_exporter/pull/714)) + ## 0.30.0 / 2026-05-28 * [CHANGE] Remove the Dockerfile `HEALTHCHECK` from published container images ([#671](https://github.com/prometheus/statsd_exporter/pull/671)) * [ENHANCEMENT] Add a distroless container image variant ([#703](https://github.com/prometheus/statsd_exporter/pull/703)) diff --git a/README.md b/README.md index 5c919353..1ec72a44 100644 --- a/README.md +++ b/README.md @@ -113,8 +113,8 @@ NOTE: Version 0.7.0 switched to the [kingpin](https://github.com/alecthomas/king ## Lifecycle API -The `statsd_exporter` has an optional lifecycle API (disabled by default) that can be used to reload or quit the exporter -by sending a `PUT` or `POST` request to the `/-/reload` or `/-/quit` endpoints. +The `statsd_exporter` has an optional lifecycle API (disabled by default) that can be used to reload, quit, or clear the exporter +by sending a `PUT` or `POST` request to the `/-/reload`, `/-/quit`, or `/-/clear` endpoints. ## Relay diff --git a/main.go b/main.go index 2f139e61..13e3960e 100644 --- a/main.go +++ b/main.go @@ -205,6 +205,20 @@ func reloadConfig(fileName string, mapper *mapper.MetricMapper, logger *slog.Log } } +type metricsClearer interface { + ClearMetrics() int +} + +func clearMetricsHandler(clearer metricsClearer, logger *slog.Logger) http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPut || r.Method == http.MethodPost { + cleared := clearer.ClearMetrics() + logger.Info("Received lifecycle api clear", "metrics", cleared) + fmt.Fprintf(w, "Cleared %d metric series", cleared) + } + } +} + func dumpFSM(mapper *mapper.MetricMapper, dumpFilename string, logger *slog.Logger) error { f, err := os.Create(dumpFilename) if err != nil { @@ -528,6 +542,7 @@ func main() { quitChan <- struct{}{} } }) + mux.HandleFunc("/-/clear", clearMetricsHandler(exporter, logger)) } mux.HandleFunc("/-/healthy", func(w http.ResponseWriter, r *http.Request) { diff --git a/main_test.go b/main_test.go new file mode 100644 index 00000000..b88df8db --- /dev/null +++ b/main_test.go @@ -0,0 +1,50 @@ +// Copyright 2026 The Prometheus 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. + +package main + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/prometheus/common/promslog" +) + +type fakeMetricsClearer struct { + cleared int + called int +} + +func (f *fakeMetricsClearer) ClearMetrics() int { + f.called++ + return f.cleared +} + +func TestClearMetricsHandler(t *testing.T) { + clearer := &fakeMetricsClearer{cleared: 3} + handler := clearMetricsHandler(clearer, promslog.NewNopLogger()) + + request := httptest.NewRequest(http.MethodPost, "/-/clear", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if clearer.called != 1 { + t.Fatalf("expected clearer to be called once, got %d", clearer.called) + } + if body := response.Body.String(); !strings.Contains(body, "Cleared 3 metric series") { + t.Fatalf("unexpected response body: %q", body) + } +} diff --git a/pkg/exporter/exporter.go b/pkg/exporter/exporter.go index 07d5a461..7f89ae64 100644 --- a/pkg/exporter/exporter.go +++ b/pkg/exporter/exporter.go @@ -37,6 +37,7 @@ type Registry interface { GetHistogram(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) GetSummary(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) RemoveStaleMetrics() + ClearMetrics() int } type Exporter struct { @@ -49,6 +50,7 @@ type Exporter struct { EventStats *prometheus.CounterVec ConflictingEventStats *prometheus.CounterVec MetricsCount *prometheus.GaugeVec + clearMetrics chan chan int } // Listen handles all events sent to the given channel sequentially. It @@ -60,6 +62,8 @@ func (b *Exporter) Listen(e <-chan event.Events) { select { case <-removeStaleMetricsTicker.C: b.Registry.RemoveStaleMetrics() + case cleared := <-b.clearMetrics: + cleared <- b.Registry.ClearMetrics() case events, ok := <-e: if !ok { b.Logger.Debug("Channel is closed. Break out of Exporter.Listener.") @@ -73,6 +77,13 @@ func (b *Exporter) Listen(e <-chan event.Events) { } } +// ClearMetrics clears all dynamically registered StatsD time series. +func (b *Exporter) ClearMetrics() int { + cleared := make(chan int) + b.clearMetrics <- cleared + return <-cleared +} + // handleEvent processes a single Event according to the configured mapping. func (b *Exporter) handleEvent(thisEvent event.Event) { mapping, labels, present := b.Mapper.GetMapping(thisEvent.MetricName(), thisEvent.MetricType()) @@ -207,5 +218,6 @@ func NewExporter(reg prometheus.Registerer, mapper *mapper.MetricMapper, logger EventStats: eventStats, ConflictingEventStats: conflictingEventStats, MetricsCount: metricsCount, + clearMetrics: make(chan chan int), } } diff --git a/pkg/exporter/exporter_test.go b/pkg/exporter/exporter_test.go index a578674b..db4fe347 100644 --- a/pkg/exporter/exporter_test.go +++ b/pkg/exporter/exporter_test.go @@ -1168,6 +1168,74 @@ mappings: } } +func TestClearMetrics(t *testing.T) { + clockInstance := clock.ClockInstance + clock.ClockInstance = nil + defer func() { + clock.ClockInstance = clockInstance + }() + + reg := prometheus.NewRegistry() + testMapper := mapper.MetricMapper{} + ex := NewExporter(reg, &testMapper, promslog.NewNopLogger(), eventsActions, eventsUnmapped, errorEventStats, eventStats, conflictingEventStats, metricsCount) + events := make(chan event.Events) + done := make(chan struct{}) + go func() { + ex.Listen(events) + close(done) + }() + defer func() { + close(events) + <-done + }() + + events <- event.Events{ + &event.GaugeEvent{ + GMetricName: "clearable_gauge", + GValue: 200, + }, + } + events <- event.Events{} + + metrics, err := reg.Gather() + if err != nil { + t.Fatal("Gather should not fail") + } + gaugeValue := getFloat64(metrics, "clearable_gauge", prometheus.Labels{}) + if gaugeValue == nil || *gaugeValue != 200 { + t.Fatalf("Gauge `clearable_gauge` should be gathered with value 200, got %v", gaugeValue) + } + + if cleared := ex.ClearMetrics(); cleared != 1 { + t.Fatalf("Expected to clear 1 metric series, cleared %d", cleared) + } + + metrics, err = reg.Gather() + if err != nil { + t.Fatal("Gather should not fail") + } + if gaugeValue = getFloat64(metrics, "clearable_gauge", prometheus.Labels{}); gaugeValue != nil { + t.Fatalf("Gauge `clearable_gauge` should be cleared, got %v", *gaugeValue) + } + + events <- event.Events{ + &event.GaugeEvent{ + GMetricName: "clearable_gauge", + GValue: 42, + }, + } + events <- event.Events{} + + metrics, err = reg.Gather() + if err != nil { + t.Fatal("Gather should not fail") + } + gaugeValue = getFloat64(metrics, "clearable_gauge", prometheus.Labels{}) + if gaugeValue == nil || *gaugeValue != 42 { + t.Fatalf("Gauge `clearable_gauge` should be gathered again with value 42, got %v", gaugeValue) + } +} + func TestHashLabelNames(t *testing.T) { r := registry.NewRegistry(prometheus.DefaultRegisterer, nil) // Validate value hash changes and name has doesn't when just the value changes. diff --git a/pkg/registry/registry.go b/pkg/registry/registry.go index 825cb14c..64da30f2 100644 --- a/pkg/registry/registry.go +++ b/pkg/registry/registry.go @@ -387,14 +387,33 @@ func (r *Registry) RemoveStaleMetrics() { continue } if rm.LastRegisteredAt.Add(rm.TTL).Before(now) { - metric.Vectors[rm.VecKey].Holder.Delete(rm.Labels) - metric.Vectors[rm.VecKey].RefCount-- - delete(metric.Metrics, hash) + r.removeMetric(metric, hash, rm) } } } } +// ClearMetrics deletes all registered time series from the registry. +func (r *Registry) ClearMetrics() int { + removed := 0 + for _, metric := range r.Metrics { + for hash, rm := range metric.Metrics { + r.removeMetric(metric, hash, rm) + removed++ + } + } + return removed +} + +func (r *Registry) removeMetric(metric metrics.Metric, hash metrics.ValueHash, rm *metrics.RegisteredMetric) { + vector := metric.Vectors[rm.VecKey] + vector.Holder.Delete(rm.Labels) + if vector.RefCount > 0 { + vector.RefCount-- + } + delete(metric.Metrics, hash) +} + // Calculates a hash of both the label names and values. func (r *Registry) HashLabels(labels prometheus.Labels) (metrics.LabelHash, []string) { r.Hasher.Reset() From 76310625efe60bcd961219e7d1e6e7311a9953a4 Mon Sep 17 00:00:00 2001 From: YuMin Kim Date: Mon, 29 Jun 2026 18:18:29 +0900 Subject: [PATCH 2/2] Address clear metrics review feedback Signed-off-by: YuMin Kim --- README.md | 2 +- main.go | 10 +++- main_test.go | 24 ++++++++- pkg/exporter/exporter.go | 94 +++++++++++++++++++++++++++++++---- pkg/exporter/exporter_test.go | 88 +++++++++++++++++++++++++++++++- 5 files changed, 203 insertions(+), 15 deletions(-) diff --git a/README.md b/README.md index 1ec72a44..4f3d9d35 100644 --- a/README.md +++ b/README.md @@ -113,7 +113,7 @@ NOTE: Version 0.7.0 switched to the [kingpin](https://github.com/alecthomas/king ## Lifecycle API -The `statsd_exporter` has an optional lifecycle API (disabled by default) that can be used to reload, quit, or clear the exporter +The `statsd_exporter` has an optional lifecycle API (disabled by default) that can be used to reload, quit, or clear dynamically registered StatsD metric series by sending a `PUT` or `POST` request to the `/-/reload`, `/-/quit`, or `/-/clear` endpoints. ## Relay diff --git a/main.go b/main.go index 13e3960e..e6dd0ac0 100644 --- a/main.go +++ b/main.go @@ -15,6 +15,7 @@ package main import ( "bufio" + "context" "fmt" "log/slog" "net" @@ -206,13 +207,18 @@ func reloadConfig(fileName string, mapper *mapper.MetricMapper, logger *slog.Log } type metricsClearer interface { - ClearMetrics() int + ClearMetrics(context.Context) (int, error) } func clearMetricsHandler(clearer metricsClearer, logger *slog.Logger) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { if r.Method == http.MethodPut || r.Method == http.MethodPost { - cleared := clearer.ClearMetrics() + cleared, err := clearer.ClearMetrics(r.Context()) + if err != nil { + logger.Error("Failed to clear metrics", "error", err) + http.Error(w, err.Error(), http.StatusServiceUnavailable) + return + } logger.Info("Received lifecycle api clear", "metrics", cleared) fmt.Fprintf(w, "Cleared %d metric series", cleared) } diff --git a/main_test.go b/main_test.go index b88df8db..6acde9f7 100644 --- a/main_test.go +++ b/main_test.go @@ -14,6 +14,8 @@ package main import ( + "context" + "errors" "net/http" "net/http/httptest" "strings" @@ -25,11 +27,12 @@ import ( type fakeMetricsClearer struct { cleared int called int + err error } -func (f *fakeMetricsClearer) ClearMetrics() int { +func (f *fakeMetricsClearer) ClearMetrics(_ context.Context) (int, error) { f.called++ - return f.cleared + return f.cleared, f.err } func TestClearMetricsHandler(t *testing.T) { @@ -48,3 +51,20 @@ func TestClearMetricsHandler(t *testing.T) { t.Fatalf("unexpected response body: %q", body) } } + +func TestClearMetricsHandlerError(t *testing.T) { + clearer := &fakeMetricsClearer{err: errors.New("clear failed")} + handler := clearMetricsHandler(clearer, promslog.NewNopLogger()) + + request := httptest.NewRequest(http.MethodPost, "/-/clear", nil) + response := httptest.NewRecorder() + + handler.ServeHTTP(response, request) + + if response.Code != http.StatusServiceUnavailable { + t.Fatalf("expected status %d, got %d", http.StatusServiceUnavailable, response.Code) + } + if body := response.Body.String(); !strings.Contains(body, "clear failed") { + t.Fatalf("unexpected response body: %q", body) + } +} diff --git a/pkg/exporter/exporter.go b/pkg/exporter/exporter.go index 7f89ae64..b4d570e9 100644 --- a/pkg/exporter/exporter.go +++ b/pkg/exporter/exporter.go @@ -14,8 +14,11 @@ package exporter import ( + "context" + "errors" "log/slog" "os" + "sync" "time" "github.com/prometheus/client_golang/prometheus" @@ -37,9 +40,28 @@ type Registry interface { GetHistogram(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) GetSummary(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) RemoveStaleMetrics() +} + +var ( + // ErrClearMetricsUnsupported indicates that the configured registry cannot clear metrics. + ErrClearMetricsUnsupported = errors.New("registry does not support clearing metrics") + // ErrExporterNotRunning indicates that the exporter event loop is not running. + ErrExporterNotRunning = errors.New("exporter listener is not running") +) + +type clearableRegistry interface { ClearMetrics() int } +type clearMetricsRequest struct { + result chan clearMetricsResult +} + +type clearMetricsResult struct { + cleared int + err error +} + type Exporter struct { Mapper *mapper.MetricMapper Registry Registry @@ -50,24 +72,40 @@ type Exporter struct { EventStats *prometheus.CounterVec ConflictingEventStats *prometheus.CounterVec MetricsCount *prometheus.GaugeVec - clearMetrics chan chan int + clearMetrics chan clearMetricsRequest + listenStateMtx sync.RWMutex + listening bool + listenDone chan struct{} } // Listen handles all events sent to the given channel sequentially. It // terminates when the channel is closed. func (b *Exporter) Listen(e <-chan event.Events) { removeStaleMetricsTicker := clock.NewTicker(time.Second) + defer removeStaleMetricsTicker.Stop() + + b.listenStateMtx.Lock() + b.listening = true + b.listenDone = make(chan struct{}) + done := b.listenDone + b.listenStateMtx.Unlock() + + defer func() { + b.listenStateMtx.Lock() + b.listening = false + close(done) + b.listenStateMtx.Unlock() + }() for { select { case <-removeStaleMetricsTicker.C: b.Registry.RemoveStaleMetrics() - case cleared := <-b.clearMetrics: - cleared <- b.Registry.ClearMetrics() + case request := <-b.clearMetrics: + request.result <- b.clearRegistryMetrics() case events, ok := <-e: if !ok { b.Logger.Debug("Channel is closed. Break out of Exporter.Listener.") - removeStaleMetricsTicker.Stop() return } for _, event := range events { @@ -78,10 +116,48 @@ func (b *Exporter) Listen(e <-chan event.Events) { } // ClearMetrics clears all dynamically registered StatsD time series. -func (b *Exporter) ClearMetrics() int { - cleared := make(chan int) - b.clearMetrics <- cleared - return <-cleared +func (b *Exporter) ClearMetrics(ctx context.Context) (int, error) { + if ctx == nil { + ctx = context.Background() + } + + done, ok := b.listenState() + if !ok { + return 0, ErrExporterNotRunning + } + + request := clearMetricsRequest{ + result: make(chan clearMetricsResult, 1), + } + + select { + case b.clearMetrics <- request: + case <-done: + return 0, ErrExporterNotRunning + case <-ctx.Done(): + return 0, ctx.Err() + } + + select { + case result := <-request.result: + return result.cleared, result.err + case <-ctx.Done(): + return 0, ctx.Err() + } +} + +func (b *Exporter) listenState() (<-chan struct{}, bool) { + b.listenStateMtx.RLock() + defer b.listenStateMtx.RUnlock() + return b.listenDone, b.listening +} + +func (b *Exporter) clearRegistryMetrics() clearMetricsResult { + clearable, ok := b.Registry.(clearableRegistry) + if !ok { + return clearMetricsResult{err: ErrClearMetricsUnsupported} + } + return clearMetricsResult{cleared: clearable.ClearMetrics()} } // handleEvent processes a single Event according to the configured mapping. @@ -218,6 +294,6 @@ func NewExporter(reg prometheus.Registerer, mapper *mapper.MetricMapper, logger EventStats: eventStats, ConflictingEventStats: conflictingEventStats, MetricsCount: metricsCount, - clearMetrics: make(chan chan int), + clearMetrics: make(chan clearMetricsRequest), } } diff --git a/pkg/exporter/exporter_test.go b/pkg/exporter/exporter_test.go index db4fe347..60691362 100644 --- a/pkg/exporter/exporter_test.go +++ b/pkg/exporter/exporter_test.go @@ -14,6 +14,8 @@ package exporter import ( + "context" + "errors" "fmt" "log/slog" "net" @@ -1168,6 +1170,32 @@ mappings: } } +type noClearRegistry struct { + inner *registry.Registry +} + +var _ Registry = (*noClearRegistry)(nil) + +func (r *noClearRegistry) GetCounter(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Counter, error) { + return r.inner.GetCounter(metricName, labels, help, mapping, metricsCount) +} + +func (r *noClearRegistry) GetGauge(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Gauge, error) { + return r.inner.GetGauge(metricName, labels, help, mapping, metricsCount) +} + +func (r *noClearRegistry) GetHistogram(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) { + return r.inner.GetHistogram(metricName, labels, help, mapping, metricsCount) +} + +func (r *noClearRegistry) GetSummary(metricName string, labels prometheus.Labels, help string, mapping *mapper.MetricMapping, metricsCount *prometheus.GaugeVec) (prometheus.Observer, error) { + return r.inner.GetSummary(metricName, labels, help, mapping, metricsCount) +} + +func (r *noClearRegistry) RemoveStaleMetrics() { + r.inner.RemoveStaleMetrics() +} + func TestClearMetrics(t *testing.T) { clockInstance := clock.ClockInstance clock.ClockInstance = nil @@ -1206,7 +1234,11 @@ func TestClearMetrics(t *testing.T) { t.Fatalf("Gauge `clearable_gauge` should be gathered with value 200, got %v", gaugeValue) } - if cleared := ex.ClearMetrics(); cleared != 1 { + cleared, err := ex.ClearMetrics(context.Background()) + if err != nil { + t.Fatalf("ClearMetrics should not fail: %v", err) + } + if cleared != 1 { t.Fatalf("Expected to clear 1 metric series, cleared %d", cleared) } @@ -1236,6 +1268,60 @@ func TestClearMetrics(t *testing.T) { } } +func TestClearMetricsBeforeListen(t *testing.T) { + reg := prometheus.NewRegistry() + testMapper := mapper.MetricMapper{} + ex := NewExporter(reg, &testMapper, promslog.NewNopLogger(), eventsActions, eventsUnmapped, errorEventStats, eventStats, conflictingEventStats, metricsCount) + + _, err := ex.ClearMetrics(context.Background()) + if !errors.Is(err, ErrExporterNotRunning) { + t.Fatalf("expected ErrExporterNotRunning, got %v", err) + } +} + +func TestClearMetricsAfterListenStops(t *testing.T) { + reg := prometheus.NewRegistry() + testMapper := mapper.MetricMapper{} + ex := NewExporter(reg, &testMapper, promslog.NewNopLogger(), eventsActions, eventsUnmapped, errorEventStats, eventStats, conflictingEventStats, metricsCount) + events := make(chan event.Events) + done := make(chan struct{}) + go func() { + ex.Listen(events) + close(done) + }() + + close(events) + <-done + + _, err := ex.ClearMetrics(context.Background()) + if !errors.Is(err, ErrExporterNotRunning) { + t.Fatalf("expected ErrExporterNotRunning, got %v", err) + } +} + +func TestClearMetricsUnsupportedRegistry(t *testing.T) { + reg := prometheus.NewRegistry() + testMapper := mapper.MetricMapper{} + ex := NewExporter(reg, &testMapper, promslog.NewNopLogger(), eventsActions, eventsUnmapped, errorEventStats, eventStats, conflictingEventStats, metricsCount) + ex.Registry = &noClearRegistry{inner: registry.NewRegistry(reg, &testMapper)} + events := make(chan event.Events) + done := make(chan struct{}) + go func() { + ex.Listen(events) + close(done) + }() + defer func() { + close(events) + <-done + }() + events <- event.Events{} + + _, err := ex.ClearMetrics(context.Background()) + if !errors.Is(err, ErrClearMetricsUnsupported) { + t.Fatalf("expected ErrClearMetricsUnsupported, got %v", err) + } +} + func TestHashLabelNames(t *testing.T) { r := registry.NewRegistry(prometheus.DefaultRegisterer, nil) // Validate value hash changes and name has doesn't when just the value changes.