274 lines
8.3 KiB
Go
274 lines
8.3 KiB
Go
// Package statuspage implementiert Core OPS-02: eine zentrale Statusseite,
|
|
// die den Health-Zustand aller registrierten Module aggregiert und den
|
|
// Verlauf vergangener Statusaenderungen speichert. Baut auf OPS-01
|
|
// (internal/health) auf, indem es GENAU die dort etablierten
|
|
// Readiness-Endpunkte je Modul abfragt — dieses Paket dupliziert keine
|
|
// Health-Check-Logik, es aggregiert nur deren Ergebnisse ueber die Zeit.
|
|
package statuspage
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
// Status ist der aggregierte Zustand EINES Moduls zu einem Zeitpunkt.
|
|
type Status string
|
|
|
|
const (
|
|
StatusUp Status = "up"
|
|
StatusDown Status = "down"
|
|
)
|
|
|
|
// Target ist ein zu ueberwachendes Modul mit seiner Readiness-URL
|
|
// (OPS-01-Endpunkt, z.B. "http://dms:8080/readyz"). Eigenstaendige
|
|
// Konfiguration statt Erweiterung von internal/moduleregistry.Module, um
|
|
// API-02 nicht anzufassen (Kein Umbau angrenzender Bereiche).
|
|
type Target struct {
|
|
Name string
|
|
HealthURL string
|
|
}
|
|
|
|
// Store persistiert Ueberwachungsziele und den Verlauf ihrer
|
|
// Statusaenderungen.
|
|
type Store struct {
|
|
pool *pgxpool.Pool
|
|
}
|
|
|
|
func NewStore(pool *pgxpool.Pool) *Store {
|
|
return &Store{pool: pool}
|
|
}
|
|
|
|
// RegisterTarget traegt ein zu ueberwachendes Modul ein oder aktualisiert
|
|
// dessen URL (Akzeptanzkriterium 1: "aller registrierten Module").
|
|
func (s *Store) RegisterTarget(ctx context.Context, t Target) error {
|
|
_, err := s.pool.Exec(ctx, `
|
|
INSERT INTO status_targets (name, health_url)
|
|
VALUES ($1, $2)
|
|
ON CONFLICT (name) DO UPDATE SET health_url = $2
|
|
`, t.Name, t.HealthURL)
|
|
if err != nil {
|
|
return fmt.Errorf("ueberwachungsziel speichern: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) ListTargets(ctx context.Context) ([]Target, error) {
|
|
rows, err := s.pool.Query(ctx, `SELECT name, health_url FROM status_targets ORDER BY name`)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("ueberwachungsziele auflisten: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []Target
|
|
for rows.Next() {
|
|
var t Target
|
|
if err := rows.Scan(&t.Name, &t.HealthURL); err != nil {
|
|
return nil, fmt.Errorf("ueberwachungsziel lesen: %w", err)
|
|
}
|
|
out = append(out, t)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// recordIfChanged schreibt NUR dann einen neuen Verlaufseintrag, wenn sich
|
|
// der Status seit dem letzten Eintrag geaendert hat (oder es der erste
|
|
// Eintrag ist) — der Verlauf zeigt Statusaenderungen (Akzeptanzkriterium 3),
|
|
// nicht jede einzelne Abfrage.
|
|
func (s *Store) recordIfChanged(ctx context.Context, name string, status Status) error {
|
|
var lastStatus string
|
|
err := s.pool.QueryRow(ctx, `
|
|
SELECT status FROM status_history WHERE name = $1 ORDER BY changed_at DESC LIMIT 1
|
|
`, name).Scan(&lastStatus)
|
|
if err != nil && err != pgx.ErrNoRows {
|
|
return fmt.Errorf("letzten status lesen: %w", err)
|
|
}
|
|
if err == nil && lastStatus == string(status) {
|
|
return nil // keine Aenderung, kein neuer Eintrag
|
|
}
|
|
|
|
if _, err := s.pool.Exec(ctx, `
|
|
INSERT INTO status_history (name, status, changed_at) VALUES ($1, $2, now())
|
|
`, name, string(status)); err != nil {
|
|
return fmt.Errorf("statuseintrag schreiben: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ModuleStatus ist der aktuelle Zustand EINES Moduls fuer die Uebersicht.
|
|
type ModuleStatus struct {
|
|
Name string `json:"name"`
|
|
Status Status `json:"status"`
|
|
LastChecked time.Time `json:"last_checked"`
|
|
}
|
|
|
|
// Overview liefert den aktuellen (letzten bekannten) Status jedes
|
|
// registrierten Ziels (Akzeptanzkriterium 1). Ziele ohne jemals erfolgte
|
|
// Pruefung erscheinen mit Status "down" — ein Modul, ueber das nichts
|
|
// bekannt ist, gilt als nicht verfuegbar (Fail-Safe-Default), nicht als
|
|
// stillschweigend "ok".
|
|
func (s *Store) Overview(ctx context.Context) ([]ModuleStatus, error) {
|
|
targets, err := s.ListTargets(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
out := make([]ModuleStatus, 0, len(targets))
|
|
for _, t := range targets {
|
|
var status string
|
|
var changedAt time.Time
|
|
err := s.pool.QueryRow(ctx, `
|
|
SELECT status, changed_at FROM status_history WHERE name = $1 ORDER BY changed_at DESC LIMIT 1
|
|
`, t.Name).Scan(&status, &changedAt)
|
|
if err == pgx.ErrNoRows {
|
|
out = append(out, ModuleStatus{Name: t.Name, Status: StatusDown})
|
|
continue
|
|
}
|
|
if err != nil {
|
|
return nil, fmt.Errorf("aktuellen status lesen (%s): %w", t.Name, err)
|
|
}
|
|
out = append(out, ModuleStatus{Name: t.Name, Status: Status(status), LastChecked: changedAt})
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// HistoryEntry ist EIN Verlaufseintrag (Akzeptanzkriterium 3).
|
|
type HistoryEntry struct {
|
|
Status Status `json:"status"`
|
|
ChangedAt time.Time `json:"changed_at"`
|
|
}
|
|
|
|
func (s *Store) History(ctx context.Context, name string) ([]HistoryEntry, error) {
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT status, changed_at FROM status_history WHERE name = $1 ORDER BY changed_at DESC
|
|
`, name)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("verlauf abfragen: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []HistoryEntry
|
|
for rows.Next() {
|
|
var e HistoryEntry
|
|
var status string
|
|
if err := rows.Scan(&status, &e.ChangedAt); err != nil {
|
|
return nil, fmt.Errorf("verlaufseintrag lesen: %w", err)
|
|
}
|
|
e.Status = Status(status)
|
|
out = append(out, e)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
// HTTPChecker fragt die Readiness-URL eines Moduls ab (OPS-01-Endpunkt) und
|
|
// liefert StatusUp NUR bei HTTP 200 — jeder andere Statuscode ODER ein
|
|
// Netzwerkfehler/Timeout gilt als StatusDown. Ein einzelnes nicht
|
|
// erreichbares Modul liefert einen FEHLERFREIEN StatusDown-Wert statt eines
|
|
// Go-Errors, damit Poller.Run ein fehlerhaftes Modul niemals mit einem
|
|
// anderen verwechseln oder den gesamten Zyklus abbrechen kann
|
|
// (Akzeptanzkriterium 2).
|
|
type HTTPChecker struct {
|
|
Client *http.Client
|
|
Timeout time.Duration
|
|
}
|
|
|
|
func NewHTTPChecker(timeout time.Duration) *HTTPChecker {
|
|
return &HTTPChecker{Client: &http.Client{}, Timeout: timeout}
|
|
}
|
|
|
|
func (c *HTTPChecker) Check(ctx context.Context, url string) Status {
|
|
ctx, cancel := context.WithTimeout(ctx, c.Timeout)
|
|
defer cancel()
|
|
|
|
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
|
if err != nil {
|
|
return StatusDown
|
|
}
|
|
resp, err := c.Client.Do(req)
|
|
if err != nil {
|
|
return StatusDown
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode == http.StatusOK {
|
|
return StatusUp
|
|
}
|
|
return StatusDown
|
|
}
|
|
|
|
// Poller fragt periodisch alle Ziele ab und schreibt Statusaenderungen fort.
|
|
type Poller struct {
|
|
store *Store
|
|
checker *HTTPChecker
|
|
}
|
|
|
|
func NewPoller(store *Store, checker *HTTPChecker) *Poller {
|
|
return &Poller{store: store, checker: checker}
|
|
}
|
|
|
|
// PollOnce prueft ALLE Ziele in einem Durchlauf. Ein fehlschlagendes Ziel
|
|
// (Netzwerkfehler, Timeout, Nicht-200) wird als StatusDown vermerkt und
|
|
// haelt die Pruefung der UEBRIGEN Ziele nicht auf — die Schleife laeuft
|
|
// sequenziell weiter, kein Ziel kann ein anderes blockieren
|
|
// (Akzeptanzkriterium 2 / Pruefung 2).
|
|
func (p *Poller) PollOnce(ctx context.Context) error {
|
|
targets, err := p.store.ListTargets(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, t := range targets {
|
|
status := p.checker.Check(ctx, t.HealthURL)
|
|
if err := p.store.recordIfChanged(ctx, t.Name, status); err != nil {
|
|
// Ein Schreibfehler fuer EIN Ziel darf die Pruefung der anderen
|
|
// nicht verhindern — dieselbe Fail-Isolation wie bei einem
|
|
// unerreichbaren Modul.
|
|
continue
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Run ruft PollOnce in festen Abstaenden auf, bis ctx beendet wird —
|
|
// dieselbe Konvention wie internal/tenant.Lifecycle.RunSweeper.
|
|
func (p *Poller) Run(ctx context.Context, interval time.Duration) {
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
_ = p.PollOnce(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// --- HTTP-Bindung fuer die Oberflaeche ---
|
|
|
|
func (s *Store) OverviewHandler(w http.ResponseWriter, r *http.Request) {
|
|
overview, err := s.Overview(r.Context())
|
|
writeJSONResult(w, overview, err)
|
|
}
|
|
|
|
func (s *Store) HistoryHandler(w http.ResponseWriter, r *http.Request) {
|
|
name := r.URL.Query().Get("name")
|
|
history, err := s.History(r.Context(), name)
|
|
writeJSONResult(w, history, err)
|
|
}
|
|
|
|
func writeJSONResult(w http.ResponseWriter, body any, err error) {
|
|
w.Header().Set("Content-Type", "application/json")
|
|
if err != nil {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
_ = json.NewEncoder(w).Encode(map[string]string{"error": err.Error()})
|
|
return
|
|
}
|
|
w.WriteHeader(http.StatusOK)
|
|
_ = json.NewEncoder(w).Encode(body)
|
|
}
|