Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 25 additions & 24 deletions rest/fleet_manager_reporter.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,33 +54,34 @@ func (sc *ServerContext) reportFleetManagerMetrics(ctx context.Context) {
}
return settings
}
go runFleetManagerReportLoop(ctx, currentSettings, report)
}

go func() {
settings := currentSettings()
ticker := time.NewTicker(settings.Interval())
defer ticker.Stop()
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting interval set to %s (enabled=%t)", settings.Interval(), settings.Enabled)
// Report once on startup rather than waiting a full interval for the first tick.
report(settings)
for {
select {
case <-ticker.C:
refreshed := currentSettings()
if refreshed.Enabled != settings.Enabled {
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting enabled changed to %t", refreshed.Enabled)
}
if refreshed.Interval() != settings.Interval() {
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting interval changed from %s to %s", settings.Interval(), refreshed.Interval())
ticker.Reset(refreshed.Interval())
}
settings = refreshed
report(settings)
case <-ctx.Done():
base.InfofCtx(ctx, base.KeyAll, "Stopping fleet manager metrics reporting: %v", context.Cause(ctx))
return
func runFleetManagerReportLoop(ctx context.Context, currentSettings func() base.FleetManagerCollectorSettings, report func(base.FleetManagerCollectorSettings)) {
settings := currentSettings()
ticker := time.NewTicker(settings.Interval())
defer ticker.Stop()
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting interval set to %s (enabled=%t)", settings.Interval(), settings.Enabled)
// Report once on startup rather than waiting a full interval for the first tick.
report(settings)
for {
select {
case <-ticker.C:
refreshed := currentSettings()
if refreshed.Enabled != settings.Enabled {
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting enabled changed to %t", refreshed.Enabled)
}
if refreshed.Interval() != settings.Interval() {
base.InfofCtx(ctx, base.KeyAll, "Fleet manager metrics reporting interval changed from %s to %s", settings.Interval(), refreshed.Interval())
ticker.Reset(refreshed.Interval())
}
settings = refreshed
report(settings)
case <-ctx.Done():
base.InfofCtx(ctx, base.KeyAll, "Stopping fleet manager metrics reporting: %v", context.Cause(ctx))
return
}
}()
}
}

// sendFleetManagerMetrics POSTs the collected metrics to the ns_server fleet manager collector
Expand Down
66 changes: 66 additions & 0 deletions rest/fleet_manager_reporter_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,13 @@ package rest

import (
"bytes"
"context"
"fmt"
"net/http"
"sync"
"testing"
"testing/synctest"
"time"

"github.com/couchbase/sync_gateway/base"
"github.com/couchbase/sync_gateway/testing/assert"
Expand Down Expand Up @@ -204,3 +208,65 @@ func TestAllFleetManagerMetricsPopulated(t *testing.T) {
assert.NotEmpty(t, metrics.ProductInfo.Version)
assert.NotEmpty(t, metrics.ProductInfo.Name)
}

func TestFleetManagerReporterLoop(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
ctx, cancel := context.WithCancelCause(base.TestCtx(t))
t.Cleanup(func() {
cancel(nil)
})

var mutex sync.Mutex
var reports []base.FleetManagerCollectorSettings
current := base.FleetManagerCollectorSettings{Enabled: true, ReportingInterval: 1}
getCollectorFunc := func() base.FleetManagerCollectorSettings {
mutex.Lock()
defer mutex.Unlock()
return current
}
reportFunc := func(settings base.FleetManagerCollectorSettings) {
mutex.Lock()
defer mutex.Unlock()
if settings.Enabled { // mirror our report closure in production code
reports = append(reports, settings)
}
}
reportCount := func() int {
mutex.Lock()
defer mutex.Unlock()
return len(reports)
}

go runFleetManagerReportLoop(ctx, getCollectorFunc, reportFunc)

synctest.Wait() // wait until first report runs and is blocking on the ticker interval
require.Equal(t, 1, reportCount())

time.Sleep(1 * time.Hour) // mock 1 hour passing, should trigger new report on ticker interval
synctest.Wait() // wait for report goroutine to block again
require.Equal(t, 2, reportCount())

mutex.Lock()
current.Enabled = false
mutex.Unlock()
time.Sleep(1 * time.Hour) // mock another hour passing
synctest.Wait() // wait for report goroutine to block again after another report interval
require.Equal(t, 2, reportCount()) // tick has fired above but enabled is false so no report should fire

// update interval to two hours
mutex.Lock()
current.ReportingInterval = 2
current.Enabled = true
mutex.Unlock()
time.Sleep(2 * time.Hour)
synctest.Wait() // wait for report goroutine to block again after another report interval
require.Equal(t, 3, reportCount()) // should only report once (on the next 1 hour tick, then tick resets to 2 hours so no 4th report on second hour)
Comment on lines +261 to +263

// assert reports only contain entries of enabled=true
mutex.Lock()
for _, report := range reports {
assert.True(t, report.Enabled)
}
mutex.Unlock()
})
}
Loading