diff --git a/cmd/felis/operator.go b/cmd/felis/operator.go index 2ac0e92..38e2a46 100644 --- a/cmd/felis/operator.go +++ b/cmd/felis/operator.go @@ -79,6 +79,15 @@ func cmdOperator(args []string, _, stderr io.Writer) int { 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 { fmt.Fprintf(stderr, "felis operator: manager exited: %v\n", err) return 1 diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go index 805eb11..7db2331 100644 --- a/internal/metrics/metrics.go +++ b/internal/metrics/metrics.go @@ -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 // tests register the same slice, so the test asserting the full set is exposed // also pins the production surface. diff --git a/internal/metrics/metrics_test.go b/internal/metrics/metrics_test.go index 5e27b2a..42bbade 100644 --- a/internal/metrics/metrics_test.go +++ b/internal/metrics/metrics_test.go @@ -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) { reg := prometheus.NewRegistry() if err := Register(reg); err != nil { diff --git a/internal/operator/gauge.go b/internal/operator/gauge.go new file mode 100644 index 0000000..ca51b6a --- /dev/null +++ b/internal/operator/gauge.go @@ -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 +} diff --git a/internal/operator/gauge_test.go b/internal/operator/gauge_test.go new file mode 100644 index 0000000..75f392a --- /dev/null +++ b/internal/operator/gauge_test.go @@ -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) + } +}