From bd6a5c900f34dfba247913005133c7d96291014c Mon Sep 17 00:00:00 2001 From: Alexander Weaver Date: Mon, 26 Sep 2022 12:35:33 -0500 Subject: [PATCH] Alerting: Extract ticker into shared package (#55703) * Move ticker files to dedicated package with no changes * Fix package naming and resolve naming conflicts * Fix up all existing references to moved objects * Remove all alerting-specific references from shared util * Rename TickerMetrics to simply Metrics * Rename base ticker type to T and rename NewTicker to simply New --- pkg/services/alerting/engine.go | 6 +-- pkg/services/ngalert/metrics/ngalert.go | 6 +-- pkg/services/ngalert/schedule/schedule.go | 8 ++-- .../metrics => util/ticker}/metrics.go | 14 +++---- .../alerting => util/ticker}/ticker.go | 17 ++++---- .../alerting => util/ticker}/ticker_test.go | 40 +++++++++---------- 6 files changed, 44 insertions(+), 47 deletions(-) rename pkg/{services/alerting/metrics => util/ticker}/metrics.go (81%) rename pkg/{services/alerting => util/ticker}/ticker.go (86%) rename pkg/{services/alerting => util/ticker}/ticker_test.go (73%) diff --git a/pkg/services/alerting/engine.go b/pkg/services/alerting/engine.go index 67a8c31c055..316b3771305 100644 --- a/pkg/services/alerting/engine.go +++ b/pkg/services/alerting/engine.go @@ -16,7 +16,6 @@ import ( "github.com/grafana/grafana/pkg/infra/tracing" "github.com/grafana/grafana/pkg/infra/usagestats" "github.com/grafana/grafana/pkg/models" - "github.com/grafana/grafana/pkg/services/alerting/metrics" "github.com/grafana/grafana/pkg/services/annotations" "github.com/grafana/grafana/pkg/services/dashboards" "github.com/grafana/grafana/pkg/services/datasources" @@ -25,6 +24,7 @@ import ( "github.com/grafana/grafana/pkg/services/rendering" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/tsdb/legacydata" + "github.com/grafana/grafana/pkg/util/ticker" ) // AlertEngine is the background process that @@ -37,7 +37,7 @@ type AlertEngine struct { Cfg *setting.Cfg execQueue chan *Job - ticker *Ticker + ticker *ticker.T scheduler scheduler evalHandler evalHandler ruleReader ruleReader @@ -90,7 +90,7 @@ func ProvideAlertEngine(renderer rendering.Service, requestValidator models.Plug // Run starts the alerting service background process. func (e *AlertEngine) Run(ctx context.Context) error { reg := prometheus.WrapRegistererWithPrefix("legacy_", prometheus.DefaultRegisterer) - e.ticker = NewTicker(clock.New(), 1*time.Second, metrics.NewTickerMetrics(reg)) + e.ticker = ticker.New(clock.New(), 1*time.Second, ticker.NewMetrics(reg, "alerting")) defer e.ticker.Stop() alertGroup, ctx := errgroup.WithContext(ctx) alertGroup.Go(func() error { return e.alertingTicker(ctx) }) diff --git a/pkg/services/ngalert/metrics/ngalert.go b/pkg/services/ngalert/metrics/ngalert.go index b9941400a27..47d5bd4f095 100644 --- a/pkg/services/ngalert/metrics/ngalert.go +++ b/pkg/services/ngalert/metrics/ngalert.go @@ -13,8 +13,8 @@ import ( "github.com/grafana/grafana/pkg/api/response" "github.com/grafana/grafana/pkg/models" - legacyMetrics "github.com/grafana/grafana/pkg/services/alerting/metrics" apimodels "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" + "github.com/grafana/grafana/pkg/util/ticker" "github.com/grafana/grafana/pkg/web" ) @@ -55,7 +55,7 @@ type Scheduler struct { SchedulableAlertRules prometheus.Gauge SchedulableAlertRulesHash prometheus.Gauge UpdateSchedulableAlertRulesDuration prometheus.Histogram - Ticker *legacyMetrics.Ticker + Ticker *ticker.Metrics EvaluationMissed *prometheus.CounterVec } @@ -199,7 +199,7 @@ func newSchedulerMetrics(r prometheus.Registerer) *Scheduler { Buckets: []float64{0.1, 0.25, 0.5, 1, 2, 5, 10}, }, ), - Ticker: legacyMetrics.NewTickerMetrics(r), + Ticker: ticker.NewMetrics(r, "alerting"), EvaluationMissed: promauto.With(r).NewCounterVec( prometheus.CounterOpts{ Namespace: Namespace, diff --git a/pkg/services/ngalert/schedule/schedule.go b/pkg/services/ngalert/schedule/schedule.go index 14c52b08c1b..fd426a59b4c 100644 --- a/pkg/services/ngalert/schedule/schedule.go +++ b/pkg/services/ngalert/schedule/schedule.go @@ -10,7 +10,6 @@ import ( prometheusModel "github.com/prometheus/common/model" "github.com/grafana/grafana/pkg/infra/log" - "github.com/grafana/grafana/pkg/services/alerting" "github.com/grafana/grafana/pkg/services/datasources" "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" "github.com/grafana/grafana/pkg/services/ngalert/eval" @@ -21,6 +20,7 @@ import ( "github.com/grafana/grafana/pkg/services/org" "github.com/grafana/grafana/pkg/services/user" "github.com/grafana/grafana/pkg/setting" + "github.com/grafana/grafana/pkg/util/ticker" "github.com/benbjohnson/clock" "golang.org/x/sync/errgroup" @@ -68,7 +68,7 @@ type schedule struct { clock clock.Clock - ticker *alerting.Ticker + ticker *ticker.T // evalApplied is only used for tests: test code can set it to non-nil // function, and then it'll be called from the event loop whenever the @@ -119,7 +119,7 @@ type SchedulerCfg struct { // NewScheduler returns a new schedule. func NewScheduler(cfg SchedulerCfg, appURL *url.URL, stateManager *state.Manager) *schedule { - ticker := alerting.NewTicker(cfg.C, cfg.Cfg.BaseInterval, cfg.Metrics.Ticker) + ticker := ticker.New(cfg.C, cfg.Cfg.BaseInterval, cfg.Metrics.Ticker) sch := schedule{ registry: alertRuleInfoRegistry{alertRuleInfo: make(map[ngmodels.AlertRuleKey]*alertRuleInfo)}, @@ -449,7 +449,7 @@ func (sch *schedule) overrideCfg(cfg SchedulerCfg) { sch.clock = cfg.C sch.baseInterval = cfg.Cfg.BaseInterval sch.ticker.Stop() - sch.ticker = alerting.NewTicker(cfg.C, cfg.Cfg.BaseInterval, cfg.Metrics.Ticker) + sch.ticker = ticker.New(cfg.C, cfg.Cfg.BaseInterval, cfg.Metrics.Ticker) sch.evalAppliedFunc = cfg.EvalAppliedFunc sch.stopAppliedFunc = cfg.StopAppliedFunc } diff --git a/pkg/services/alerting/metrics/metrics.go b/pkg/util/ticker/metrics.go similarity index 81% rename from pkg/services/alerting/metrics/metrics.go rename to pkg/util/ticker/metrics.go index e93087aa8dd..9c8a60e3850 100644 --- a/pkg/services/alerting/metrics/metrics.go +++ b/pkg/util/ticker/metrics.go @@ -1,33 +1,33 @@ -package metrics +package ticker import ( "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promauto" ) -type Ticker struct { +type Metrics struct { LastTickTime prometheus.Gauge NextTickTime prometheus.Gauge IntervalSeconds prometheus.Gauge } -func NewTickerMetrics(reg prometheus.Registerer) *Ticker { - return &Ticker{ +func NewMetrics(reg prometheus.Registerer, subsystem string) *Metrics { + return &Metrics{ LastTickTime: promauto.With(reg).NewGauge(prometheus.GaugeOpts{ Namespace: "grafana", - Subsystem: "alerting", + Subsystem: subsystem, Name: "ticker_last_consumed_tick_timestamp_seconds", Help: "Timestamp of the last consumed tick in seconds.", }), NextTickTime: promauto.With(reg).NewGauge(prometheus.GaugeOpts{ Namespace: "grafana", - Subsystem: "alerting", + Subsystem: subsystem, Name: "ticker_next_tick_timestamp_seconds", Help: "Timestamp of the next tick in seconds before it is consumed.", }), IntervalSeconds: promauto.With(reg).NewGauge(prometheus.GaugeOpts{ Namespace: "grafana", - Subsystem: "alerting", + Subsystem: subsystem, Name: "ticker_interval_seconds", Help: "Interval at which the ticker is meant to tick.", }), diff --git a/pkg/services/alerting/ticker.go b/pkg/util/ticker/ticker.go similarity index 86% rename from pkg/services/alerting/ticker.go rename to pkg/util/ticker/ticker.go index f7aa36def7c..502325ca07f 100644 --- a/pkg/services/alerting/ticker.go +++ b/pkg/util/ticker/ticker.go @@ -1,4 +1,4 @@ -package alerting +package ticker import ( "fmt" @@ -7,29 +7,28 @@ import ( "github.com/benbjohnson/clock" "github.com/grafana/grafana/pkg/infra/log" - "github.com/grafana/grafana/pkg/services/alerting/metrics" ) -// Ticker is a ticker to power the alerting scheduler. it's like a time.Ticker, except: +// Ticker emits ticks at regular time intervals. it's like a time.Ticker, except: // - it doesn't drop ticks for slow receivers, rather, it queues up. so that callers are in control to instrument what's going on. // - it ticks on interval marks or very shortly after. this provides a predictable load pattern // (this shouldn't cause too much load contention issues because the next steps in the pipeline just process at their own pace) // - the timestamps are used to mark "last datapoint to query for" and as such, are a configurable amount of seconds in the past -type Ticker struct { +type T struct { C chan time.Time clock clock.Clock last time.Time interval time.Duration - metrics *metrics.Ticker + metrics *Metrics stopCh chan struct{} } // NewTicker returns a Ticker that ticks on interval marks (or very shortly after) starting at c.Now(), and never drops ticks. interval should not be negative or zero. -func NewTicker(c clock.Clock, interval time.Duration, metric *metrics.Ticker) *Ticker { +func New(c clock.Clock, interval time.Duration, metric *Metrics) *T { if interval <= 0 { panic(fmt.Errorf("non-positive interval [%v] is not allowed", interval)) } - t := &Ticker{ + t := &T{ C: make(chan time.Time), clock: c, last: getStartTick(c, interval), @@ -47,7 +46,7 @@ func getStartTick(clk clock.Clock, interval time.Duration) time.Time { return time.Unix(0, nano-(nano%interval.Nanoseconds())) } -func (t *Ticker) run() { +func (t *T) run() { logger := log.New("ticker") logger.Info("starting", "first_tick", t.last.Add(t.interval)) LOOP: @@ -77,7 +76,7 @@ LOOP: } // Stop stops the ticker. It does not close the C channel -func (t *Ticker) Stop() { +func (t *T) Stop() { select { case t.stopCh <- struct{}{}: default: diff --git a/pkg/services/alerting/ticker_test.go b/pkg/util/ticker/ticker_test.go similarity index 73% rename from pkg/services/alerting/ticker_test.go rename to pkg/util/ticker/ticker_test.go index f759d5019ab..39cbca221e7 100644 --- a/pkg/services/alerting/ticker_test.go +++ b/pkg/util/ticker/ticker_test.go @@ -1,4 +1,4 @@ -package alerting +package ticker import ( "bytes" @@ -13,8 +13,6 @@ import ( "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/testutil" "github.com/stretchr/testify/require" - - "github.com/grafana/grafana/pkg/services/alerting/metrics" ) func TestTicker(t *testing.T) { @@ -53,7 +51,7 @@ func TestTicker(t *testing.T) { interval := time.Duration(rand.Int63n(100)+10) * time.Second clk := clock.NewMock() clk.Add(interval) // align clock with the start tick - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) ticks := rand.Intn(9) + 1 jitter := rand.Int63n(int64(interval) - 1) @@ -87,7 +85,7 @@ func TestTicker(t *testing.T) { t.Run("should not put anything to channel until it's time", func(t *testing.T) { clk := clock.NewMock() interval := time.Duration(rand.Int63n(9)+1) * time.Second - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) expectedTick := clk.Now().Add(interval) for { require.Empty(t, ticker.C) @@ -104,7 +102,7 @@ func TestTicker(t *testing.T) { t.Run("should put the tick in the channel immediately if it is behind", func(t *testing.T) { clk := clock.NewMock() interval := time.Duration(rand.Int63n(9)+1) * time.Second - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) // We can expect the first tick to be at a consistent interval. Take a snapshot of the clock now, before we advance it. expectedTick := clk.Now().Add(interval) @@ -133,25 +131,25 @@ func TestTicker(t *testing.T) { clk.Set(time.Now()) interval := time.Duration(rand.Int63n(9)+1) * time.Second registry := prometheus.NewPedanticRegistry() - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(registry)) + ticker := New(clk, interval, NewMetrics(registry, "test")) expectedTick := getStartTick(clk, interval).Add(interval) - expectedMetricFmt := `# HELP grafana_alerting_ticker_interval_seconds Interval at which the ticker is meant to tick. - # TYPE grafana_alerting_ticker_interval_seconds gauge - grafana_alerting_ticker_interval_seconds %v - # HELP grafana_alerting_ticker_last_consumed_tick_timestamp_seconds Timestamp of the last consumed tick in seconds. - # TYPE grafana_alerting_ticker_last_consumed_tick_timestamp_seconds gauge - grafana_alerting_ticker_last_consumed_tick_timestamp_seconds %v - # HELP grafana_alerting_ticker_next_tick_timestamp_seconds Timestamp of the next tick in seconds before it is consumed. - # TYPE grafana_alerting_ticker_next_tick_timestamp_seconds gauge - grafana_alerting_ticker_next_tick_timestamp_seconds %v + expectedMetricFmt := `# HELP grafana_test_ticker_interval_seconds Interval at which the ticker is meant to tick. + # TYPE grafana_test_ticker_interval_seconds gauge + grafana_test_ticker_interval_seconds %v + # HELP grafana_test_ticker_last_consumed_tick_timestamp_seconds Timestamp of the last consumed tick in seconds. + # TYPE grafana_test_ticker_last_consumed_tick_timestamp_seconds gauge + grafana_test_ticker_last_consumed_tick_timestamp_seconds %v + # HELP grafana_test_ticker_next_tick_timestamp_seconds Timestamp of the next tick in seconds before it is consumed. + # TYPE grafana_test_ticker_next_tick_timestamp_seconds gauge + grafana_test_ticker_next_tick_timestamp_seconds %v ` expectedMetric := fmt.Sprintf(expectedMetricFmt, interval.Seconds(), 0, float64(expectedTick.UnixNano())/1e9) errs := make(map[string]error, 1) require.Eventuallyf(t, func() bool { - err := testutil.GatherAndCompare(registry, bytes.NewBufferString(expectedMetric), "grafana_alerting_ticker_last_consumed_tick_timestamp_seconds", "grafana_alerting_ticker_next_tick_timestamp_seconds", "grafana_alerting_ticker_interval_seconds") + err := testutil.GatherAndCompare(registry, bytes.NewBufferString(expectedMetric), "grafana_test_ticker_last_consumed_tick_timestamp_seconds", "grafana_test_ticker_next_tick_timestamp_seconds", "grafana_test_ticker_interval_seconds") if err != nil { errs["error"] = err } @@ -164,7 +162,7 @@ func TestTicker(t *testing.T) { expectedMetric = fmt.Sprintf(expectedMetricFmt, interval.Seconds(), float64(actual.UnixNano())/1e9, float64(expectedTick.Add(interval).UnixNano())/1e9) require.Eventuallyf(t, func() bool { - err := testutil.GatherAndCompare(registry, bytes.NewBufferString(expectedMetric), "grafana_alerting_ticker_last_consumed_tick_timestamp_seconds", "grafana_alerting_ticker_next_tick_timestamp_seconds", "grafana_alerting_ticker_interval_seconds") + err := testutil.GatherAndCompare(registry, bytes.NewBufferString(expectedMetric), "grafana_test_ticker_last_consumed_tick_timestamp_seconds", "grafana_test_ticker_next_tick_timestamp_seconds", "grafana_test_ticker_interval_seconds") if err != nil { errs["error"] = err } @@ -176,7 +174,7 @@ func TestTicker(t *testing.T) { t.Run("when it waits for the next tick", func(t *testing.T) { clk := clock.NewMock() interval := time.Duration(rand.Int63n(9)+1) * time.Second - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) clk.Add(interval) readChanOrFail(t, ticker.C) ticker.Stop() @@ -187,7 +185,7 @@ func TestTicker(t *testing.T) { t.Run("when it waits for the tick to be consumed", func(t *testing.T) { clk := clock.NewMock() interval := time.Duration(rand.Int63n(9)+1) * time.Second - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) clk.Add(interval) ticker.Stop() require.Empty(t, ticker.C) @@ -196,7 +194,7 @@ func TestTicker(t *testing.T) { t.Run("multiple times", func(t *testing.T) { clk := clock.NewMock() interval := time.Duration(rand.Int63n(9)+1) * time.Second - ticker := NewTicker(clk, interval, metrics.NewTickerMetrics(prometheus.NewRegistry())) + ticker := New(clk, interval, NewMetrics(prometheus.NewRegistry(), "test")) ticker.Stop() ticker.Stop() ticker.Stop()