Fix backup check notifications and workload recovery; release 2026.09.21.r001
This commit is contained in:
@@ -0,0 +1,120 @@
|
||||
package platform
|
||||
|
||||
import (
|
||||
"context"
|
||||
"git.casaderoll.de/michael/urbm/internal/model"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func fakeDocker(t *testing.T, body string) string {
|
||||
t.Helper()
|
||||
dir := t.TempDir()
|
||||
path := filepath.Join(dir, "docker")
|
||||
if err := os.WriteFile(path, []byte("#!/bin/sh\n"+body), 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Setenv("PATH", dir+":"+os.Getenv("PATH"))
|
||||
return dir
|
||||
}
|
||||
|
||||
func TestPreparePreservesStoppedContainersAndIncludesVolumes(t *testing.T) {
|
||||
dir := fakeDocker(t, `case "$1" in
|
||||
inspect)
|
||||
if [ "$2" = --format ]; then
|
||||
if [ "$4" = off ]; then echo '{"Running":false}'; else echo '{"Running":true}'; fi
|
||||
else echo '[{"Mounts":[{"Type":"bind","Source":"/appdata"},{"Type":"volume","Source":"/docker/volumes/db/_data"},{"Type":"tmpfs","Source":""}]}]'; fi;;
|
||||
stop|start) echo "$1 $2" >> "$ACTIONS";;
|
||||
esac
|
||||
`)
|
||||
actions := filepath.Join(dir, "actions")
|
||||
t.Setenv("ACTIONS", actions)
|
||||
w := &WorkloadManager{RuntimeDir: dir}
|
||||
job := model.Job{ID: "test", Type: model.JobDocker, Sources: []model.Source{{WorkloadID: "off"}, {WorkloadID: "on"}}}
|
||||
job.Consistency.Mode = "stop"
|
||||
p, err := w.Prepare(context.Background(), job)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(p.Stopped) != 1 || p.Stopped[0].WorkloadID != "on" {
|
||||
t.Fatalf("stopped: %+v", p.Stopped)
|
||||
}
|
||||
if !strings.Contains(strings.Join(p.Sources, "\n"), "/docker/volumes/db/_data") {
|
||||
t.Fatal("named volume omitted")
|
||||
}
|
||||
if err = w.Cleanup(context.Background(), p); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
data, _ := os.ReadFile(actions)
|
||||
if strings.Contains(string(data), "start off") || !strings.Contains(string(data), "start on") {
|
||||
t.Fatalf("actions: %s", data)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecoveryContinuesAfterFailureAndCancellation(t *testing.T) {
|
||||
dir := fakeDocker(t, `echo "$1 $2" >> "$ACTIONS"
|
||||
if [ "$2" = broken ]; then exit 1; fi
|
||||
`)
|
||||
actions := filepath.Join(dir, "actions")
|
||||
t.Setenv("ACTIONS", actions)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
w := &WorkloadManager{}
|
||||
err := w.Cleanup(ctx, Prepared{Kind: model.JobDocker, Stopped: []model.Source{{WorkloadID: "healthy"}, {WorkloadID: "broken"}}})
|
||||
if err == nil || !strings.Contains(err.Error(), "broken") {
|
||||
t.Fatalf("error: %v", err)
|
||||
}
|
||||
data, _ := os.ReadFile(actions)
|
||||
if !strings.Contains(string(data), "start healthy") {
|
||||
t.Fatalf("recovery aborted: %s", data)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPrepareRollsBackPreviouslyStoppedContainers(t *testing.T) {
|
||||
dir := fakeDocker(t, `case "$1" in
|
||||
inspect)
|
||||
if [ "$2" = bad ]; then exit 1; fi
|
||||
if [ "$2" = --format ]; then echo '{"Running":true}'; else echo '[{"Mounts":[]}]'; fi;;
|
||||
stop|start) echo "$1 $2" >> "$ACTIONS";;
|
||||
esac
|
||||
`)
|
||||
actions := filepath.Join(dir, "actions")
|
||||
t.Setenv("ACTIONS", actions)
|
||||
w := &WorkloadManager{RuntimeDir: dir}
|
||||
job := model.Job{ID: "test", Type: model.JobDocker, Sources: []model.Source{{WorkloadID: "first"}, {WorkloadID: "bad"}}}
|
||||
job.Consistency.Mode = "stop"
|
||||
if _, err := w.Prepare(context.Background(), job); err == nil {
|
||||
t.Fatal("expected failure")
|
||||
}
|
||||
data, _ := os.ReadFile(actions)
|
||||
if !strings.Contains(string(data), "start first") {
|
||||
t.Fatalf("missing rollback: %s", data)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSlowRecoveryDoesNotExhaustOtherStarts(t *testing.T) {
|
||||
dir := fakeDocker(t, `if [ "$2" = slow ]; then sleep 6; fi
|
||||
echo "$2" >> "$ACTIONS"
|
||||
`)
|
||||
actions := filepath.Join(dir, "actions")
|
||||
t.Setenv("ACTIONS", actions)
|
||||
err := (&WorkloadManager{}).Cleanup(context.Background(), Prepared{Kind: model.JobDocker, Stopped: []model.Source{{WorkloadID: "next"}, {WorkloadID: "slow"}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
data, _ := os.ReadFile(actions)
|
||||
if string(data) != "slow\nnext\n" {
|
||||
t.Fatalf("incomplete recovery: %s", data)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPersistentMountWithoutPathFails(t *testing.T) {
|
||||
dir := fakeDocker(t, `echo '[{"Mounts":[{"Type":"volume","Source":""}]}]'
|
||||
`)
|
||||
w := &WorkloadManager{}
|
||||
if _, err := w.captureMetadata(context.Background(), model.JobDocker, "db", dir); err == nil {
|
||||
t.Fatal("missing volume path silently omitted")
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"encoding/xml"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
@@ -138,16 +139,21 @@ func (w *WorkloadManager) Prepare(ctx context.Context, job model.Job) (Prepared,
|
||||
}
|
||||
discovered, err := w.captureMetadata(ctx, job.Type, source.WorkloadID, metadataDir)
|
||||
if err != nil {
|
||||
w.Cleanup(context.Background(), prepared)
|
||||
return Prepared{}, err
|
||||
return Prepared{}, errors.Join(err, w.Cleanup(context.Background(), prepared))
|
||||
}
|
||||
prepared.Sources = append(prepared.Sources, discovered...)
|
||||
if job.Consistency.Mode == "stop" {
|
||||
if err := w.stop(ctx, job.Type, source.WorkloadID, time.Duration(job.ShutdownSecs)*time.Second); err != nil {
|
||||
w.Cleanup(context.Background(), prepared)
|
||||
return Prepared{}, err
|
||||
running, err := w.running(ctx, job.Type, source.WorkloadID)
|
||||
if err != nil {
|
||||
return Prepared{}, errors.Join(err, w.Cleanup(context.Background(), prepared))
|
||||
}
|
||||
if running {
|
||||
// Record before stopping: cancellation can happen after the daemon stopped it.
|
||||
prepared.Stopped = append(prepared.Stopped, source)
|
||||
if err := w.stop(ctx, job.Type, source.WorkloadID, time.Duration(job.ShutdownSecs)*time.Second); err != nil {
|
||||
return Prepared{}, errors.Join(err, w.Cleanup(context.Background(), prepared))
|
||||
}
|
||||
}
|
||||
prepared.Stopped = append(prepared.Stopped, source)
|
||||
}
|
||||
}
|
||||
prepared.Sources = uniqueStrings(prepared.Sources)
|
||||
@@ -203,22 +209,56 @@ func (w *WorkloadManager) FlashDevice(ctx context.Context) (string, error) {
|
||||
}
|
||||
|
||||
func (w *WorkloadManager) Cleanup(ctx context.Context, prepared Prepared) error {
|
||||
if _, hasDeadline := ctx.Deadline(); !hasDeadline {
|
||||
var cancel context.CancelFunc
|
||||
ctx, cancel = context.WithTimeout(ctx, 5*time.Second)
|
||||
defer cancel()
|
||||
}
|
||||
var first error
|
||||
var failures []error
|
||||
// Recovery must continue even when the backup was cancelled. Each workload
|
||||
// receives its own deadline so one slow start cannot starve the others.
|
||||
for i := len(prepared.Stopped) - 1; i >= 0; i-- {
|
||||
source := prepared.Stopped[i]
|
||||
if err := w.start(ctx, prepared.Kind, source.WorkloadID); err != nil && first == nil {
|
||||
first = err
|
||||
startCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 120*time.Second)
|
||||
if err := w.start(startCtx, prepared.Kind, source.WorkloadID); err != nil {
|
||||
failures = append(failures, fmt.Errorf("restart %s: %w", source.WorkloadID, err))
|
||||
}
|
||||
cancel()
|
||||
}
|
||||
if prepared.MetadataDir != "" {
|
||||
_ = os.RemoveAll(prepared.MetadataDir)
|
||||
if err := os.RemoveAll(prepared.MetadataDir); err != nil {
|
||||
failures = append(failures, err)
|
||||
}
|
||||
}
|
||||
return errors.Join(failures...)
|
||||
}
|
||||
|
||||
func (w *WorkloadManager) running(ctx context.Context, kind model.JobType, id string) (bool, error) {
|
||||
if kind == model.JobDocker {
|
||||
output, err := exec.CommandContext(ctx, "docker", "inspect", "--format", "{{json .State}}", id).Output()
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("source: inspect state of %s: %w", id, err)
|
||||
}
|
||||
var state struct {
|
||||
Running bool
|
||||
Paused bool
|
||||
Restarting bool
|
||||
}
|
||||
if err := json.Unmarshal(output, &state); err != nil {
|
||||
return false, err
|
||||
}
|
||||
if state.Paused || state.Restarting {
|
||||
return false, fmt.Errorf("source: %s is paused or restarting; refusing to change its state", id)
|
||||
}
|
||||
return state.Running, nil
|
||||
}
|
||||
output, err := exec.CommandContext(ctx, "virsh", "domstate", id).Output()
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("source: inspect VM state %s: %w", id, err)
|
||||
}
|
||||
switch strings.TrimSpace(string(output)) {
|
||||
case "running":
|
||||
return true, nil
|
||||
case "shut off":
|
||||
return false, nil
|
||||
default:
|
||||
return false, fmt.Errorf("source: VM %s is in an unsupported state", id)
|
||||
}
|
||||
return first
|
||||
}
|
||||
|
||||
func (w *WorkloadManager) captureMetadata(ctx context.Context, kind model.JobType, id, dir string) ([]string, error) {
|
||||
@@ -249,7 +289,10 @@ func (w *WorkloadManager) captureMetadata(ctx context.Context, kind model.JobTyp
|
||||
var paths []string
|
||||
if len(inspected) > 0 {
|
||||
for _, mount := range inspected[0].Mounts {
|
||||
if mount.Type == "bind" && filepath.IsAbs(mount.Source) {
|
||||
if mount.Type == "bind" || mount.Type == "volume" {
|
||||
if !filepath.IsAbs(mount.Source) {
|
||||
return nil, fmt.Errorf("source: container %s has a %s mount without an absolute backup path", id, mount.Type)
|
||||
}
|
||||
paths = append(paths, mount.Source)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,15 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"git.casaderoll.de/michael/urbm/internal/notify"
|
||||
"git.casaderoll.de/michael/urbm/internal/platform"
|
||||
"git.casaderoll.de/michael/urbm/internal/queue"
|
||||
"git.casaderoll.de/michael/urbm/internal/restic"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -42,3 +51,41 @@ func TestNotificationMessageIncludesPruneSummary(t *testing.T) {
|
||||
t.Fatalf("message contains redundant prune completion details: %q", message)
|
||||
}
|
||||
}
|
||||
|
||||
type checkSecrets struct{}
|
||||
|
||||
func (checkSecrets) Get(string) (string, error) { return "test", nil }
|
||||
|
||||
func TestCheckSendsSuccessAndFailureNotifications(t *testing.T) {
|
||||
for _, code := range []int{0, 1} {
|
||||
t.Run(fmt.Sprint(code), func(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
binary := filepath.Join(dir, "restic")
|
||||
if err := os.WriteFile(binary, []byte(fmt.Sprintf("#!/bin/sh\necho check-result >&2\nexit %d\n", code)), 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var sent []string
|
||||
s := &Service{
|
||||
config: model.Config{Repositories: []model.Repository{{ID: "repo", Name: "test", Type: model.RepositoryLocal, Location: dir}}, Notifications: []model.NotificationTarget{{Type: "unraid", Enabled: true, Events: []string{"success", "failed"}}}},
|
||||
repositoryLocks: map[string]chan struct{}{},
|
||||
restic: &restic.Runner{Binary: binary, RuntimeDir: dir, Secrets: checkSecrets{}},
|
||||
mounts: &platform.MountManager{}, queue: queue.New(nil, nil), log: slog.Default(),
|
||||
notifier: ¬ify.Sender{RunCommand: func(_ context.Context, _ string, args ...string) error {
|
||||
sent = append(sent, strings.Join(args, " "))
|
||||
return nil
|
||||
}},
|
||||
}
|
||||
result := s.execute(context.Background(), model.Run{ID: "run", JobID: "repo", TaskType: "check"})
|
||||
want := "success"
|
||||
if code != 0 {
|
||||
want = "failed"
|
||||
}
|
||||
if result.Status != want || len(sent) != 1 {
|
||||
t.Fatalf("status=%s messages=%v", result.Status, sent)
|
||||
}
|
||||
if !strings.Contains(sent[0], "URBM") {
|
||||
t.Fatal(sent)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1052,8 +1052,12 @@ func (s *Service) execute(ctx context.Context, run model.Run) (result model.Run)
|
||||
ctx = restic.WithRunID(ctx, run.ID)
|
||||
maintenanceRepoName := ""
|
||||
defer func() {
|
||||
if result.TaskType == "prune" && (result.Status == "success" || result.Status == "warning" || result.Status == "failed") {
|
||||
go s.notify("Repository-Bereinigung", maintenanceRepoName, result)
|
||||
if (result.TaskType == "prune" || result.TaskType == "check") && (result.Status == "success" || result.Status == "warning" || result.Status == "failed") {
|
||||
name := "Repository-Bereinigung"
|
||||
if result.TaskType == "check" {
|
||||
name = "Repository-Prüfung"
|
||||
}
|
||||
s.notify(name, maintenanceRepoName, result)
|
||||
}
|
||||
}()
|
||||
if run.TaskType == "backup" {
|
||||
@@ -1543,18 +1547,20 @@ func (s *Service) finishJob(run model.Run, job model.Job, result model.Run) mode
|
||||
repoName = repo.Name
|
||||
}
|
||||
}
|
||||
go s.notify(job.Name, repoName, result)
|
||||
s.notify(job.Name, repoName, result)
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
func (s *Service) notify(name, repository string, run model.Run) {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
title := notificationTitle(name, repository, run)
|
||||
message := notificationMessage(name, repository, run, time.Now().UTC())
|
||||
for _, target := range s.Config().Notifications {
|
||||
for _, event := range target.Events {
|
||||
if event == run.Status {
|
||||
if err := s.notifier.Send(context.Background(), target, title, message, run.Status); err != nil {
|
||||
if err := s.notifier.Send(ctx, target, title, message, run.Status); err != nil {
|
||||
s.log.Warn("send notification failed", "target", target.ID, "type", target.Type, "error", err)
|
||||
}
|
||||
break
|
||||
|
||||
Reference in New Issue
Block a user