diff --git a/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic.go b/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic.go index 150cf37b607..91be6b119ca 100644 --- a/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic.go +++ b/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic.go @@ -20,32 +20,40 @@ import ( const defaultCacheTTL = time.Hour * 1 +// backgroundJob implements the dynamic angular detectors provider job that is periodically executed in the background. type backgroundJob interface { - backgroundJob(ctx context.Context) + // runBackgroundJob updates the dynamic angular detectors from a source. + runBackgroundJob(ctx context.Context) } // Dynamic is an angulardetector.DetectorsProvider that calls GCOM to get Angular detection patterns, // converts them to detectors and caches them for all future calls. // It also provides a background service that will periodically refresh the patterns from GCOM. type Dynamic struct { - log log.Logger - store angularpatternsstore.Service + log log.Logger httpClient *http.Client baseURL string - cacheTTL time.Duration + // store is the underlying angular patterns store used as a cache. + store angularpatternsstore.Service // initialRestoreDone is a channel that will be closed when the first restore from db is done by the // background service. It can be used to wait for the first restore to be done by reading a value from this channel. initialRestoreDone chan struct{} + // detectors contains the cached angular detectors, which are created from the remote angular patterns. + // Use setDetectors and ProvideDetectors to write/read this value. detectors []angulardetector.AngularDetector - mux sync.RWMutex - // bgJob is the implementation of the background job. - // This is called when scheduled by Run(). - bgJob backgroundJob + // mux is the mutex used to read/write the cached detectors in a concurrency-safe way. + mux sync.RWMutex + + // backgroundJob is the implementation of the background job. This is periodically invoked by Run(). + backgroundJob backgroundJob + + // backgroundJobInterval is the interval between the periodic background job calls. + backgroundJobInterval time.Duration } func ProvideDynamic(cfg *config.Cfg, store angularpatternsstore.Service) (*Dynamic, error) { @@ -55,16 +63,15 @@ func ProvideDynamic(cfg *config.Cfg, store angularpatternsstore.Service) (*Dynam return nil, fmt.Errorf("httpclient new: %w", err) } d := &Dynamic{ - log: log.New("plugins.angulardetector.gcom"), - store: store, - httpClient: cl, - baseURL: cfg.GrafanaComURL, - cacheTTL: defaultCacheTTL, - - initialRestoreDone: make(chan struct{}), + log: log.New("plugins.angulardetector.gcom"), + store: store, + httpClient: cl, + baseURL: cfg.GrafanaComURL, + backgroundJobInterval: defaultCacheTTL, + initialRestoreDone: make(chan struct{}), } - // By default, use ourselves as bgJob - d.bgJob = d + // By default, use ourselves as backgroundJob + d.backgroundJob = d return d, nil } @@ -150,6 +157,7 @@ func (d *Dynamic) fetchAndStoreDetectors(ctx context.Context) ([]angulardetector return newDetectors, nil } +// setDetectors sets the detectors by acquiring the lock first. func (d *Dynamic) setDetectors(newDetectors []angulardetector.AngularDetector) { d.mux.Lock() d.detectors = newDetectors @@ -181,7 +189,7 @@ func (d *Dynamic) tryUpdateDetectors(ctx context.Context) { // backgroundJob is the function executed periodically in the background by the background service. // It calls tryUpdateDetectors. -func (d *Dynamic) backgroundJob(ctx context.Context) { +func (d *Dynamic) runBackgroundJob(ctx context.Context) { d.tryUpdateDetectors(ctx) } @@ -210,9 +218,9 @@ func (d *Dynamic) setDetectorsFromCache(ctx context.Context) error { } } -// notifyInitialRestoreDone sets the initial restore as "done" +// notifyInitialRestoreDone sets the initial restore as "done", and will unblock all goroutines waiting +// for the initial restore to be completed. func (d *Dynamic) notifyInitialRestoreDone() { - // Notify that the initial restore is done (see docstring for d.initialRestoreDone) close(d.initialRestoreDone) } @@ -231,11 +239,11 @@ func (d *Dynamic) Run(ctx context.Context) error { if err != nil { return fmt.Errorf("get last updated: %w", err) } - nextRunUntil := time.Until(lastUpdate.Add(d.cacheTTL)) + nextRunUntil := time.Until(lastUpdate.Add(d.backgroundJobInterval)) if nextRunUntil <= 0 { // Do first run immediately - d.bgJob.backgroundJob(ctx) - nextRunUntil = d.cacheTTL + d.backgroundJob.runBackgroundJob(ctx) + nextRunUntil = d.backgroundJobInterval } // Keep running periodically @@ -244,9 +252,9 @@ func (d *Dynamic) Run(ctx context.Context) error { for { select { case <-ticker.C: - d.bgJob.backgroundJob(ctx) + d.backgroundJob.runBackgroundJob(ctx) // Restore default TTL if we run with a shorter interval the first time - ticker.Reset(d.cacheTTL) + ticker.Reset(d.backgroundJobInterval) case <-ctx.Done(): return ctx.Err() } @@ -256,7 +264,7 @@ func (d *Dynamic) Run(ctx context.Context) error { // ProvideDetectors returns the cached detectors. It returns an empty slice if there's no value. // TODO: remove context here? func (d *Dynamic) ProvideDetectors(_ context.Context) []angulardetector.AngularDetector { - // Wait for channel to be closed, which is done after the restore from db is done + // Block until channel is closed, which is done after the restore from db is done. <-d.initialRestoreDone d.mux.RLock() diff --git a/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic_test.go b/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic_test.go index 536e5ca7ab3..bd26d57625b 100644 --- a/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic_test.go +++ b/pkg/services/pluginsintegration/angulardetectorsprovider/dynamic_test.go @@ -134,9 +134,7 @@ func TestDynamicAngularDetectorsProvider(t *testing.T) { require.NotEmpty(t, dbV, "new store should not be empty") var patterns GCOMPatterns require.NoError(t, json.Unmarshal([]byte(dbV), &patterns), "could not unmarshal stored value") - detectors, err := svc.patternsToDetectors(patterns) - require.NoError(t, err, "could not convert patterns to detectors") - checkMockDetectors(t, detectors) + require.Equal(t, mockGCOMPatterns, patterns) // Check that last updated has been updated in the kv store (which is used for cache ttl) lastUpdated, err = svc.store.GetLastUpdated(context.Background()) @@ -405,7 +403,7 @@ func provideDynamic(t *testing.T, gcomURL string, cacheTTL time.Duration, opts . for _, opt := range opts { opt(d) } - d.cacheTTL = cacheTTL + d.backgroundJobInterval = cacheTTL return d } @@ -431,8 +429,8 @@ func (j *fakeBackgroundJob) close() { close(j.callback) } -func (j *fakeBackgroundJob) backgroundJob(ctx context.Context) { - j.inner.backgroundJob(ctx) +func (j *fakeBackgroundJob) runBackgroundJob(ctx context.Context) { + j.inner.runBackgroundJob(ctx) j.callback <- struct{}{} } @@ -462,7 +460,7 @@ func newBackgroundServiceScenario(svc *Dynamic, callback func()) *backgroundServ bgDone: make(chan struct{}), fakeBgJob: newFakeJober(svc), } - svc.bgJob = s.fakeBgJob + svc.backgroundJob = s.fakeBgJob go func() { for range s.fakeBgJob.callback { callback()