syncova-backup/packages/metrics/collector.go
Jerrit Fritzsche 610719c316
Some checks failed
CI / Backend (Go) (push) Failing after 3m7s
CI / Frontend (React/TypeScript) (push) Successful in 37s
CI / Sicherheitsprüfungen (push) Successful in 44s
Syncova Backups V1
Enterprise-Backup-, Recovery-, Verification-, Security- und
Monitoring-Plattform fuer Proxmox VE, Windows, Linux und Dateisysteme.

Der Leitsatz, der fast jede Entscheidung erklaert: Ein Backup gilt erst als
vertrauenswuerdig, wenn Integritaet geprueft und Wiederherstellbarkeit
nachgewiesen wurde. Deshalb steigt ein Wiederherstellungspunkt erst nach einem
tatsaechlich durchgefuehrten Restore-Test auf "recoverable", und Unbekanntes
geht in keine Bewertung als "gut" ein.

Umfang (Phasen 0-23):

- Repository Engine: inhaltsadressierte Bloecke, atomares Commit-Protokoll,
  Katalogaufbau allein aus den Manifesten — ohne Datenbank
- Backup Engine: inhaltsabhaengiges Chunking, Deduplizierung trotz
  Verschluesselung, zstd, AES-256-GCM, Streaming mit Gegendruck
- Agenten fuer Windows und Linux mit Auftragsabholung (Pull-Modell)
- Proxmox-Provider mit beiden Zugriffswegen auf die Sicherungsarchive
- Scheduler, Recovery Engine mit Pruefpunkt, Verification, Unveraenderlichkeit
- Weboberflaeche, Kennzahlen, Meldungen, Berichte, Security Center,
  Ransomware-Heuristik (meldet, handelt nie)
- Disaster Recovery, Haertung, Leistungsmessung, Chaos Testing
- Eingefrorene Vertraege fuer API, Migrationen, Backup-Format und Repository
- Auslieferungspaket fuer linux/amd64, linux/arm64 und windows/amd64

Nicht enthalten und als solches gekennzeichnet: Kapazitaetsprognose, Backup
Copy, Changed Block Tracking bei Proxmox, erweiterte Attribute und ACLs.

Gebaut, aber nie auf echter Hardware gefahren: der Windows-Dienst, die
systemd-Einheit und der verpflichtende Proxmox-Meilenstein — ob eine
wiederhergestellte VM startet, ist ungeprueft. Einzelheiten in CHANGELOG.md
und docs/release-candidate.md.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-17 09:10:54 +02:00

420 lines
15 KiB
Go

package metrics
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"github.com/google/uuid"
"github.com/syncova/syncova/packages/platform/logging"
)
// Namen der gesammelten Momentaufnahmen.
const (
// sampleStorageUsed ist der belegte Speicher eines Repositorys.
sampleStorageUsed = "repository_used_bytes"
// sampleStorageCapacity ist die Gesamtkapazitaet eines Repositorys.
sampleStorageCapacity = "repository_capacity_bytes"
// sampleStorageUsedPercent ist die Auslastung eines Repositorys.
sampleStorageUsedPercent = "repository_used_percent"
// sampleSecurityScore ist die Sicherheitsbewertung der Anlage.
sampleSecurityScore = "security_score"
)
// systemSubjectID ist der Gegenstand anlagenweiter Messungen.
//
// Eine feste Kennung statt NULL: Die Reihe braucht einen Gegenstand, und die
// Anlage als Ganzes ist einer. Ein NULL brauchte eine Sonderbehandlung in jeder
// Abfrage und im Eindeutigkeitsindex.
var systemSubjectID = uuid.MustParse("00000000-0000-0000-0000-000000000001")
// RecordSecurityScore legt eine Sicherheitsbewertung als Messung ab.
//
// Der Score selbst wird bei jedem Aufruf neu berechnet (er haengt am Zustand).
// Was hier entsteht, ist sein **Verlauf** — und der geht sonst verloren, weil
// niemand die Lage von gestern noch kennt.
func (store *Store) RecordSecurityScore(recordContext context.Context, scorePercentage float64, measuredAt time.Time) error {
roundedTime := measuredAt.UTC().Truncate(time.Minute)
const insertStatement = `
INSERT INTO metric_samples (metric_name, subject_type, subject_id, value, unit, measured_at)
VALUES ($1, 'system', $2, $3, 'percent', $4)
ON CONFLICT (metric_name, subject_id, measured_at) DO NOTHING`
if _, execError := store.connectionPool.Exec(recordContext, insertStatement,
sampleSecurityScore, systemSubjectID, scorePercentage, roundedTime); execError != nil {
return fmt.Errorf("die sicherheitsbewertung konnte nicht abgelegt werden: %w", execError)
}
return nil
}
// buildSampleSeries bildet eine Reihe aus gesammelten Momentaufnahmen.
//
// Anders als bei den Ereignistabellen wird hier **der letzte Wert je
// Zeitfenster** genommen, nicht die Summe: Ein Fuellstand ist ein Zustand, kein
// Vorgang. Zwei Messungen in derselben Stunde zu addieren ergaebe die doppelte
// Belegung.
func (store *Store) buildSampleSeries(buildContext context.Context, window Window, metricName string, seriesLabel string, seriesUnit SeriesUnit, repositoryFilter *uuid.UUID) ([]Series, error) {
// DISTINCT ON liefert je Zeitfenster und Repository die juengste Messung;
// die aeussere Abfrage summiert dann ueber die Repositories. Ohne den
// inneren Schritt zaehlte ein haeufiger gemessenes Repository mehrfach.
const sampleStatement = `
WITH letzte_je_fenster AS (
SELECT DISTINCT ON (date_bin($3::interval, measured_at, $1::timestamptz), subject_id)
date_bin($3::interval, measured_at, $1::timestamptz) AS bucket,
subject_id,
value
FROM metric_samples
WHERE metric_name = $4
AND measured_at >= $1 AND measured_at < $2
AND ($5::uuid IS NULL OR subject_id = $5::uuid)
ORDER BY date_bin($3::interval, measured_at, $1::timestamptz), subject_id, measured_at DESC
)
SELECT bucket, sum(value)::float8, count(*)::int
FROM letzte_je_fenster
GROUP BY bucket
ORDER BY bucket`
collectedRows, queryError := store.queryBuckets(buildContext, sampleStatement,
window.From, window.To, window.BucketWidth.String(), metricName, repositoryFilter)
if queryError != nil {
return nil, queryError
}
// Eine Auslastung ist ein Anteil und darf nicht ueber Repositories summiert
// werden — zwei zu 80 % gefuellte Repositories ergaeben sonst 160 %.
if seriesUnit == UnitPercent {
for rowIndex := range collectedRows {
if collectedRows[rowIndex].SampleCount > 0 {
collectedRows[rowIndex].Value /= float64(collectedRows[rowIndex].SampleCount)
}
}
}
return []Series{fillSeries(metricName, seriesLabel, seriesUnit, window, collectedRows)}, nil
}
// RepositoryMeasurement ist eine Messung an einem Repository.
type RepositoryMeasurement struct {
// RepositoryID ist das gemessene Repository.
RepositoryID uuid.UUID
// UsedBytes ist der belegte Speicher.
UsedBytes int64
// CapacityBytes ist die Gesamtkapazitaet; 0 bedeutet unbekannt.
CapacityBytes int64
}
// RecordRepositoryMeasurement legt eine Momentaufnahme ab.
//
// Der Zeitpunkt wird auf die volle Minute gerundet. Ohne diese Rundung
// entstuenden bei zwei gleichzeitig sammelnden Control-Servern zwei Messungen
// mit Millisekundenabstand — der Eindeutigkeitsindex griffe nicht, und die
// Belegung erschiene doppelt.
func (store *Store) RecordRepositoryMeasurement(recordContext context.Context, measurement RepositoryMeasurement, measuredAt time.Time) error {
roundedTime := measuredAt.UTC().Truncate(time.Minute)
// ON CONFLICT DO NOTHING statt eines Fehlers: Zwei Control-Server, die
// dieselbe Minute erfassen, sind kein Betriebsfehler. Der zweite hat nichts
// beizutragen.
const insertStatement = `
INSERT INTO metric_samples (metric_name, subject_type, subject_id, value, unit, measured_at)
VALUES ($1, 'repository', $2, $3, $4, $5)
ON CONFLICT (metric_name, subject_id, measured_at) DO NOTHING`
samplesToWrite := []struct {
name string
value float64
unit string
}{
{sampleStorageUsed, float64(measurement.UsedBytes), "bytes"},
}
if measurement.CapacityBytes > 0 {
// Ohne hinterlegte Kapazitaet gibt es keine Auslastung. Sie zu schaetzen
// waere eine erfundene Statistik; die Reihe bleibt an dieser Stelle leer.
usedPercentage := float64(measurement.UsedBytes) * 100 / float64(measurement.CapacityBytes)
samplesToWrite = append(samplesToWrite,
struct {
name string
value float64
unit string
}{sampleStorageCapacity, float64(measurement.CapacityBytes), "bytes"},
struct {
name string
value float64
unit string
}{sampleStorageUsedPercent, usedPercentage, "percent"},
)
}
for _, sampleToWrite := range samplesToWrite {
if _, execError := store.connectionPool.Exec(recordContext, insertStatement,
sampleToWrite.name, measurement.RepositoryID, sampleToWrite.value,
sampleToWrite.unit, roundedTime); execError != nil {
return fmt.Errorf("die messung %s konnte nicht abgelegt werden: %w", sampleToWrite.name, execError)
}
}
const markStatement = `UPDATE repositories SET metrics_collected_at = $2 WHERE id = $1`
if _, execError := store.connectionPool.Exec(recordContext, markStatement,
measurement.RepositoryID, roundedTime); execError != nil {
return fmt.Errorf("der erfassungszeitpunkt konnte nicht vermerkt werden: %w", execError)
}
return nil
}
// PruneSamples entfernt Messungen aelter als die Aufbewahrungsfrist.
//
// Ohne Aufraeumen waechst die Tabelle unbegrenzt: Bei fuenf Repositories, drei
// Reihen und einer Messung je fuenf Minuten sind das gut anderthalb Millionen
// Zeilen im Jahr. Die Frist deckt den laengsten Auswertungszeitraum ab.
func (store *Store) PruneSamples(pruneContext context.Context, retentionPeriod time.Duration) (int64, error) {
const deleteStatement = `DELETE FROM metric_samples WHERE measured_at < now() - $1::interval`
commandTag, execError := store.connectionPool.Exec(pruneContext, deleteStatement,
retentionPeriod.String())
if execError != nil {
return 0, fmt.Errorf("alte messungen konnten nicht entfernt werden: %w", execError)
}
return commandTag.RowsAffected(), nil
}
// RepositorySource liefert die zu erfassenden Repositories.
//
// Die Schnittstelle haelt das Metrikpaket frei von der Auftragsverwaltung —
// dieselbe Naht wie bei Wiederherstellung und Pruefung.
type RepositorySource interface {
// MeasurableRepositories liefert Kennung und Pfad aller aktiven Repositories.
MeasurableRepositories(sourceContext context.Context) ([]MeasurableRepository, error)
}
// MeasurableRepository ist ein zu erfassendes Repository.
type MeasurableRepository struct {
// ID ist die Kennung in der Control Plane.
ID uuid.UUID
// Name ist die sprechende Bezeichnung.
Name string
// Location ist der Pfad der Ablage.
Location string
}
// CapacityProbe ermittelt Belegung und Kapazitaet eines Pfades.
type CapacityProbe func(repositoryPath string) (usedBytes int64, capacityBytes int64, probeError error)
// SecurityScoreSource liefert die aktuelle Sicherheitsbewertung.
//
// Eine Schnittstelle statt eines direkten Aufrufs: Sie haelt das Metrikpaket
// frei vom Security Center — dieselbe Naht wie bei der Repository-Quelle.
type SecurityScoreSource interface {
// CurrentSecurityScore liefert die Bewertung in Prozent.
CurrentSecurityScore(scoreContext context.Context) (float64, error)
}
// CollectorOptions steuern die Sammelschleife.
type CollectorOptions struct {
// Interval ist der Abstand zweier Erfassungen.
//
// Fuenf Minuten sind der Ausgleich zwischen Aufloesung und Last: Ein
// Fuellstand aendert sich nicht sekuendlich, und jede Erfassung liest das
// Dateisystem.
Interval time.Duration
// RetentionPeriod ist die Aufbewahrungsfrist der Messungen.
RetentionPeriod time.Duration
}
// Standardwerte der Sammelschleife.
const (
// defaultCollectionInterval ist der Standardabstand der Erfassung.
defaultCollectionInterval = 5 * time.Minute
// defaultSampleRetention ist die Standardaufbewahrung der Messungen.
//
// Vierhundert Tage decken den laengsten Auswertungszeitraum von einem Jahr
// ab und lassen Raum fuer einen Jahresvergleich.
defaultSampleRetention = 400 * 24 * time.Hour
)
// applyDefaults fuellt fehlende Werte.
func (options *CollectorOptions) applyDefaults() {
if options.Interval <= 0 {
options.Interval = defaultCollectionInterval
}
if options.RetentionPeriod <= 0 {
options.RetentionPeriod = defaultSampleRetention
}
}
// Collector erfasst regelmaessig die Kennzahlen der Repositories.
type Collector struct {
// store legt die Messungen ab.
store *Store
// repositorySource liefert die zu erfassenden Repositories.
repositorySource RepositorySource
// capacityProbe ermittelt Belegung und Kapazitaet.
capacityProbe CapacityProbe
// securityScoreSource liefert die Sicherheitsbewertung; darf nil sein.
securityScoreSource SecurityScoreSource
// options sind die Einstellungen.
options CollectorOptions
// logger protokolliert den Verlauf.
logger *slog.Logger
}
// NewCollector erzeugt die Sammelschleife.
func NewCollector(store *Store, repositorySource RepositorySource, capacityProbe CapacityProbe, options CollectorOptions, baseLogger *slog.Logger) (*Collector, error) {
if store == nil {
return nil, errors.New("die erfassung braucht eine datenzugriffsschicht")
}
if repositorySource == nil {
return nil, errors.New("die erfassung braucht eine quelle fuer repositories")
}
if capacityProbe == nil {
return nil, errors.New("die erfassung braucht eine messung fuer belegung und kapazitaet")
}
options.applyDefaults()
return &Collector{
store: store,
repositorySource: repositorySource,
capacityProbe: capacityProbe,
options: options,
logger: logging.WithComponent(baseLogger, "metrics-collector"),
}, nil
}
// WithSecurityScoreSource ergaenzt die Erfassung der Sicherheitsbewertung.
//
// Getrennt vom Konstruktor, weil die Erfassung ohne sie vollstaendig arbeitet:
// Ein Dienst ohne Security Center soll trotzdem Kennzahlen sammeln.
func (collector *Collector) WithSecurityScoreSource(scoreSource SecurityScoreSource) *Collector {
collector.securityScoreSource = scoreSource
return collector
}
// Run erfasst bis zum Abbruch des Kontexts.
func (collector *Collector) Run(runContext context.Context) error {
collector.logger.Info("die kennzahlenerfassung beginnt",
slog.Duration("abstand", collector.options.Interval))
// Der erste Durchgang laeuft sofort: Sonst begaenne jede Verlaufsreihe erst
// fuenf Minuten nach dem Start, und ein kurz laufender Dienst erfasste nie
// etwas.
collector.collectOnce(runContext)
collectionTicker := time.NewTicker(collector.options.Interval)
defer collectionTicker.Stop()
pruneTicker := time.NewTicker(24 * time.Hour)
defer pruneTicker.Stop()
for {
select {
case <-runContext.Done():
collector.logger.Info("die kennzahlenerfassung wird beendet")
return nil
case <-collectionTicker.C:
collector.collectOnce(runContext)
case <-pruneTicker.C:
collector.pruneOldSamples(runContext)
}
}
}
// collectOnce erfasst alle Repositories einmal.
func (collector *Collector) collectOnce(collectContext context.Context) {
measurableRepositories, sourceError := collector.repositorySource.MeasurableRepositories(collectContext)
if sourceError != nil {
collector.logger.Error("die repositories konnten nicht ermittelt werden",
slog.String("grund", sourceError.Error()))
return
}
measurementTime := time.Now()
for _, measurableRepository := range measurableRepositories {
if contextError := collectContext.Err(); contextError != nil {
return
}
usedBytes, capacityBytes, probeError := collector.capacityProbe(measurableRepository.Location)
if probeError != nil {
// Ein nicht erreichbares Repository ist eine Auskunft, kein Grund,
// die Erfassung der uebrigen abzubrechen. Eine Null einzutragen waere
// schlimmer: Die Kurve fiele auf den Boden, obwohl die Daten da sind.
collector.logger.Warn("ein repository liess sich nicht vermessen",
slog.String("repository", measurableRepository.Name),
slog.String("grund", probeError.Error()))
continue
}
if recordError := collector.store.RecordRepositoryMeasurement(collectContext,
RepositoryMeasurement{
RepositoryID: measurableRepository.ID,
UsedBytes: usedBytes,
CapacityBytes: capacityBytes,
}, measurementTime); recordError != nil {
collector.logger.Error("eine messung konnte nicht abgelegt werden",
slog.String("repository", measurableRepository.Name),
slog.String("grund", recordError.Error()))
}
}
collector.collectSecurityScore(collectContext, measurementTime)
}
// collectSecurityScore haelt die Sicherheitsbewertung als Messung fest.
//
// Der Score wird bei jedem Abruf neu berechnet — er haengt am Zustand der
// Anlage. Was hier entsteht, ist sein Verlauf: Ohne ihn liesse sich nicht
// sagen, ob die Lage besser oder schlechter geworden ist.
func (collector *Collector) collectSecurityScore(collectContext context.Context, measurementTime time.Time) {
if collector.securityScoreSource == nil {
return
}
scorePercentage, scoreError := collector.securityScoreSource.CurrentSecurityScore(collectContext)
if scoreError != nil {
collector.logger.Warn("die sicherheitsbewertung liess sich nicht erheben",
slog.String("grund", scoreError.Error()))
return
}
if recordError := collector.store.RecordSecurityScore(collectContext,
scorePercentage, measurementTime); recordError != nil {
collector.logger.Error("die sicherheitsbewertung konnte nicht abgelegt werden",
slog.String("grund", recordError.Error()))
}
}
// pruneOldSamples entfernt Messungen ausserhalb der Aufbewahrungsfrist.
func (collector *Collector) pruneOldSamples(pruneContext context.Context) {
removedCount, pruneError := collector.store.PruneSamples(pruneContext, collector.options.RetentionPeriod)
if pruneError != nil {
collector.logger.Error("alte messungen konnten nicht entfernt werden",
slog.String("grund", pruneError.Error()))
return
}
if removedCount > 0 {
collector.logger.Info("alte messungen entfernt", slog.Int64("anzahl", removedCount))
}
}