// 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_ — Kennzahlen des Core-Dienstes selbst // nexarch_module__ — von einem Modul gescrapte Kennzahl // , 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 } } } }