diff --git a/pkg/poller/inputs.go b/pkg/poller/inputs.go index 0c64cab6..c37816b1 100644 --- a/pkg/poller/inputs.go +++ b/pkg/poller/inputs.go @@ -315,10 +315,24 @@ func AppendMetrics(existing *Metrics, m *Metrics) *Metrics { return existing } + // The aggregate starts as a bare &Metrics{}, so without this its TS stays zero and + // influxunifi's per-point fallback (see influxdb.go collect) stamps every point that + // carries no timestamp of its own with the zero time. + // + // First writer wins, and that is whichever input's result is merged first -- in + // collectMetrics the inputs run concurrently and are drained in completion order, so with + // more than one input configured it is not deterministic which one supplies the TS. That + // is fine here: every input stamps the moment it began polling within the same tick, so + // any of them places the batch in the right poll interval. Do not read this as "earliest". + if existing.TS.IsZero() { + existing.TS = m.TS + } + existing.SitesDPI = append(existing.SitesDPI, m.SitesDPI...) existing.Sites = append(existing.Sites, m.Sites...) existing.ClientsDPI = append(existing.ClientsDPI, m.ClientsDPI...) existing.RogueAPs = append(existing.RogueAPs, m.RogueAPs...) + existing.SpeedTests = append(existing.SpeedTests, m.SpeedTests...) existing.Clients = append(existing.Clients, m.Clients...) existing.Devices = append(existing.Devices, m.Devices...) existing.CountryTraffic = append(existing.CountryTraffic, m.CountryTraffic...) diff --git a/pkg/poller/inputs_test.go b/pkg/poller/inputs_test.go index 6a41f2d1..577d7e7f 100644 --- a/pkg/poller/inputs_test.go +++ b/pkg/poller/inputs_test.go @@ -1,7 +1,9 @@ package poller_test import ( + "reflect" "testing" + "time" "github.com/unpoller/unpoller/pkg/poller" @@ -96,3 +98,60 @@ func TestCollectEventsRecoversPanickingInput(t *testing.T) { require.Error(t, err) assert.Contains(t, err.Error(), "panic-input") } + +// TestAppendMetricsCoversEverySliceField walks the Metrics struct by reflection and fails if +// any slice field is not merged by AppendMetrics. +// +// This is a guard against a silent, invisible failure mode: a new metric family needs both a +// field here and an append line there, and omitting the second loses every one of those +// metrics with no error, no log line, and a passing build. SpeedTests was lost that way -- +// collected by inputunifi and exported by all three output plugins, but never merged, so the +// export code was dead for every user. Asserting field-by-field would just repeat the bug; +// reflection means a new field is covered the moment it is declared. +func TestAppendMetricsCoversEverySliceField(t *testing.T) { + t.Parallel() + + // Put one zero-valued element in every slice, so a merged field has two and a dropped + // field has one. The element values are irrelevant -- only the count is under test. + fill := func(m *poller.Metrics) { + v := reflect.ValueOf(m).Elem() + for i := 0; i < v.NumField(); i++ { + if f := v.Field(i); f.Kind() == reflect.Slice { + f.Set(reflect.Append(f, reflect.Zero(f.Type().Elem()))) + } + } + } + + existing, incoming := &poller.Metrics{}, &poller.Metrics{} + fill(existing) + fill(incoming) + + got := reflect.ValueOf(poller.AppendMetrics(existing, incoming)).Elem() + + for i := 0; i < got.NumField(); i++ { + f := got.Field(i) + if f.Kind() != reflect.Slice { + continue + } + + assert.Equal(t, 2, f.Len(), + "Metrics.%s is missing an append line in AppendMetrics, so those metrics are silently dropped", + got.Type().Field(i).Name) + } +} + +func TestAppendMetricsPropagatesTS(t *testing.T) { + t.Parallel() + + first := time.Now().Add(-time.Minute) + second := time.Now() + + // The aggregate starts bare, so it must adopt the first input's timestamp -- otherwise it + // stays zero and influxunifi stamps untimed points with the zero time. + got := poller.AppendMetrics(&poller.Metrics{}, &poller.Metrics{TS: first}) + assert.Equal(t, first, got.TS) + + // A timestamp already set wins: it marks when the batch began. + got = poller.AppendMetrics(&poller.Metrics{TS: first}, &poller.Metrics{TS: second}) + assert.Equal(t, first, got.TS) +}