mirror of
https://github.com/unpoller/unpoller.git
synced 2026-09-30 03:21:28 +02:00
fix: honor scrape-cache disable and export locate-mode devices
interval=0 now turns off the Prometheus scrape cache as PR #1014 documented, and sub-15s intervals warn instead of clamping (#1083). Adopted devices stay in Prometheus, Influx, OTel, and Datadog exports while locate/identify is on (#1075). Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -37,7 +37,9 @@
|
||||
dead_ports = false
|
||||
# How often the background poller refreshes the cache served to /metrics.
|
||||
# Decouples scrape cadence from UniFi API calls so 429 backoff loops no
|
||||
# longer block /metrics. Default: 60s. Values below 15s are clamped to 15s.
|
||||
# longer block /metrics. Omitted: cache on, 60s. Set to 0 / "0" / "0s" to
|
||||
# disable the cache so /metrics fetches live. Values below 15s log a
|
||||
# warning but are not clamped. Env: UP_PROMETHEUS_INTERVAL.
|
||||
interval = "60s"
|
||||
|
||||
[influxdb]
|
||||
|
||||
@@ -20,6 +20,8 @@ prometheus:
|
||||
ssl_cert_path: ""
|
||||
ssl_key_path: ""
|
||||
report_errors: false
|
||||
# Cache refresh. Omitted: 60s. 0 / "0s" disables the cache (live /metrics).
|
||||
# Values below 15s warn but are not clamped. Env: UP_PROMETHEUS_INTERVAL.
|
||||
interval: "60s"
|
||||
|
||||
influxdb:
|
||||
|
||||
@@ -10,7 +10,7 @@ const pduT = item("PDU")
|
||||
// batchPDU generates Unifi PDU datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchPDU(r report, s *unifi.PDU) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ const ubbT = item("UBB")
|
||||
// - wifi0: 5GHz radio (802.11ac)
|
||||
// - terra2/wlan0/ad: 60GHz radio (802.11ad - Terragraph/WiGig)
|
||||
func (u *DatadogUnifi) batchUBB(r report, s *unifi.UBB) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uciT = item("UCI")
|
||||
// batchUCI generates UCI datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchUCI(r report, s *unifi.UCI) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,7 @@ const udbT = item("UDB")
|
||||
// UDB-Switch is a hybrid device combining switch ports with WiFi 7
|
||||
// wireless bridge capability.
|
||||
func (u *DatadogUnifi) batchUDB(r report, s *unifi.UDB) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -106,7 +106,7 @@ func (u *DatadogUnifi) batchUDMstorage(storage []*unifi.Storage) map[string]floa
|
||||
// batchUDM generates Unifi Gateway datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchUDM(r report, s *unifi.UDM) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const usgT = item("USG")
|
||||
// batchUSG generates Unifi Gateway datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchUSG(r report, s *unifi.USG) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uswT = item("USW")
|
||||
// batchUSW generates Unifi Switch datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchUSW(r report, s *unifi.USW) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uxgT = item("UXG")
|
||||
// batchUXG generates 10Gb Unifi Gateway datapoints for Datadog.
|
||||
// These points can be passed directly to datadog.
|
||||
func (u *DatadogUnifi) batchUXG(r report, s *unifi.UXG) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const pduT = item("PDU")
|
||||
// batchPDU generates Unifi PDU data points for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchPDU(r report, s *unifi.PDU) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -42,7 +42,7 @@ func (u *InfluxUnifi) batchRogueAP(r report, s *unifi.RogueAP) {
|
||||
// batchUAP generates Wireless-Access-Point datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUAP(r report, s *unifi.UAP) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ const ubbT = item("UBB")
|
||||
// - wifi0: 5GHz radio (802.11ac)
|
||||
// - terra2/wlan0/ad: 60GHz radio (802.11ad - Terragraph/WiGig)
|
||||
func (u *InfluxUnifi) batchUBB(r report, s *unifi.UBB) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uciT = item("UCI")
|
||||
// batchUCI generates UCI datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUCI(r report, s *unifi.UCI) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,7 @@ const udbT = item("UDB")
|
||||
// UDB-Switch is a hybrid device combining switch ports with WiFi 7
|
||||
// wireless bridge capability.
|
||||
func (u *InfluxUnifi) batchUDB(r report, s *unifi.UDB) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -81,7 +81,7 @@ func (u *InfluxUnifi) batchUDMstorage(storage []*unifi.Storage) map[string]any {
|
||||
// batchUDM generates Unifi Gateway datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUDM(r report, s *unifi.UDM) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const usgT = item("USG")
|
||||
// batchUSG generates Unifi Gateway datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUSG(r report, s *unifi.USG) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uswT = item("USW")
|
||||
// batchUSW generates Unifi Switch datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUSW(r report, s *unifi.USW) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const uxgT = item("UXG")
|
||||
// batchUXG generates 10Gb Unifi Gateway datapoints for InfluxDB.
|
||||
// These points can be passed directly to influx.
|
||||
func (u *InfluxUnifi) batchUXG(r report, s *unifi.UXG) { // nolint: funlen
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
|
||||
// exportUAP emits metrics for a wireless access point.
|
||||
func (u *OtelOutput) exportUAP(ctx context.Context, meter metric.Meter, r *Report, s *unifi.UAP) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
|
||||
// exportUDM emits metrics for a UniFi Dream Machine (all variants).
|
||||
func (u *OtelOutput) exportUDM(ctx context.Context, meter metric.Meter, r *Report, s *unifi.UDM) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -49,7 +49,7 @@ func (u *OtelOutput) exportUDM(ctx context.Context, meter metric.Meter, r *Repor
|
||||
|
||||
// exportUXG emits metrics for a UniFi Next-Gen Gateway.
|
||||
func (u *OtelOutput) exportUXG(ctx context.Context, meter metric.Meter, r *Report, s *unifi.UXG) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
|
||||
// exportUSG emits metrics for a UniFi Security Gateway.
|
||||
func (u *OtelOutput) exportUSG(ctx context.Context, meter metric.Meter, r *Report, s *unifi.USG) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -11,7 +11,7 @@ import (
|
||||
|
||||
// exportUSW emits metrics for a UniFi switch.
|
||||
func (u *OtelOutput) exportUSW(ctx context.Context, meter metric.Meter, r *Report, s *unifi.USW) {
|
||||
if !s.Adopted.Val || s.Locating.Val {
|
||||
if !s.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -5,10 +5,11 @@ exported metrics. Requires the poller package for actual UniFi data collection.
|
||||
|
||||
## Scrape cache
|
||||
|
||||
Prometheus scrapes are served from an in-memory cache refreshed by a background
|
||||
poller on a fixed interval. This decouples the scrape cadence from the UniFi
|
||||
API call cadence: scrapes always return immediately, and upstream backpressure
|
||||
(e.g. `429 Too Many Requests`) no longer stalls `/metrics`.
|
||||
By default, Prometheus scrapes are served from an in-memory cache refreshed by a
|
||||
background poller on a fixed interval. This decouples the scrape cadence from the
|
||||
UniFi API call cadence: scrapes always return immediately, and upstream backpressure
|
||||
(e.g. `429 Too Many Requests`) no longer stalls `/metrics`. Set `interval = 0` to
|
||||
disable the cache and restore on-demand `/metrics` fetches.
|
||||
|
||||
Config (TOML):
|
||||
|
||||
@@ -16,11 +17,12 @@ Config (TOML):
|
||||
[prometheus]
|
||||
http_listen = "0.0.0.0:9130"
|
||||
# How often the background poller refreshes the cache served to /metrics.
|
||||
# Default: 60s. Values below 15s are clamped to 15s.
|
||||
# Omitted: cache on, 60s. Set to 0 / "0" / "0s" to disable the cache so
|
||||
# /metrics fetches live. Values below 15s log a warning but are not clamped.
|
||||
interval = "60s"
|
||||
```
|
||||
|
||||
Environment variable: `UP_PROMETHEUS_INTERVAL=60s`.
|
||||
Environment variable: `UP_PROMETHEUS_INTERVAL` (same contract; `0` / `0s` disables the cache).
|
||||
|
||||
On poll error the last successful snapshot is preserved, so a transient 429 no
|
||||
longer empties `/metrics`. To monitor cache staleness, scrape the
|
||||
|
||||
+136
-17
@@ -4,6 +4,8 @@ package promunifi
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
@@ -16,6 +18,7 @@ import (
|
||||
"github.com/stretchr/testify/require"
|
||||
"github.com/unpoller/unpoller/pkg/poller"
|
||||
"golift.io/cnfg"
|
||||
"golift.io/cnfgfile"
|
||||
)
|
||||
|
||||
// stubCollect is a minimal poller.Collect implementation for cache tests.
|
||||
@@ -58,6 +61,10 @@ func (s *stubCollect) Logf(string, ...any) {}
|
||||
func (s *stubCollect) LogErrorf(string, ...any) { s.errLogs.Add(1) }
|
||||
func (s *stubCollect) LogDebugf(string, ...any) {}
|
||||
|
||||
func interval(d time.Duration) *cnfg.Duration {
|
||||
return &cnfg.Duration{Duration: d}
|
||||
}
|
||||
|
||||
func TestMetricsCacheSetKeepsLastGoodOnError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
@@ -115,21 +122,55 @@ func TestFetchMetricsReturnsErrorWhenCacheEmpty(t *testing.T) {
|
||||
assert.Error(t, err)
|
||||
}
|
||||
|
||||
// TestFetchMetricsBypassesNilCache covers the nil-cache defensive branch.
|
||||
// Run() always initializes the cache in production; this branch exists only
|
||||
// for tests that exercise fetchMetrics directly without invoking Run.
|
||||
func TestFetchMetricsBypassesNilCache(t *testing.T) {
|
||||
// TestFetchMetricsCacheOffHitsUpstream is the production cache-off path:
|
||||
// Interval = 0 leaves cache nil, so /metrics live-fetches and must not
|
||||
// serve a leftover snapshot from a previous cache.
|
||||
func TestFetchMetricsCacheOffHitsUpstream(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
stub := &stubCollect{metrics: &poller.Metrics{}}
|
||||
u := &promUnifi{Config: &Config{}, Collector: stub}
|
||||
leftover := &poller.Metrics{}
|
||||
live := &poller.Metrics{}
|
||||
stub := &stubCollect{metrics: live}
|
||||
u := &promUnifi{Config: &Config{Interval: interval(0)}, Collector: stub}
|
||||
|
||||
abandoned := &metricsCache{}
|
||||
abandoned.set(leftover, nil)
|
||||
|
||||
u.cache = nil
|
||||
|
||||
got, err := u.fetchMetrics(nil)
|
||||
require.NoError(t, err)
|
||||
assert.Same(t, stub.metrics, got)
|
||||
assert.Same(t, live, got)
|
||||
assert.NotSame(t, leftover, got)
|
||||
assert.EqualValues(t, 1, stub.calls.Load())
|
||||
}
|
||||
|
||||
func TestFetchMetricsCacheOffSingleflightCoalesces(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
stub := &stubCollect{metrics: &poller.Metrics{}, delay: 50 * time.Millisecond}
|
||||
u := &promUnifi{Config: &Config{Interval: interval(0)}, Collector: stub}
|
||||
|
||||
const concurrent = 20
|
||||
|
||||
var wg sync.WaitGroup
|
||||
|
||||
wg.Add(concurrent)
|
||||
|
||||
for i := 0; i < concurrent; i++ {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
|
||||
_, _ = u.fetchMetrics(nil)
|
||||
}()
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
assert.LessOrEqual(t, stub.calls.Load(), int64(2),
|
||||
"singleflight should coalesce concurrent cache-off /metrics scrapes to ~1 upstream call")
|
||||
}
|
||||
|
||||
func TestFetchMetricsSingleflightCoalescesScrapes(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
@@ -186,15 +227,17 @@ func TestNormalizeInterval(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
in time.Duration
|
||||
want time.Duration
|
||||
name string
|
||||
in *cnfg.Duration
|
||||
want time.Duration
|
||||
cacheOn bool
|
||||
}{
|
||||
{"unset uses default", 0, defaultInterval},
|
||||
{"negative uses default", -1, defaultInterval},
|
||||
{"below minimum clamps", time.Second, minimumInterval},
|
||||
{"exact minimum unchanged", minimumInterval, minimumInterval},
|
||||
{"above minimum unchanged", 2 * time.Minute, 2 * time.Minute},
|
||||
{"nil uses default, cache on", nil, defaultInterval, true},
|
||||
{"explicit 0 disables cache", interval(0), 0, false},
|
||||
{"2s kept, cache on", interval(2 * time.Second), 2 * time.Second, true},
|
||||
{"exact minimum unchanged", interval(minimumInterval), minimumInterval, true},
|
||||
{"above minimum unchanged", interval(2 * time.Minute), 2 * time.Minute, true},
|
||||
{"negative uses default, cache on", interval(-time.Second), defaultInterval, true},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
@@ -203,13 +246,89 @@ func TestNormalizeInterval(t *testing.T) {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
u := &promUnifi{Config: &Config{Interval: cnfg.Duration{Duration: tc.in}}}
|
||||
u := &promUnifi{Config: &Config{Interval: tc.in}}
|
||||
u.normalizeInterval()
|
||||
assert.Equal(t, tc.want, u.Interval.Duration)
|
||||
assert.Equal(t, tc.want, u.refreshInterval())
|
||||
assert.Equal(t, tc.cacheOn, u.scrapeCacheEnabled())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestPrometheusIntervalOmittedFromFile(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
require.NoError(t, cnfgfile.Unmarshal(u, writeConfig(t, "up.conf", "[prometheus]\n http_listen = \"0.0.0.0:9130\"\n")))
|
||||
assert.Nil(t, u.Interval)
|
||||
assert.True(t, u.scrapeCacheEnabled())
|
||||
assert.Equal(t, defaultInterval, u.refreshInterval())
|
||||
}
|
||||
|
||||
func TestPrometheusIntervalZeroFromFile(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
require.NoError(t, cnfgfile.Unmarshal(u, writeConfig(t, "up.conf", "[prometheus]\n interval = \"0s\"\n")))
|
||||
require.NotNil(t, u.Interval)
|
||||
assert.Equal(t, time.Duration(0), u.Interval.Duration)
|
||||
assert.False(t, u.scrapeCacheEnabled())
|
||||
}
|
||||
|
||||
func TestPrometheusIntervalTwoSecondsFromFile(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
require.NoError(t, cnfgfile.Unmarshal(u, writeConfig(t, "up.conf", "[prometheus]\n interval = \"2s\"\n")))
|
||||
require.NotNil(t, u.Interval)
|
||||
assert.Equal(t, 2*time.Second, u.Interval.Duration)
|
||||
assert.True(t, u.scrapeCacheEnabled())
|
||||
assert.Equal(t, 2*time.Second, u.refreshInterval())
|
||||
}
|
||||
|
||||
// Not parallel: t.Setenv is incompatible with a parallel test or parent.
|
||||
func TestPrometheusIntervalBindsFromEnvironment(t *testing.T) {
|
||||
t.Run("omitted stays nil", func(t *testing.T) {
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
_, err := cnfg.UnmarshalENV(u, "UP")
|
||||
require.NoError(t, err)
|
||||
assert.Nil(t, u.Interval)
|
||||
assert.True(t, u.scrapeCacheEnabled())
|
||||
assert.Equal(t, defaultInterval, u.refreshInterval())
|
||||
})
|
||||
|
||||
t.Run("UP_PROMETHEUS_INTERVAL=0 disables cache", func(t *testing.T) {
|
||||
t.Setenv("UP_PROMETHEUS_INTERVAL", "0")
|
||||
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
_, err := cnfg.UnmarshalENV(u, "UP")
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, u.Interval)
|
||||
assert.Equal(t, time.Duration(0), u.Interval.Duration)
|
||||
assert.False(t, u.scrapeCacheEnabled())
|
||||
})
|
||||
|
||||
t.Run("UP_PROMETHEUS_INTERVAL=2s", func(t *testing.T) {
|
||||
t.Setenv("UP_PROMETHEUS_INTERVAL", "2s")
|
||||
|
||||
u := &promUnifi{Config: &Config{}}
|
||||
_, err := cnfg.UnmarshalENV(u, "UP")
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, u.Interval)
|
||||
assert.Equal(t, 2*time.Second, u.Interval.Duration)
|
||||
assert.True(t, u.scrapeCacheEnabled())
|
||||
assert.Equal(t, 2*time.Second, u.refreshInterval())
|
||||
})
|
||||
}
|
||||
|
||||
func writeConfig(t *testing.T, name, body string) string {
|
||||
t.Helper()
|
||||
|
||||
path := filepath.Join(t.TempDir(), name)
|
||||
require.NoError(t, os.WriteFile(path, []byte(body), 0o600))
|
||||
|
||||
return path
|
||||
}
|
||||
|
||||
func TestRefreshCachePreservesLastGoodAcrossError(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
+67
-25
@@ -96,9 +96,9 @@ type promUnifi struct {
|
||||
// so operators can alert on failure rate independently of cache staleness.
|
||||
refreshFailures prometheus.Counter
|
||||
// cache holds the last successful metrics snapshot from the background
|
||||
// poller. Run() always initializes it; the nil-guards in cache-using
|
||||
// methods exist only for tests that exercise those methods directly
|
||||
// without invoking Run().
|
||||
// poller. Run() initializes it when the scrape cache is enabled
|
||||
// (Interval omitted or > 0). When Interval is explicitly 0 the cache
|
||||
// stays nil and /metrics fetches live.
|
||||
cache *metricsCache
|
||||
// scrapeFlight coalesces concurrent /scrape requests targeting the same
|
||||
// controller URL so a noisy scraper can't multiply upstream load.
|
||||
@@ -167,10 +167,12 @@ type Config struct {
|
||||
// Interval controls how often the background poller refreshes the cached
|
||||
// metrics that /metrics scrapes are served from. Decouples Prometheus
|
||||
// scrape cadence from upstream UniFi API calls so 429 backoff loops cannot
|
||||
// stall scrapes. Defaults to defaultInterval; values below minimumInterval
|
||||
// are clamped up. Must be > 0 before use; normalizeInterval applies the
|
||||
// default and floor during Run().
|
||||
Interval cnfg.Duration `json:"interval" toml:"interval" xml:"interval" yaml:"interval"`
|
||||
// stall scrapes. Omitted (nil) enables the cache at defaultInterval.
|
||||
// An explicit 0 disables the cache so /metrics fetches live. Any duration
|
||||
// greater than 0 is used as-is; values below minimumInterval log a warning
|
||||
// but are not rewritten. Negative values are invalid and fall back to
|
||||
// defaultInterval.
|
||||
Interval *cnfg.Duration `json:"interval" toml:"interval" xml:"interval" yaml:"interval"`
|
||||
}
|
||||
|
||||
type metric struct {
|
||||
@@ -361,15 +363,19 @@ func (u *promUnifi) Run(c poller.Collect) error {
|
||||
prometheus.MustRegister(u.refreshFailures)
|
||||
prometheus.MustRegister(u)
|
||||
|
||||
u.cache = &metricsCache{}
|
||||
prometheus.MustRegister(u.cacheAgeGauge())
|
||||
// safeRefresh (not refreshCache) because a panic in the initial upstream
|
||||
// fetch must not kill Run() before the HTTP listener starts.
|
||||
u.safeRefresh()
|
||||
if u.scrapeCacheEnabled() {
|
||||
u.cache = &metricsCache{}
|
||||
prometheus.MustRegister(u.cacheAgeGauge())
|
||||
// safeRefresh (not refreshCache) because a panic in the initial upstream
|
||||
// fetch must not kill Run() before the HTTP listener starts.
|
||||
u.safeRefresh()
|
||||
|
||||
go u.backgroundPoll()
|
||||
go u.backgroundPoll()
|
||||
|
||||
u.Logf("Prometheus scrape cache enabled, refresh interval: %v", u.Interval.Duration)
|
||||
u.Logf("Prometheus scrape cache enabled, refresh interval: %v", u.refreshInterval())
|
||||
} else {
|
||||
u.Logf("Prometheus scrape cache disabled; /metrics fetches live")
|
||||
}
|
||||
|
||||
mux.Handle("/metrics", promhttp.HandlerFor(prometheus.DefaultGatherer,
|
||||
promhttp.HandlerOpts{ErrorHandling: promhttp.ContinueOnError},
|
||||
@@ -389,20 +395,48 @@ func (u *promUnifi) Run(c poller.Collect) error {
|
||||
}
|
||||
}
|
||||
|
||||
// normalizeInterval applies defaults and the minimum-interval floor to the
|
||||
// configured scrape cache refresh interval. Values <= 0 use the default.
|
||||
// scrapeCacheEnabled reports whether /metrics is served from the
|
||||
// background-refreshed cache. False only when Interval is explicitly 0.
|
||||
func (u *promUnifi) scrapeCacheEnabled() bool {
|
||||
return u.Interval == nil || u.Interval.Duration != 0
|
||||
}
|
||||
|
||||
// refreshInterval is the cache refresh period. Omitted or negative
|
||||
// values use defaultInterval; any configured duration (including 0) is
|
||||
// returned as-is.
|
||||
func (u *promUnifi) refreshInterval() time.Duration {
|
||||
if u.Interval == nil || u.Interval.Duration < 0 {
|
||||
return defaultInterval
|
||||
}
|
||||
|
||||
return u.Interval.Duration
|
||||
}
|
||||
|
||||
// normalizeInterval applies the scrape-cache interval contract: omitted/nil
|
||||
// keeps the default (cache on, 60s); explicit 0 disables the cache; negative
|
||||
// values are invalid and fall back to the default; values below
|
||||
// minimumInterval log a warning but are not rewritten.
|
||||
func (u *promUnifi) normalizeInterval() {
|
||||
if u.Interval.Duration <= 0 {
|
||||
if u.Interval == nil {
|
||||
return
|
||||
}
|
||||
|
||||
if u.Interval.Duration < 0 {
|
||||
u.Logf("Prometheus interval %v is invalid; using default %v",
|
||||
u.Interval.Duration, defaultInterval)
|
||||
|
||||
u.Interval.Duration = defaultInterval
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
if u.Interval.Duration < minimumInterval {
|
||||
u.Logf("Prometheus interval %v is below minimum %v; clamping to minimum",
|
||||
u.Interval.Duration, minimumInterval)
|
||||
if u.Interval.Duration == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
u.Interval.Duration = minimumInterval
|
||||
if u.Interval.Duration < minimumInterval {
|
||||
u.Logf("Prometheus interval %v is below recommended minimum %v; this may increase UniFi API load",
|
||||
u.Interval.Duration, minimumInterval)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -415,7 +449,7 @@ func (u *promUnifi) backgroundPoll() {
|
||||
return
|
||||
}
|
||||
|
||||
ticker := time.NewTicker(u.Interval.Duration)
|
||||
ticker := time.NewTicker(u.refreshInterval())
|
||||
defer ticker.Stop()
|
||||
|
||||
for range ticker.C {
|
||||
@@ -501,12 +535,14 @@ func (u *promUnifi) cacheAgeGauge() prometheus.Collector {
|
||||
}
|
||||
|
||||
// fetchMetrics returns the metrics for a scrape, using the cache for global
|
||||
// /metrics scrapes and singleflight-coalesced live calls for per-target
|
||||
// /scrape requests.
|
||||
// /metrics scrapes when the cache is enabled. A nil cache (interval = 0)
|
||||
// live-fetches under a singleflight key so concurrent /metrics scrapes
|
||||
// cannot multiply upstream load. Per-target /scrape requests always
|
||||
// live-fetch, coalesced by filter path.
|
||||
func (u *promUnifi) fetchMetrics(filter *poller.Filter) (*poller.Metrics, error) {
|
||||
if filter == nil {
|
||||
if u.cache == nil {
|
||||
return u.Collector.Metrics(nil)
|
||||
return u.liveMetrics("global", nil)
|
||||
}
|
||||
|
||||
m, _, err := u.cache.get()
|
||||
@@ -529,6 +565,12 @@ func (u *promUnifi) fetchMetrics(filter *poller.Filter) (*poller.Metrics, error)
|
||||
key = filter.Name
|
||||
}
|
||||
|
||||
return u.liveMetrics(key, filter)
|
||||
}
|
||||
|
||||
// liveMetrics fetches from the collector once per singleflight key so
|
||||
// concurrent scrapes (cache-off /metrics or /scrape) share one upstream call.
|
||||
func (u *promUnifi) liveMetrics(key string, filter *poller.Filter) (*poller.Metrics, error) {
|
||||
result, err, _ := u.scrapeFlight.Do(key, func() (any, error) {
|
||||
return u.Collector.Metrics(filter)
|
||||
})
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
//nolint:testpackage // white-box: exercises the unexported export skip path.
|
||||
package promunifi
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/unpoller/unifi/v6"
|
||||
)
|
||||
|
||||
// Locate/identify only blinks the LED. An adopted device in that state must
|
||||
// still scrape; dropping it was the silent Grafana gap in #1075.
|
||||
func TestExportUCILocatingStillExports(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
r := &fakeReport{}
|
||||
u := &promUnifi{Device: descDevice("unifi_")}
|
||||
u.exportUCI(r, &unifi.UCI{
|
||||
Adopted: unifi.FlexBool{Val: true},
|
||||
Locating: unifi.FlexBool{Val: true},
|
||||
Type: "uci",
|
||||
SiteName: "default",
|
||||
Name: "UCI-1",
|
||||
SourceName: "https://controller.example",
|
||||
})
|
||||
|
||||
assert.NotEmpty(t, r.sent, "an adopted device in locate mode must still export metrics")
|
||||
}
|
||||
|
||||
func TestExportUCIUnadoptedIsSkipped(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
r := &fakeReport{}
|
||||
u := &promUnifi{Device: descDevice("unifi_")}
|
||||
u.exportUCI(r, &unifi.UCI{
|
||||
Adopted: unifi.FlexBool{Val: false},
|
||||
Locating: unifi.FlexBool{Val: true},
|
||||
Type: "uci",
|
||||
SiteName: "default",
|
||||
Name: "UCI-1",
|
||||
SourceName: "https://controller.example",
|
||||
})
|
||||
|
||||
assert.Empty(t, r.sent, "unadopted devices remain excluded")
|
||||
}
|
||||
@@ -163,7 +163,7 @@ func descPDU(ns string) *pdu {
|
||||
}
|
||||
|
||||
func (u *promUnifi) exportPDU(r report, d *unifi.PDU) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -211,7 +211,7 @@ func (u *promUnifi) exportRogueAP(r report, d *unifi.RogueAP) {
|
||||
}
|
||||
|
||||
func (u *promUnifi) exportUAP(r report, d *unifi.UAP) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,7 @@ import (
|
||||
// - wifi0: 5GHz radio (802.11ac)
|
||||
// - terra2/wlan0/ad: 60GHz radio (802.11ad - Terragraph/WiGig)
|
||||
func (u *promUnifi) exportUBB(r report, d *unifi.UBB) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
|
||||
// exportUCI is a collection of stats from UCI.
|
||||
func (u *promUnifi) exportUCI(r report, d *unifi.UCI) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@ import "github.com/unpoller/unifi/v6"
|
||||
// UDB-Switch is a hybrid device combining switch ports (8 PoE ports)
|
||||
// with WiFi 7 wireless bridge capability (5GHz + 6GHz radios).
|
||||
func (u *promUnifi) exportUDB(r report, d *unifi.UDB) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -79,7 +79,7 @@ func descDevice(ns string) *unifiDevice {
|
||||
|
||||
// UDM is a collection of stats from USG, USW and UAP. It has no unique stats.
|
||||
func (u *promUnifi) exportUDM(r report, d *unifi.UDM) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -78,7 +78,7 @@ func descUSG(ns string) *usg {
|
||||
}
|
||||
|
||||
func (u *promUnifi) exportUSG(r report, d *unifi.USG) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -112,7 +112,7 @@ func descUSW(ns string) *usw {
|
||||
}
|
||||
|
||||
func (u *promUnifi) exportUSW(r report, d *unifi.USW) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
|
||||
// exportUXG is a collection of stats from USG and USW. It has no unique stats.
|
||||
func (u *promUnifi) exportUXG(r report, d *unifi.UXG) {
|
||||
if !d.Adopted.Val || d.Locating.Val {
|
||||
if !d.Adopted.Val {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user