feat(metrics): publish felis_servers_total from a fleet snapshot
A per-object reconcile cannot maintain felis_servers_total (spec §23): it sees one server per call, so it could never Set a correct fleet-wide gauge and inc/dec on transitions would drift on any missed event. Add a snapshot producer instead. metrics.SyncServerGauge Resets the GaugeVec then Sets one child per state, so a state that drains to zero reports 0 rather than a stale last value. operator.GaugeSyncer is a manager.Runnable that periodically Lists the fleet and republishes from it, defaulting an unset desiredState to Stopped. SyncOnce is exercised end-to-end against a fake client (List, default, republish); the ticker loop in Start is the only untested I/O edge.
This commit is contained in:
5 files changed
+194
No files matched your search
@@ -79,6 +79,15 @@ func cmdOperator(args []string, _, stderr io.Writer) int {
|
|||||||
return 1
|
return 1
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Republish felis_servers_total from a periodic full List of the fleet. A
|
||||||
|
// per-object reconcile can never maintain a fleet-wide gauge correctly, so a
|
||||||
|
// snapshot Runnable owns it; it shares the manager's cached client and stops
|
||||||
|
// with the manager.
|
||||||
|
if err := mgr.Add(&operator.GaugeSyncer{Client: mgr.GetClient()}); err != nil {
|
||||||
|
fmt.Fprintf(stderr, "felis operator: add gauge syncer: %v\n", err)
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
|
||||||
if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
|
if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
|
||||||
fmt.Fprintf(stderr, "felis operator: manager exited: %v\n", err)
|
fmt.Fprintf(stderr, "felis operator: manager exited: %v\n", err)
|
||||||
return 1
|
return 1
|
||||||
|
|||||||
@@ -57,6 +57,27 @@ var (
|
|||||||
})
|
})
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// SyncServerGauge republishes felis_servers_total from a full snapshot of the
|
||||||
|
// fleet's per-server states. states holds one entry per MinecraftServer the
|
||||||
|
// operator knows about (its desiredState).
|
||||||
|
//
|
||||||
|
// It Resets the GaugeVec before Setting one child per distinct state, so a state
|
||||||
|
// that has drained to zero reports 0 rather than its stale last value. That is
|
||||||
|
// the whole reason a periodic full-snapshot is used instead of inc/dec on
|
||||||
|
// reconcile transitions: a snapshot is self-correcting and cannot drift on a
|
||||||
|
// missed event. Producing the states slice (a cached List of MinecraftServers)
|
||||||
|
// is the untestable I/O edge; this Reset+tally+Set logic is pure and unit-tested.
|
||||||
|
func SyncServerGauge(states []string) {
|
||||||
|
ServersTotal.Reset()
|
||||||
|
counts := make(map[string]int, len(states))
|
||||||
|
for _, s := range states {
|
||||||
|
counts[s]++
|
||||||
|
}
|
||||||
|
for state, n := range counts {
|
||||||
|
ServersTotal.WithLabelValues(state).Set(float64(n))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Collectors returns every felis_* collector in a stable order. Production and
|
// Collectors returns every felis_* collector in a stable order. Production and
|
||||||
// tests register the same slice, so the test asserting the full set is exposed
|
// tests register the same slice, so the test asserting the full set is exposed
|
||||||
// also pins the production surface.
|
// also pins the production surface.
|
||||||
|
|||||||
@@ -54,6 +54,30 @@ func TestRegisterExposesNamedFelisMetrics(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSyncServerGaugeResetsStaleStates(t *testing.T) {
|
||||||
|
// Self-contained (Reset->sync->assert within this one test) so it never
|
||||||
|
// clobbers another test's ServersTotal children. The correctness point is the
|
||||||
|
// Reset inside SyncServerGauge: a state present in one snapshot but absent from
|
||||||
|
// the next must drain to 0, not retain its stale last value forever.
|
||||||
|
SyncServerGauge([]string{"Running", "Running", "Stopped"})
|
||||||
|
if got := testutil.ToFloat64(ServersTotal.WithLabelValues("Running")); got != 2 {
|
||||||
|
t.Errorf("Running = %v, want 2", got)
|
||||||
|
}
|
||||||
|
if got := testutil.ToFloat64(ServersTotal.WithLabelValues("Stopped")); got != 1 {
|
||||||
|
t.Errorf("Stopped = %v, want 1", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Next snapshot: the two Running servers are gone. Without the Reset the gauge
|
||||||
|
// would still report Running=2; with it, the child drains to 0.
|
||||||
|
SyncServerGauge([]string{"Stopped"})
|
||||||
|
if got := testutil.ToFloat64(ServersTotal.WithLabelValues("Running")); got != 0 {
|
||||||
|
t.Errorf("Running after drain = %v, want 0", got)
|
||||||
|
}
|
||||||
|
if got := testutil.ToFloat64(ServersTotal.WithLabelValues("Stopped")); got != 1 {
|
||||||
|
t.Errorf("Stopped = %v, want 1", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestRegisterIsIdempotent(t *testing.T) {
|
func TestRegisterIsIdempotent(t *testing.T) {
|
||||||
reg := prometheus.NewRegistry()
|
reg := prometheus.NewRegistry()
|
||||||
if err := Register(reg); err != nil {
|
if err := Register(reg); err != nil {
|
||||||
|
|||||||
@@ -0,0 +1,90 @@
|
|||||||
|
package operator
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||||
|
"felis.lolicon.best/internal/metrics"
|
||||||
|
ctrl "sigs.k8s.io/controller-runtime"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||||
|
)
|
||||||
|
|
||||||
|
// defaultGaugeInterval is the republish cadence used when none is configured.
|
||||||
|
const defaultGaugeInterval = 30 * time.Second
|
||||||
|
|
||||||
|
// GaugeSyncer periodically republishes felis_servers_total (spec §23) from a
|
||||||
|
// full List of MinecraftServers.
|
||||||
|
//
|
||||||
|
// felis_servers_total is a fleet-wide gauge partitioned by desiredState, which a
|
||||||
|
// per-object Reconcile fundamentally cannot maintain: one reconcile observes a
|
||||||
|
// single server, so it could never Set a correct fleet count and inc/dec on
|
||||||
|
// transitions would drift on any missed event. A periodic full-snapshot sidesteps
|
||||||
|
// that — every tick re-derives the entire gauge from the live list (and
|
||||||
|
// metrics.SyncServerGauge Resets first), so the published value is always exactly
|
||||||
|
// the current fleet and never accumulates stale state.
|
||||||
|
//
|
||||||
|
// It is registered as a manager.Runnable (mgr.Add), sharing the manager's cached,
|
||||||
|
// namespace-scoped client and lifecycle; it stops when the manager's context is
|
||||||
|
// cancelled.
|
||||||
|
type GaugeSyncer struct {
|
||||||
|
// Client is the manager's cached client (namespace-scoped like the operator).
|
||||||
|
Client client.Client
|
||||||
|
// Interval is the republish cadence; defaults to defaultGaugeInterval.
|
||||||
|
Interval time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
// Start runs the republish loop until ctx is cancelled, satisfying
|
||||||
|
// manager.Runnable. The ticker loop and the live List are the untestable I/O
|
||||||
|
// edge; the List->translate->gauge work it delegates to SyncOnce is exercised
|
||||||
|
// against a fake client. A failed sync is logged best-effort, not fatal — the
|
||||||
|
// next tick re-derives the gauge from scratch regardless.
|
||||||
|
func (g *GaugeSyncer) Start(ctx context.Context) error {
|
||||||
|
interval := g.Interval
|
||||||
|
if interval <= 0 {
|
||||||
|
interval = defaultGaugeInterval
|
||||||
|
}
|
||||||
|
t := time.NewTicker(interval)
|
||||||
|
defer t.Stop()
|
||||||
|
// Publish once up front so the gauge is populated before the first tick.
|
||||||
|
if err := g.SyncOnce(ctx); err != nil {
|
||||||
|
ctrl.LoggerFrom(ctx).Error(err, "gauge: initial fleet sync failed")
|
||||||
|
}
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return nil
|
||||||
|
case <-t.C:
|
||||||
|
if err := g.SyncOnce(ctx); err != nil {
|
||||||
|
ctrl.LoggerFrom(ctx).Error(err, "gauge: fleet sync failed")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SyncOnce Lists the fleet once and republishes felis_servers_total from it.
|
||||||
|
// Separated from Start so the List->translate->gauge path is unit-testable
|
||||||
|
// against a fake client, leaving only the ticker loop untested.
|
||||||
|
func (g *GaugeSyncer) SyncOnce(ctx context.Context) error {
|
||||||
|
var list v1alpha1.MinecraftServerList
|
||||||
|
if err := g.Client.List(ctx, &list); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
metrics.SyncServerGauge(serverStates(list.Items))
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// serverStates maps each server to its desiredState, defaulting an unset state
|
||||||
|
// to Stopped (the CRD default, spec §4). Pure, so the labelling rule is covered
|
||||||
|
// by SyncOnce's fake-client test rather than only at runtime.
|
||||||
|
func serverStates(items []v1alpha1.MinecraftServer) []string {
|
||||||
|
states := make([]string, 0, len(items))
|
||||||
|
for i := range items {
|
||||||
|
s := string(items[i].Spec.DesiredState)
|
||||||
|
if s == "" {
|
||||||
|
s = string(v1alpha1.DesiredStopped)
|
||||||
|
}
|
||||||
|
states = append(states, s)
|
||||||
|
}
|
||||||
|
return states
|
||||||
|
}
|
||||||
@@ -0,0 +1,50 @@
|
|||||||
|
package operator_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"felis.lolicon.best/internal/apis/felis/v1alpha1"
|
||||||
|
"felis.lolicon.best/internal/metrics"
|
||||||
|
"felis.lolicon.best/internal/operator"
|
||||||
|
"github.com/prometheus/client_golang/prometheus/testutil"
|
||||||
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||||
|
"sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||||
|
)
|
||||||
|
|
||||||
|
func gaugeServer(name string, state v1alpha1.DesiredState) *v1alpha1.MinecraftServer {
|
||||||
|
return &v1alpha1.MinecraftServer{
|
||||||
|
ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "minecraft"},
|
||||||
|
Spec: v1alpha1.MinecraftServerSpec{Subdomain: name, DesiredState: state},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestGaugeSyncerPublishesFleetStateFromCluster drives SyncOnce against a fake
|
||||||
|
// client to verify the whole production path bar the ticker: List the fleet,
|
||||||
|
// translate each server to its desiredState (an unset state defaulting to
|
||||||
|
// Stopped), and republish felis_servers_total partitioned by that state.
|
||||||
|
func TestGaugeSyncerPublishesFleetStateFromCluster(t *testing.T) {
|
||||||
|
scheme := newScheme(t)
|
||||||
|
c := fake.NewClientBuilder().
|
||||||
|
WithScheme(scheme).
|
||||||
|
WithObjects(
|
||||||
|
gaugeServer("a", v1alpha1.DesiredRunning),
|
||||||
|
gaugeServer("b", v1alpha1.DesiredRunning),
|
||||||
|
gaugeServer("c", v1alpha1.DesiredStopped),
|
||||||
|
gaugeServer("d", ""), // unset desiredState must be counted as Stopped
|
||||||
|
).
|
||||||
|
Build()
|
||||||
|
|
||||||
|
g := &operator.GaugeSyncer{Client: c}
|
||||||
|
if err := g.SyncOnce(context.Background()); err != nil {
|
||||||
|
t.Fatalf("SyncOnce: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := testutil.ToFloat64(metrics.ServersTotal.WithLabelValues("Running")); got != 2 {
|
||||||
|
t.Errorf("servers_total{state=Running} = %v, want 2", got)
|
||||||
|
}
|
||||||
|
// One explicit Stopped plus one unset (defaulted to Stopped) = 2.
|
||||||
|
if got := testutil.ToFloat64(metrics.ServersTotal.WithLabelValues("Stopped")); got != 2 {
|
||||||
|
t.Errorf("servers_total{state=Stopped} = %v, want 2", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user