From 978f1119d7d6774a4d6fc9acb779a32cb82c8aa0 Mon Sep 17 00:00:00 2001 From: Yuri Tseretyan Date: Fri, 4 Nov 2022 17:06:47 -0400 Subject: [PATCH] Alerting: Run state manager as regular sub-service (#58246) --- pkg/services/ngalert/ngalert.go | 4 ++ pkg/services/ngalert/schedule/schedule.go | 2 - pkg/services/ngalert/state/manager.go | 46 +++++++++-------------- 3 files changed, 22 insertions(+), 30 deletions(-) diff --git a/pkg/services/ngalert/ngalert.go b/pkg/services/ngalert/ngalert.go index b1e23e36420..d61b07fe9c0 100644 --- a/pkg/services/ngalert/ngalert.go +++ b/pkg/services/ngalert/ngalert.go @@ -282,6 +282,10 @@ func (ng *AlertNG) Run(ctx context.Context) error { children, subCtx := errgroup.WithContext(ctx) + children.Go(func() error { + return ng.stateManager.Run(subCtx) + }) + children.Go(func() error { return ng.MultiOrgAlertmanager.Run(subCtx) }) diff --git a/pkg/services/ngalert/schedule/schedule.go b/pkg/services/ngalert/schedule/schedule.go index ea9a3e5eb1c..b33c0de0d6d 100644 --- a/pkg/services/ngalert/schedule/schedule.go +++ b/pkg/services/ngalert/schedule/schedule.go @@ -298,8 +298,6 @@ func (sch *schedule) schedulePeriodic(ctx context.Context, t *ticker.T) error { case <-ctx.Done(): // waiting for all rule evaluation routines to stop waitErr := dispatcherGroup.Wait() - // close the state manager and flush the state - sch.stateManager.Close() return waitErr } } diff --git a/pkg/services/ngalert/state/manager.go b/pkg/services/ngalert/state/manager.go index ad27bb68278..d46b5414a6a 100644 --- a/pkg/services/ngalert/state/manager.go +++ b/pkg/services/ngalert/state/manager.go @@ -15,7 +15,10 @@ import ( ngModels "github.com/grafana/grafana/pkg/services/ngalert/models" ) -var ResendDelay = 30 * time.Second +var ( + ResendDelay = 30 * time.Second + MetricsScrapeInterval = 15 * time.Second // TODO: parameterize? // Setting to a reasonable default scrape interval for Prometheus. +) // AlertInstanceManager defines the interface for querying the current alert instances. type AlertInstanceManager interface { @@ -29,7 +32,6 @@ type Manager struct { clock clock.Clock cache *cache - quit chan struct{} ResendDelay time.Duration instanceStore InstanceStore @@ -39,9 +41,8 @@ type Manager struct { } func NewManager(metrics *metrics.State, externalURL *url.URL, instanceStore InstanceStore, imageService image.ImageService, clock clock.Clock, historian Historian) *Manager { - manager := &Manager{ + return &Manager{ cache: newCache(), - quit: make(chan struct{}), ResendDelay: ResendDelay, // TODO: make this configurable log: log.New("ngalert.state.manager"), metrics: metrics, @@ -51,14 +52,21 @@ func NewManager(metrics *metrics.State, externalURL *url.URL, instanceStore Inst clock: clock, externalURL: externalURL, } - if manager.metrics != nil { - go manager.recordMetrics() - } - return manager } -func (st *Manager) Close() { - st.quit <- struct{}{} +func (st *Manager) Run(ctx context.Context) error { + ticker := st.clock.Ticker(MetricsScrapeInterval) + for { + select { + case <-ticker.C: + st.log.Debug("Recording state cache metrics", "now", st.clock.Now()) + st.cache.recordMetrics(st.metrics) + case <-ctx.Done(): + st.log.Debug("Stopping") + ticker.Stop() + return ctx.Err() + } + } } func (st *Manager) Warm(ctx context.Context, rulesReader RuleReader) { @@ -269,24 +277,6 @@ func (st *Manager) GetStatesForRuleUID(orgID int64, alertRuleUID string) []*Stat return st.cache.getStatesForRuleUID(orgID, alertRuleUID) } -func (st *Manager) recordMetrics() { - // TODO: parameterize? - // Setting to a reasonable default scrape interval for Prometheus. - dur := time.Duration(15) * time.Second - ticker := st.clock.Ticker(dur) - for { - select { - case <-ticker.C: - st.log.Debug("Recording state cache metrics", "now", st.clock.Now()) - st.cache.recordMetrics(st.metrics) - case <-st.quit: - st.log.Debug("Stopping state cache metrics recording", "now", st.clock.Now()) - ticker.Stop() - return - } - } -} - func (st *Manager) Put(states []*State) { for _, s := range states { st.cache.set(s)