Merge branch 'feature/ops-03-metrics-aggregation-ueber-module-hinweg' into feature/qa-05-abnahme-compliance-pruefung-core

This commit is contained in:
sysops
2026-08-29 17:05:03 +02:00
8 changed files with 529 additions and 0 deletions
+43
View File
@@ -0,0 +1,43 @@
// metrics-devserver stellt den OPS-03-Metrics-Aggregator (internal/metrics)
// unter /metrics bereit, damit ein echter Prometheus-Scrape-Vorgang gegen
// den Core-Dienst geprueft werden kann (Pruefung 3). Getrennt von cmd/core
// aus demselben Grund wie die anderen *-devserver.
package main
import (
"context"
"log"
"net/http"
"os"
"gitea.perlbach24.de/scripte/nexarch/internal/db"
"gitea.perlbach24.de/scripte/nexarch/internal/metrics"
)
func main() {
dsn := os.Getenv("NEXARCH_REGISTRY_DSN")
if dsn == "" {
log.Fatal("NEXARCH_REGISTRY_DSN nicht gesetzt")
}
addr := os.Getenv("NEXARCH_METRICS_LISTEN_ADDR")
if addr == "" {
addr = ":8085"
}
ctx := context.Background()
pool, err := db.Connect(ctx, dsn)
if err != nil {
log.Fatalf("db: %v", err)
}
defer pool.Close()
sourceStore := metrics.NewSourceStore(pool)
agg := metrics.NewAggregator(metrics.NewCoreRegistry(), sourceStore.Provide)
mux := http.NewServeMux()
mux.HandleFunc("/metrics", agg.Handler())
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })
log.Printf("metrics-devserver listening on %s", addr)
log.Fatal(http.ListenAndServe(addr, mux))
}
+19
View File
@@ -0,0 +1,19 @@
package metrics
import "github.com/prometheus/client_golang/prometheus"
// NewCoreRegistry liefert das Prometheus-Registry fuer die EIGENEN
// Kennzahlen des Core-Dienstes (Akzeptanzkriterium 1) — alle Namen tragen
// das Praefix "nexarch_core_" gemaess der im Paketkommentar dokumentierten
// Namenskonvention (Akzeptanzkriterium 3). Ein eigenes Registry statt des
// globalen DefaultRegisterer, damit Tests unabhaengig voneinander sind.
func NewCoreRegistry() *prometheus.Registry {
reg := prometheus.NewRegistry()
reg.MustRegister(
prometheus.NewGaugeFunc(prometheus.GaugeOpts{
Name: "nexarch_core_up",
Help: "1, solange der Core-Dienst laeuft und Metriken liefern kann.",
}, func() float64 { return 1 }),
)
return reg
}
+168
View File
@@ -0,0 +1,168 @@
// Package metrics implementiert Core OPS-03: einen zentralen /metrics-
// Endpunkt im Prometheus-Textformat, der Kennzahlen des Core-Dienstes UND
// aggregierte Kennzahlen aller registrierten Module bereitstellt — offenes
// Pull-Modell nach Prometheus-Vorbild, kein proprietaerer Push-Mechanismus.
//
// Namenskonvention (Akzeptanzkriterium 3, modulübergreifend konsistent):
//
// nexarch_core_<name> — Kennzahlen des Core-Dienstes selbst
// nexarch_module_<modul>_<name> — von einem Modul gescrapte Kennzahl
// <name>, umbenannt mit dem
// Modulnamen als Praefix
//
// Ein Modul liefert seine eigenen Kennzahlen unter EIGENEM Namen (z.B.
// "requests_total") unter seinem eigenen /metrics-Endpunkt — dieses Paket
// benennt sie beim Einsammeln konsistent um, damit im aggregierten Core-
// Endpunkt niemals zwei Module denselben Metrik-Namen kollidieren lassen.
package metrics
import (
"context"
"fmt"
"net/http"
"time"
"github.com/prometheus/client_golang/prometheus"
dto "github.com/prometheus/client_model/go"
"github.com/prometheus/common/expfmt"
"github.com/prometheus/common/model"
)
// init erzwingt das klassische Prometheus-Namensschema (a-z, A-Z, 0-9, _)
// fuer die Namensvalidierung von expfmt/model — ohne diese explizite
// Festlegung liefert die Bibliothek "Invalid name validation scheme
// requested: unset" beim Parsen/Kodieren, da sie den globalen Default in
// dieser Version nicht mehr implizit setzt.
func init() {
model.NameValidationScheme = model.LegacyValidation
}
// Source ist EIN registriertes Modul mit seinem eigenen /metrics-Endpunkt
// (siehe internal/health fuer das analoge Muster bei Readiness-Checks).
type Source struct {
ModuleName string
MetricsURL string
}
// SourceProvider liefert die aktuell registrierten Module — typischerweise
// rueckgebunden an internal/moduleregistry.Registry.List (API-02) ueber
// einen kleinen Adapter im aufrufenden Code, damit dieses Paket
// internal/moduleregistry nicht direkt importieren muss (Kein Umbau
// angrenzender Bereiche). Ein NEU registriertes Modul erscheint automatisch
// beim naechsten Aufruf von Aggregator.Handler, OHNE Codeaenderung an diesem
// Paket (Akzeptanzkriterium 2 / Pruefung 2).
type SourceProvider func(ctx context.Context) ([]Source, error)
// FetchTimeout begrenzt, wie lange EIN Modul-Scrape maximal dauern darf —
// ein haengendes Modul darf den gesamten Aggregations-Request nicht
// verzoegern (Pruefung 1: Antwort unter Last innerhalb definierter Zeit).
const FetchTimeout = 2 * time.Second
// Aggregator sammelt Core-eigene Metriken (coreGatherer) und die Metriken
// aller ueber sourceProvider gemeldeten Module in EINER Antwort ein.
type Aggregator struct {
coreGatherer prometheus.Gatherer
sourceProvider SourceProvider
client *http.Client
}
func NewAggregator(coreGatherer prometheus.Gatherer, sourceProvider SourceProvider) *Aggregator {
return &Aggregator{
coreGatherer: coreGatherer,
sourceProvider: sourceProvider,
client: &http.Client{Timeout: FetchTimeout},
}
}
// Gather implementiert prometheus.Gatherer: liefert Core-Metriken PLUS alle
// erreichbaren Modul-Metriken (umbenannt gemaess Namenskonvention) in einer
// gemeinsamen Liste von MetricFamilies.
func (a *Aggregator) Gather(ctx context.Context) ([]*dto.MetricFamily, error) {
families, err := a.coreGatherer.Gather()
if err != nil {
return nil, fmt.Errorf("core-metriken einsammeln: %w", err)
}
sources, err := a.sourceProvider(ctx)
if err != nil {
return nil, fmt.Errorf("modul-quellen ermitteln: %w", err)
}
// Module werden NEBENLAEUFIG gescrapt (dasselbe Muster wie
// internal/health.Registry.CheckAll) — ein langsames/nicht erreichbares
// Modul haelt weder andere Module noch den Gesamt-Request auf.
type fetchResult struct {
families []*dto.MetricFamily
}
resultCh := make(chan fetchResult, len(sources))
for _, src := range sources {
go func(src Source) {
fetchCtx, cancel := context.WithTimeout(ctx, FetchTimeout)
defer cancel()
mf, err := a.fetchAndRename(fetchCtx, src)
if err != nil {
resultCh <- fetchResult{} // Fehlerfall: einfach nichts beitragen, Aggregation laeuft weiter
return
}
resultCh <- fetchResult{families: mf}
}(src)
}
for range sources {
r := <-resultCh
families = append(families, r.families...)
}
return families, nil
}
// fetchAndRename ruft die /metrics-URL eines Moduls ab, parst das
// Prometheus-Textformat und benennt jede Metrik gemaess der
// Namenskonvention um (Akzeptanzkriterium 2 + 3).
func (a *Aggregator) fetchAndRename(ctx context.Context, src Source) ([]*dto.MetricFamily, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, src.MetricsURL, nil)
if err != nil {
return nil, err
}
resp, err := a.client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("modul %s: unerwarteter status %d", src.ModuleName, resp.StatusCode)
}
parser := expfmt.NewTextParser(model.LegacyValidation)
parsed, err := parser.TextToMetricFamilies(resp.Body)
if err != nil {
return nil, fmt.Errorf("modul %s: metrik-text nicht parsebar: %w", src.ModuleName, err)
}
out := make([]*dto.MetricFamily, 0, len(parsed))
for name, mf := range parsed {
renamed := fmt.Sprintf("nexarch_module_%s_%s", src.ModuleName, name)
mf.Name = &renamed
out = append(out, mf)
}
return out, nil
}
// Handler liefert einen HTTP-Handler, der Gather aufruft und das Ergebnis im
// Prometheus-Textformat ausgibt (Akzeptanzkriterium 1).
func (a *Aggregator) Handler() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
families, err := a.Gather(r.Context())
if err != nil {
http.Error(w, "metriken konnten nicht eingesammelt werden", http.StatusInternalServerError)
return
}
w.Header().Set("Content-Type", string(expfmt.NewFormat(expfmt.TypeTextPlain)))
enc := expfmt.NewEncoder(w, expfmt.NewFormat(expfmt.TypeTextPlain))
for _, mf := range families {
if err := enc.Encode(mf); err != nil {
return
}
}
}
}
+179
View File
@@ -0,0 +1,179 @@
package metrics
import (
"context"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
"github.com/prometheus/common/expfmt"
"github.com/prometheus/common/model"
)
func fakeModuleServer(metricName string) *httptest.Server {
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/plain; version=0.0.4")
fmt.Fprintf(w, "# HELP %s ein test-zaehler\n# TYPE %s counter\n%s 42\n", metricName, metricName, metricName)
}))
}
// Akzeptanzkriterium 1: Core liefert unter dem Handler valides
// Prometheus-Textformat mit den eigenen Kennzahlen.
func TestHandler_ServesCoreMetricsInPrometheusFormat(t *testing.T) {
agg := NewAggregator(NewCoreRegistry(), func(ctx context.Context) ([]Source, error) { return nil, nil })
req := httptest.NewRequest(http.MethodGet, "/metrics", nil)
rec := httptest.NewRecorder()
agg.Handler()(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", rec.Code)
}
parser := expfmt.NewTextParser(model.LegacyValidation)
families, err := parser.TextToMetricFamilies(strings.NewReader(rec.Body.String()))
if err != nil {
t.Fatalf("antwort ist kein valides prometheus-textformat: %v", err)
}
if _, ok := families["nexarch_core_up"]; !ok {
t.Fatalf("erwartet 'nexarch_core_up' unter den core-metriken, habe: %v", keysOf(families))
}
}
func keysOf[V any](m map[string]V) []string {
out := make([]string, 0, len(m))
for k := range m {
out = append(out, k)
}
return out
}
// Akzeptanzkriterium 2 + Pruefung 2: ein NEU registriertes Modul erscheint
// in der Aggregation, OHNE dass dieses Paket oder der Aufrufer Code
// aendern muss — die Quelle kommt ausschliesslich aus sourceProvider.
func TestHandler_NewlyRegisteredModuleAppearsWithoutCodeChange(t *testing.T) {
moduleServer := fakeModuleServer("requests_total")
defer moduleServer.Close()
// Simuliert eine sich zur Laufzeit aendernde Modul-Liste (z.B. aus
// SourceStore.Provide) — zunaechst LEER, dann mit einem Eintrag.
var sources []Source
var mu sync.Mutex
provider := func(ctx context.Context) ([]Source, error) {
mu.Lock()
defer mu.Unlock()
out := make([]Source, len(sources))
copy(out, sources)
return out, nil
}
agg := NewAggregator(NewCoreRegistry(), provider)
// Vor der Registrierung: Modul-Metrik nicht vorhanden.
rec1 := httptest.NewRecorder()
agg.Handler()(rec1, httptest.NewRequest(http.MethodGet, "/metrics", nil))
if strings.Contains(rec1.Body.String(), "requests_total") {
t.Fatal("modul-metrik haette vor registrierung nicht erscheinen duerfen")
}
// Modul wird "registriert" (kein Code hier oder in metrics.go aendert sich).
mu.Lock()
sources = append(sources, Source{ModuleName: "dms", MetricsURL: moduleServer.URL})
mu.Unlock()
rec2 := httptest.NewRecorder()
agg.Handler()(rec2, httptest.NewRequest(http.MethodGet, "/metrics", nil))
body := rec2.Body.String()
if !strings.Contains(body, "nexarch_module_dms_requests_total") {
t.Fatalf("erwartet umbenannte modul-metrik 'nexarch_module_dms_requests_total' nach registrierung, body:\n%s", body)
}
}
// Akzeptanzkriterium 3: Namenskonvention "nexarch_module_<modul>_<name>"
// wird tatsaechlich angewendet.
func TestFetchAndRename_AppliesNamingConvention(t *testing.T) {
moduleServer := fakeModuleServer("queue_depth")
defer moduleServer.Close()
agg := NewAggregator(NewCoreRegistry(), nil)
families, err := agg.fetchAndRename(context.Background(), Source{ModuleName: "mail", MetricsURL: moduleServer.URL})
if err != nil {
t.Fatalf("fetchAndRename: %v", err)
}
if len(families) != 1 || families[0].GetName() != "nexarch_module_mail_queue_depth" {
t.Fatalf("erwartet genau 1 metrik 'nexarch_module_mail_queue_depth', habe: %+v", families)
}
}
// Ein nicht erreichbares Modul darf die Aggregation der uebrigen und die
// Gesamtantwort nicht verhindern (dieselbe Resilienz wie OPS-02).
func TestHandler_UnreachableModuleDoesNotBreakAggregation(t *testing.T) {
reachable := fakeModuleServer("healthy_metric")
defer reachable.Close()
unreachable := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {}))
unreachableURL := unreachable.URL
unreachable.Close() // sofort schliessen -> Verbindung schlaegt fehl
provider := func(ctx context.Context) ([]Source, error) {
return []Source{
{ModuleName: "ok", MetricsURL: reachable.URL},
{ModuleName: "kaputt", MetricsURL: unreachableURL},
}, nil
}
agg := NewAggregator(NewCoreRegistry(), provider)
rec := httptest.NewRecorder()
agg.Handler()(rec, httptest.NewRequest(http.MethodGet, "/metrics", nil))
if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200 trotz einem nicht erreichbaren modul", rec.Code)
}
body := rec.Body.String()
if !strings.Contains(body, "nexarch_module_ok_healthy_metric") {
t.Fatal("erreichbares modul haette trotz ausfall des anderen aggregiert werden sollen")
}
if strings.Contains(body, "kaputt") {
t.Fatal("nicht erreichbares modul haette keine metrik beitragen duerfen")
}
}
// Pruefung 1: Endpunkt antwortet unter mehreren gleichzeitigen Anfragen
// innerhalb definierter Zeit — kein unbeschraenktes Blockieren durch
// langsame Module (FetchTimeout begrenzt jeden Scrape).
func TestHandler_RespondsWithinBoundedTimeUnderLoad(t *testing.T) {
hangingServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
time.Sleep(10 * time.Second) // wuerde ohne timeout jede anfrage blockieren
}))
defer hangingServer.Close()
provider := func(ctx context.Context) ([]Source, error) {
return []Source{{ModuleName: "haengend", MetricsURL: hangingServer.URL}}, nil
}
agg := NewAggregator(NewCoreRegistry(), provider)
const concurrentRequests = 10
var wg sync.WaitGroup
start := time.Now()
for i := 0; i < concurrentRequests; i++ {
wg.Add(1)
go func() {
defer wg.Done()
rec := httptest.NewRecorder()
agg.Handler()(rec, httptest.NewRequest(http.MethodGet, "/metrics", nil))
if rec.Code != http.StatusOK {
t.Errorf("status = %d, want 200", rec.Code)
}
}()
}
wg.Wait()
elapsed := time.Since(start)
if elapsed > FetchTimeout+3*time.Second {
t.Fatalf("%d gleichzeitige anfragen brauchten %s, erwartet deutlich unter %s durch FetchTimeout",
concurrentRequests, elapsed, FetchTimeout+3*time.Second)
}
}
+52
View File
@@ -0,0 +1,52 @@
package metrics
import (
"context"
"fmt"
"github.com/jackc/pgx/v5/pgxpool"
)
// SourceStore persistiert, welche Module ihre Metriken unter welcher URL
// bereitstellen — dieselbe Postgres-basierte "kein Code-Deploy noetig"-
// Konvention wie internal/statuspage.Store.RegisterTarget (OPS-02): ein neu
// registriertes Modul erscheint automatisch in der Aggregation, sobald es
// hier eingetragen ist (Akzeptanzkriterium 2 / Pruefung 2).
type SourceStore struct {
pool *pgxpool.Pool
}
func NewSourceStore(pool *pgxpool.Pool) *SourceStore {
return &SourceStore{pool: pool}
}
func (s *SourceStore) RegisterSource(ctx context.Context, moduleName, metricsURL string) error {
_, err := s.pool.Exec(ctx, `
INSERT INTO metrics_sources (module_name, metrics_url)
VALUES ($1, $2)
ON CONFLICT (module_name) DO UPDATE SET metrics_url = $2
`, moduleName, metricsURL)
if err != nil {
return fmt.Errorf("metrik-quelle speichern: %w", err)
}
return nil
}
// Provide implementiert SourceProvider direkt aus der Datenbank.
func (s *SourceStore) Provide(ctx context.Context) ([]Source, error) {
rows, err := s.pool.Query(ctx, `SELECT module_name, metrics_url FROM metrics_sources ORDER BY module_name`)
if err != nil {
return nil, fmt.Errorf("metrik-quellen auflisten: %w", err)
}
defer rows.Close()
var out []Source
for rows.Next() {
var src Source
if err := rows.Scan(&src.ModuleName, &src.MetricsURL); err != nil {
return nil, fmt.Errorf("metrik-quelle lesen: %w", err)
}
out = append(out, src)
}
return out, rows.Err()
}
+60
View File
@@ -0,0 +1,60 @@
package metrics
import (
"context"
"fmt"
"os"
"testing"
"time"
"github.com/jackc/pgx/v5/pgxpool"
)
func setupSourceStoreTest(t *testing.T) (*SourceStore, func()) {
t.Helper()
adminDSN := os.Getenv("TEST_ADMIN_DSN")
if adminDSN == "" {
t.Skip("TEST_ADMIN_DSN nicht gesetzt, Integrationstest uebersprungen")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, adminDSN)
if err != nil {
t.Fatalf("pool: %v", err)
}
if _, err := pool.Exec(ctx, `
CREATE TABLE IF NOT EXISTS metrics_sources (module_name TEXT PRIMARY KEY, metrics_url TEXT NOT NULL)
`); err != nil {
t.Fatalf("schema: %v", err)
}
cleanup := func() { pool.Close() }
return NewSourceStore(pool), cleanup
}
// Akzeptanzkriterium 2 / Pruefung 2 auf Persistenz-Ebene: eine ueber die
// Datenbank registrierte Quelle ist sofort ueber Provide() sichtbar — genau
// der Mechanismus, der ein neues Modul ohne Core-Codeaenderung erscheinen
// laesst.
func TestSourceStore_RegisterSourceAppearsInProvide(t *testing.T) {
store, cleanup := setupSourceStoreTest(t)
defer cleanup()
ctx := context.Background()
name := fmt.Sprintf("modul-%d", time.Now().UnixNano())
if err := store.RegisterSource(ctx, name, "http://example.invalid/metrics"); err != nil {
t.Fatalf("registersource: %v", err)
}
sources, err := store.Provide(ctx)
if err != nil {
t.Fatalf("provide: %v", err)
}
found := false
for _, s := range sources {
if s.ModuleName == name {
found = true
}
}
if !found {
t.Fatalf("erwartet %s in provide()-ergebnis, habe: %+v", name, sources)
}
}
+1
View File
@@ -0,0 +1 @@
DROP TABLE metrics_sources;
+7
View File
@@ -0,0 +1,7 @@
-- Metrics-Aggregation ueber Module hinweg (OPS-03, siehe
-- core-kanban/tickets/OPS-03.md) — welches Modul liefert seine Kennzahlen
-- unter welcher /metrics-URL.
CREATE TABLE metrics_sources (
module_name TEXT PRIMARY KEY,
metrics_url TEXT NOT NULL
);