Reduce metric writes and preserve scheduler history
This commit is contained in:
@@ -9,6 +9,7 @@ import (
|
||||
"time"
|
||||
|
||||
"git.casaderoll.de/michael/urbm/internal/model"
|
||||
"git.casaderoll.de/michael/urbm/internal/queue"
|
||||
"git.casaderoll.de/michael/urbm/internal/store"
|
||||
)
|
||||
|
||||
@@ -39,7 +40,8 @@ func TestRepositoryLockHonorsContextCancellation(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestBackupMetricsSurviveClearedRunSnapshot(t *testing.T) {
|
||||
st := store.New(t.TempDir())
|
||||
dir := t.TempDir()
|
||||
st := store.New(dir)
|
||||
if err := st.Init(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -47,13 +49,36 @@ func TestBackupMetricsSurviveClearedRunSnapshot(t *testing.T) {
|
||||
finished := time.Date(2026, 7, 13, 1, 30, 0, 0, time.UTC)
|
||||
run := model.Run{ID: "run-1", JobID: "job-1", TaskType: "backup", Status: "success", FinishedAt: &finished, BytesProcessed: 1024, FilesProcessed: 12}
|
||||
s.queuePersistence([]model.Run{run})
|
||||
before, err := os.ReadFile(filepath.Join(dir, "backup-metrics.jsonl"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s.queuePersistence(nil)
|
||||
after, err := os.ReadFile(filepath.Join(dir, "backup-metrics.jsonl"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(after) != string(before) {
|
||||
t.Fatal("clearing activities rewrote the backup metric file")
|
||||
}
|
||||
metrics := s.BackupMetrics()
|
||||
if len(metrics) != 1 || metrics[0].RunID != run.ID || metrics[0].FilesProcessed != 12 {
|
||||
t.Fatalf("backup metrics after clear = %#v", metrics)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLastRunUsesDurableMetricAfterActivityClear(t *testing.T) {
|
||||
created := time.Date(2026, 7, 13, 1, 0, 0, 0, time.UTC)
|
||||
finished := created.Add(3 * time.Hour)
|
||||
q := queue.New(nil, nil)
|
||||
q.RestoreHistory([]model.Run{{ID: "run-1", JobID: "job-1", Status: "success", CreatedAt: created}})
|
||||
s := &Service{queue: q, backupMetrics: []model.BackupMetric{{RunID: "run-1", JobID: "job-1", Status: "success", FinishedAt: finished}}}
|
||||
q.ClearHistory()
|
||||
if got := s.LastRun("job-1"); !got.Equal(finished) {
|
||||
t.Fatalf("last run after activity clear = %v, want %v", got, finished)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBackupMetricsOnlyCaptureCompletedSuccessfulBackupsOnce(t *testing.T) {
|
||||
st := store.New(t.TempDir())
|
||||
s := &Service{store: st}
|
||||
|
||||
+73
-11
@@ -54,8 +54,15 @@ type Service struct {
|
||||
persistStopOnce sync.Once
|
||||
metricsMu sync.RWMutex
|
||||
backupMetrics []model.BackupMetric
|
||||
metricRecords int
|
||||
metricsAppendOK bool
|
||||
}
|
||||
|
||||
const (
|
||||
maxBackupMetrics = 5000
|
||||
compactBackupMetricsAt = 5500
|
||||
)
|
||||
|
||||
type SnapshotBrowserNode struct {
|
||||
Path string `json:"path"`
|
||||
Name string `json:"name"`
|
||||
@@ -101,10 +108,20 @@ func New(st *store.Store, secrets SecretStore, rr *restic.Runner, rs *rsync.Runn
|
||||
store: st, secrets: secrets, restic: rr, rsync: rs, mounts: mounts, workloads: workloads, notifier: notifier, log: log, config: config,
|
||||
repositoryLocks: map[string]chan struct{}{}, persistRuns: make(chan []model.Run, 1), persistStop: make(chan struct{}), persistDone: make(chan struct{}),
|
||||
}
|
||||
if metrics, err := st.LoadBackupMetrics(); err != nil {
|
||||
if metrics, legacy, err := st.LoadBackupMetrics(); err != nil {
|
||||
log.Error("load backup metrics", "error", err)
|
||||
} else {
|
||||
s.backupMetrics = metrics
|
||||
s.backupMetrics = normalizeBackupMetrics(metrics)
|
||||
s.metricRecords = len(s.backupMetrics)
|
||||
if legacy || len(metrics) != len(s.backupMetrics) {
|
||||
if err := st.SaveBackupMetrics(s.backupMetrics); err != nil {
|
||||
log.Error("migrate backup metrics to append format", "error", err)
|
||||
} else {
|
||||
s.metricsAppendOK = true
|
||||
}
|
||||
} else {
|
||||
s.metricsAppendOK = true
|
||||
}
|
||||
}
|
||||
go s.runPersistence()
|
||||
s.queue = queue.New(s.execute, s.queuePersistence)
|
||||
@@ -202,9 +219,6 @@ func (s *Service) saveRunState(runs []model.Run) {
|
||||
if err := s.persistRunLogs(runs); err != nil {
|
||||
s.log.Error("save persistent run logs", "error", err)
|
||||
}
|
||||
if err := s.store.SaveBackupMetrics(s.BackupMetrics()); err != nil {
|
||||
s.log.Error("save backup metrics", "error", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Service) captureBackupMetrics(runs []model.Run) {
|
||||
@@ -212,11 +226,11 @@ func (s *Service) captureBackupMetrics(runs []model.Run) {
|
||||
return
|
||||
}
|
||||
s.metricsMu.Lock()
|
||||
defer s.metricsMu.Unlock()
|
||||
known := make(map[string]struct{}, len(s.backupMetrics))
|
||||
for _, metric := range s.backupMetrics {
|
||||
known[metric.RunID] = struct{}{}
|
||||
}
|
||||
added := make([]model.BackupMetric, 0, 1)
|
||||
for _, run := range runs {
|
||||
if run.TaskType != "backup" || run.FinishedAt == nil || (run.Status != "success" && run.Status != "warning") {
|
||||
continue
|
||||
@@ -224,17 +238,58 @@ func (s *Service) captureBackupMetrics(runs []model.Run) {
|
||||
if _, exists := known[run.ID]; exists {
|
||||
continue
|
||||
}
|
||||
s.backupMetrics = append(s.backupMetrics, model.BackupMetric{
|
||||
metric := model.BackupMetric{
|
||||
SchemaVersion: model.SchemaVersion, RunID: run.ID, JobID: run.JobID, SnapshotID: run.SnapshotID,
|
||||
FinishedAt: *run.FinishedAt, Status: run.Status, BytesProcessed: run.BytesProcessed,
|
||||
FilesProcessed: run.FilesProcessed, BytesAdded: run.BytesAdded, FilesNew: run.FilesNew, FilesChanged: run.FilesChanged,
|
||||
})
|
||||
}
|
||||
s.backupMetrics = append(s.backupMetrics, metric)
|
||||
added = append(added, metric)
|
||||
known[run.ID] = struct{}{}
|
||||
}
|
||||
sort.SliceStable(s.backupMetrics, func(i, j int) bool { return s.backupMetrics[i].FinishedAt.Before(s.backupMetrics[j].FinishedAt) })
|
||||
if len(s.backupMetrics) > 5000 {
|
||||
s.backupMetrics = append([]model.BackupMetric(nil), s.backupMetrics[len(s.backupMetrics)-5000:]...)
|
||||
if len(added) == 0 {
|
||||
s.metricsMu.Unlock()
|
||||
return
|
||||
}
|
||||
s.backupMetrics = normalizeBackupMetrics(s.backupMetrics)
|
||||
if len(added) == 1 && s.metricsAppendOK && s.metricRecords < compactBackupMetricsAt {
|
||||
if err := s.store.AppendBackupMetric(added[0]); err == nil {
|
||||
s.metricRecords++
|
||||
s.metricsMu.Unlock()
|
||||
return
|
||||
} else if s.log != nil {
|
||||
s.log.Error("append backup metric", "error", err)
|
||||
}
|
||||
}
|
||||
if err := s.store.SaveBackupMetrics(s.backupMetrics); err != nil {
|
||||
if s.log != nil {
|
||||
s.log.Error("compact backup metrics", "error", err)
|
||||
}
|
||||
} else {
|
||||
s.metricRecords = len(s.backupMetrics)
|
||||
s.metricsAppendOK = true
|
||||
}
|
||||
s.metricsMu.Unlock()
|
||||
}
|
||||
|
||||
func normalizeBackupMetrics(metrics []model.BackupMetric) []model.BackupMetric {
|
||||
known := make(map[string]struct{}, len(metrics))
|
||||
normalized := make([]model.BackupMetric, 0, len(metrics))
|
||||
for _, metric := range metrics {
|
||||
if metric.RunID == "" {
|
||||
continue
|
||||
}
|
||||
if _, exists := known[metric.RunID]; exists {
|
||||
continue
|
||||
}
|
||||
known[metric.RunID] = struct{}{}
|
||||
normalized = append(normalized, metric)
|
||||
}
|
||||
sort.SliceStable(normalized, func(i, j int) bool { return normalized[i].FinishedAt.Before(normalized[j].FinishedAt) })
|
||||
if len(normalized) > maxBackupMetrics {
|
||||
normalized = append([]model.BackupMetric(nil), normalized[len(normalized)-maxBackupMetrics:]...)
|
||||
}
|
||||
return normalized
|
||||
}
|
||||
|
||||
func (s *Service) persistRunLogs(runs []model.Run) error {
|
||||
@@ -353,6 +408,13 @@ func (s *Service) LastRun(jobID string) time.Time {
|
||||
latest = run.CreatedAt
|
||||
}
|
||||
}
|
||||
s.metricsMu.RLock()
|
||||
defer s.metricsMu.RUnlock()
|
||||
for _, metric := range s.backupMetrics {
|
||||
if metric.JobID == jobID && metric.FinishedAt.After(latest) {
|
||||
latest = metric.FinishedAt
|
||||
}
|
||||
}
|
||||
return latest
|
||||
}
|
||||
|
||||
|
||||
+78
-8
@@ -1,6 +1,7 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -70,16 +71,48 @@ func (s *Store) SaveRuns(runs []model.Run) error {
|
||||
return writeJSONAtomic(s.runsPath(), runs, 0600)
|
||||
}
|
||||
|
||||
func (s *Store) LoadBackupMetrics() ([]model.BackupMetric, error) {
|
||||
func (s *Store) LoadBackupMetrics() ([]model.BackupMetric, bool, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
var metrics []model.BackupMetric
|
||||
if err := readJSON(s.backupMetricsPath(), &metrics); errors.Is(err, os.ErrNotExist) {
|
||||
return []model.BackupMetric{}, nil
|
||||
} else if err != nil {
|
||||
return nil, err
|
||||
path := s.backupMetricsPath()
|
||||
b, err := os.ReadFile(path)
|
||||
legacySource := false
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
path = s.legacyBackupMetricsPath()
|
||||
b, err = os.ReadFile(path)
|
||||
legacySource = true
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
return []model.BackupMetric{}, false, nil
|
||||
}
|
||||
}
|
||||
return metrics, nil
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
b = bytes.TrimSpace(b)
|
||||
if len(b) == 0 {
|
||||
return []model.BackupMetric{}, false, nil
|
||||
}
|
||||
if b[0] == '[' {
|
||||
var metrics []model.BackupMetric
|
||||
if err := json.Unmarshal(b, &metrics); err != nil {
|
||||
return nil, true, fmt.Errorf("decode %s: %w", path, err)
|
||||
}
|
||||
return metrics, true, nil
|
||||
}
|
||||
lines := bytes.Split(b, []byte{'\n'})
|
||||
metrics := make([]model.BackupMetric, 0, len(lines))
|
||||
for index, line := range lines {
|
||||
line = bytes.TrimSpace(line)
|
||||
if len(line) == 0 {
|
||||
continue
|
||||
}
|
||||
var metric model.BackupMetric
|
||||
if err := json.Unmarshal(line, &metric); err != nil {
|
||||
return nil, legacySource, fmt.Errorf("decode %s line %d: %w", path, index+1, err)
|
||||
}
|
||||
metrics = append(metrics, metric)
|
||||
}
|
||||
return metrics, legacySource, nil
|
||||
}
|
||||
|
||||
func (s *Store) SaveBackupMetrics(metrics []model.BackupMetric) error {
|
||||
@@ -88,12 +121,45 @@ func (s *Store) SaveBackupMetrics(metrics []model.BackupMetric) error {
|
||||
if len(metrics) > 5000 {
|
||||
metrics = metrics[len(metrics)-5000:]
|
||||
}
|
||||
return writeJSONAtomic(s.backupMetricsPath(), metrics, 0600)
|
||||
var data bytes.Buffer
|
||||
encoder := json.NewEncoder(&data)
|
||||
for _, metric := range metrics {
|
||||
if err := encoder.Encode(metric); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return writeBytesAtomic(s.backupMetricsPath(), data.Bytes(), 0600)
|
||||
}
|
||||
|
||||
func (s *Store) AppendBackupMetric(metric model.BackupMetric) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
b, err := json.Marshal(metric)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
b = append(b, '\n')
|
||||
file, err := os.OpenFile(s.backupMetricsPath(), os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0600)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := file.Write(b); err != nil {
|
||||
file.Close()
|
||||
return err
|
||||
}
|
||||
if err := file.Sync(); err != nil {
|
||||
file.Close()
|
||||
return err
|
||||
}
|
||||
return file.Close()
|
||||
}
|
||||
|
||||
func (s *Store) configPath() string { return filepath.Join(s.dir, "config.json") }
|
||||
func (s *Store) runsPath() string { return filepath.Join(s.dir, "runs.json") }
|
||||
func (s *Store) backupMetricsPath() string {
|
||||
return filepath.Join(s.dir, "backup-metrics.jsonl")
|
||||
}
|
||||
func (s *Store) legacyBackupMetricsPath() string {
|
||||
return filepath.Join(s.dir, "backup-metrics.json")
|
||||
}
|
||||
|
||||
@@ -114,6 +180,10 @@ func writeJSONAtomic(path string, value any, mode os.FileMode) error {
|
||||
return err
|
||||
}
|
||||
b = append(b, '\n')
|
||||
return writeBytesAtomic(path, b, mode)
|
||||
}
|
||||
|
||||
func writeBytesAtomic(path string, b []byte, mode os.FileMode) error {
|
||||
tmp, err := os.CreateTemp(filepath.Dir(path), ".urbm-*")
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -43,10 +43,13 @@ func TestBackupMetricsPersistSeparatelyFromRuns(t *testing.T) {
|
||||
if err := s.SaveRuns(nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
loaded, err := s.LoadBackupMetrics()
|
||||
loaded, legacy, err := s.LoadBackupMetrics()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if legacy {
|
||||
t.Fatal("compacted metric file was detected as legacy JSON array")
|
||||
}
|
||||
if len(loaded) != 1 || loaded[0].RunID != "run-1" || loaded[0].BytesProcessed != 42 {
|
||||
t.Fatalf("backup metrics = %#v", loaded)
|
||||
}
|
||||
@@ -58,3 +61,32 @@ func TestBackupMetricsPersistSeparatelyFromRuns(t *testing.T) {
|
||||
t.Fatalf("backup metrics mode = %o", info.Mode().Perm())
|
||||
}
|
||||
}
|
||||
|
||||
func TestBackupMetricsMigrateFromJSONArrayAndAppend(t *testing.T) {
|
||||
s := New(t.TempDir())
|
||||
finished := time.Date(2026, 7, 13, 1, 30, 0, 0, time.UTC)
|
||||
legacyJSON := `[{"schemaVersion":1,"runId":"old","jobId":"job-1","finishedAt":"2026-07-13T01:30:00Z","status":"success","bytesProcessed":42,"filesProcessed":7}]`
|
||||
if err := os.WriteFile(s.legacyBackupMetricsPath(), []byte(legacyJSON), 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
metrics, legacy, err := s.LoadBackupMetrics()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !legacy || len(metrics) != 1 || metrics[0].RunID != "old" {
|
||||
t.Fatalf("legacy metrics = %#v, legacy = %v", metrics, legacy)
|
||||
}
|
||||
if err := s.SaveBackupMetrics(metrics); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.AppendBackupMetric(model.BackupMetric{SchemaVersion: model.SchemaVersion, RunID: "new", JobID: "job-1", FinishedAt: finished.Add(time.Hour), Status: "success"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
metrics, legacy, err = s.LoadBackupMetrics()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if legacy || len(metrics) != 2 || metrics[1].RunID != "new" {
|
||||
t.Fatalf("appended metrics = %#v, legacy = %v", metrics, legacy)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user