Add live activity progress and logs
This commit is contained in:
+26
-17
@@ -112,23 +112,32 @@ type NotificationTarget struct {
|
||||
}
|
||||
|
||||
type Run struct {
|
||||
SchemaVersion int `json:"schemaVersion"`
|
||||
ID string `json:"id"`
|
||||
JobID string `json:"jobId,omitempty"`
|
||||
TaskType string `json:"taskType"`
|
||||
Status string `json:"status"`
|
||||
Priority int `json:"priority"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
StartedAt *time.Time `json:"startedAt,omitempty"`
|
||||
FinishedAt *time.Time `json:"finishedAt,omitempty"`
|
||||
SnapshotID string `json:"snapshotId,omitempty"`
|
||||
Message string `json:"message,omitempty"`
|
||||
ErrorCode string `json:"errorCode,omitempty"`
|
||||
BytesAdded int64 `json:"bytesAdded,omitempty"`
|
||||
FilesNew int64 `json:"filesNew,omitempty"`
|
||||
FilesChanged int64 `json:"filesChanged,omitempty"`
|
||||
BytesProcessed int64 `json:"bytesProcessed,omitempty"`
|
||||
FilesProcessed int64 `json:"filesProcessed,omitempty"`
|
||||
SchemaVersion int `json:"schemaVersion"`
|
||||
ID string `json:"id"`
|
||||
JobID string `json:"jobId,omitempty"`
|
||||
TaskType string `json:"taskType"`
|
||||
Status string `json:"status"`
|
||||
Priority int `json:"priority"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
StartedAt *time.Time `json:"startedAt,omitempty"`
|
||||
FinishedAt *time.Time `json:"finishedAt,omitempty"`
|
||||
SnapshotID string `json:"snapshotId,omitempty"`
|
||||
Message string `json:"message,omitempty"`
|
||||
ErrorCode string `json:"errorCode,omitempty"`
|
||||
BytesAdded int64 `json:"bytesAdded,omitempty"`
|
||||
FilesNew int64 `json:"filesNew,omitempty"`
|
||||
FilesChanged int64 `json:"filesChanged,omitempty"`
|
||||
BytesProcessed int64 `json:"bytesProcessed,omitempty"`
|
||||
FilesProcessed int64 `json:"filesProcessed,omitempty"`
|
||||
ProgressPercent float64 `json:"progressPercent,omitempty"`
|
||||
ProgressBytes int64 `json:"progressBytes,omitempty"`
|
||||
ProgressTotal int64 `json:"progressTotal,omitempty"`
|
||||
ProgressFiles int64 `json:"progressFiles,omitempty"`
|
||||
ProgressFileTotal int64 `json:"progressFileTotal,omitempty"`
|
||||
BytesPerSecond float64 `json:"bytesPerSecond,omitempty"`
|
||||
SecondsRemaining int64 `json:"secondsRemaining,omitempty"`
|
||||
CurrentFile string `json:"currentFile,omitempty"`
|
||||
LiveLog []string `json:"liveLog,omitempty"`
|
||||
}
|
||||
|
||||
type RestoreTask struct {
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/backupper-unraid/backupper/internal/model"
|
||||
)
|
||||
|
||||
func TestUpdateActiveIsVisibleWithoutPersistingIntermediateState(t *testing.T) {
|
||||
started := make(chan struct{})
|
||||
release := make(chan struct{})
|
||||
persistCalls := 0
|
||||
q := New(func(_ context.Context, run model.Run) model.Run {
|
||||
close(started)
|
||||
<-release
|
||||
run.Status = "success"
|
||||
return run
|
||||
}, func([]model.Run) { persistCalls++ })
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go q.Run(ctx)
|
||||
if err := q.Enqueue(model.Run{ID: "run", JobID: "job"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
<-started
|
||||
before := persistCalls
|
||||
if !q.UpdateActive("run", func(run *model.Run) { run.ProgressPercent = 42 }) {
|
||||
t.Fatal("active run was not updated")
|
||||
}
|
||||
if runs := q.Snapshot(); len(runs) != 1 || runs[0].ProgressPercent != 42 {
|
||||
t.Fatalf("runs = %#v", runs)
|
||||
}
|
||||
if persistCalls != before {
|
||||
t.Fatal("intermediate progress was persisted")
|
||||
}
|
||||
close(release)
|
||||
deadline := time.Now().Add(time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
if runs := q.Snapshot(); len(runs) == 1 && runs[0].Status == "success" {
|
||||
return
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
t.Fatal("run did not finish")
|
||||
}
|
||||
@@ -141,6 +141,16 @@ func (q *Queue) Snapshot() []model.Run {
|
||||
return q.snapshotLocked()
|
||||
}
|
||||
|
||||
func (q *Queue) UpdateActive(id string, update func(*model.Run)) bool {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
if q.active == nil || q.active.ID != id {
|
||||
return false
|
||||
}
|
||||
update(q.active)
|
||||
return true
|
||||
}
|
||||
|
||||
func (q *Queue) snapshotLocked() []model.Run {
|
||||
result := append([]model.Run{}, q.history...)
|
||||
if q.active != nil {
|
||||
|
||||
@@ -35,11 +35,24 @@ type Summary struct {
|
||||
TotalFilesProcessed int64 `json:"total_files_processed"`
|
||||
}
|
||||
|
||||
type Progress struct {
|
||||
MessageType string `json:"message_type"`
|
||||
PercentDone float64 `json:"percent_done"`
|
||||
TotalFiles int64 `json:"total_files"`
|
||||
FilesDone int64 `json:"files_done"`
|
||||
TotalBytes int64 `json:"total_bytes"`
|
||||
BytesDone int64 `json:"bytes_done"`
|
||||
SecondsRemaining int64 `json:"seconds_remaining"`
|
||||
CurrentFiles []string `json:"current_files"`
|
||||
}
|
||||
|
||||
type ProgressCallback func(Progress)
|
||||
|
||||
func (r *Runner) Init(ctx context.Context, repo model.Repository) error {
|
||||
return r.run(ctx, repo, []string{"init"}, nil, nil)
|
||||
}
|
||||
|
||||
func (r *Runner) Backup(ctx context.Context, repo model.Repository, job model.Job, sources []string) (Summary, error) {
|
||||
func (r *Runner) Backup(ctx context.Context, repo model.Repository, job model.Job, sources []string, progress ProgressCallback) (Summary, error) {
|
||||
args := []string{"backup", "--json", "--compression", job.Compression}
|
||||
for _, tag := range append([]string{"backupper", "job:" + job.ID}, job.Tags...) {
|
||||
args = append(args, "--tag", tag)
|
||||
@@ -50,6 +63,13 @@ func (r *Runner) Backup(ctx context.Context, repo model.Repository, job model.Jo
|
||||
args = append(args, sources...)
|
||||
var summary Summary
|
||||
err := r.run(ctx, repo, args, func(line []byte) {
|
||||
var status Progress
|
||||
if json.Unmarshal(line, &status) == nil && status.MessageType == "status" {
|
||||
if progress != nil {
|
||||
progress(status)
|
||||
}
|
||||
return
|
||||
}
|
||||
var candidate Summary
|
||||
if json.Unmarshal(line, &candidate) == nil && candidate.MessageType == "summary" {
|
||||
summary = candidate
|
||||
|
||||
@@ -20,20 +20,24 @@ func TestBackupUsesPasswordFileAndStructuredArguments(t *testing.T) {
|
||||
argsPath := filepath.Join(dir, "args")
|
||||
passwordPath := filepath.Join(dir, "password")
|
||||
script := filepath.Join(dir, "restic")
|
||||
body := fmt.Sprintf("#!/bin/sh\nprintf '%%s\\n' \"$@\" > '%s'\ncat \"$RESTIC_PASSWORD_FILE\" > '%s'\nprintf '%%s\\n' '{\"message_type\":\"summary\",\"snapshot_id\":\"abc123\",\"data_added\":42,\"files_new\":3,\"files_changed\":2,\"total_bytes_processed\":1024,\"total_files_processed\":12}'\n", argsPath, passwordPath)
|
||||
body := fmt.Sprintf("#!/bin/sh\nprintf '%%s\\n' \"$@\" > '%s'\ncat \"$RESTIC_PASSWORD_FILE\" > '%s'\nprintf '%%s\\n' '{\"message_type\":\"status\",\"percent_done\":0.5,\"total_files\":12,\"files_done\":6,\"total_bytes\":1024,\"bytes_done\":512,\"seconds_remaining\":4,\"current_files\":[\"/source path/file.txt\"]}'\nprintf '%%s\\n' '{\"message_type\":\"summary\",\"snapshot_id\":\"abc123\",\"data_added\":42,\"files_new\":3,\"files_changed\":2,\"total_bytes_processed\":1024,\"total_files_processed\":12}'\n", argsPath, passwordPath)
|
||||
if err := os.WriteFile(script, []byte(body), 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
runner := &Runner{Binary: script, RuntimeDir: dir, Secrets: fakeSecrets{"password": "top-secret"}}
|
||||
repo := model.Repository{Type: model.RepositoryLocal, Location: "/repo path", PasswordRef: "password"}
|
||||
job := model.Job{ID: "job", Compression: "auto", Tags: []string{"daily"}, Excludes: []string{"*.tmp"}}
|
||||
summary, err := runner.Backup(context.Background(), repo, job, []string{"/source path"})
|
||||
var progress Progress
|
||||
summary, err := runner.Backup(context.Background(), repo, job, []string{"/source path"}, func(value Progress) { progress = value })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if summary.SnapshotID != "abc123" || summary.DataAdded != 42 || summary.FilesChanged != 2 || summary.TotalBytesProcessed != 1024 || summary.TotalFilesProcessed != 12 {
|
||||
t.Fatalf("unexpected summary: %+v", summary)
|
||||
}
|
||||
if progress.PercentDone != 0.5 || progress.FilesDone != 6 || progress.BytesDone != 512 || len(progress.CurrentFiles) != 1 {
|
||||
t.Fatalf("unexpected progress: %+v", progress)
|
||||
}
|
||||
args, _ := os.ReadFile(argsPath)
|
||||
for _, expected := range []string{"/repo path", "backup", "--json", "/source path"} {
|
||||
if !strings.Contains(string(args), expected) {
|
||||
|
||||
@@ -221,20 +221,33 @@ func (s *Service) execute(ctx context.Context, run model.Run) model.Run {
|
||||
if !ok {
|
||||
return failed(run, "validation", "repository no longer exists")
|
||||
}
|
||||
live := run
|
||||
appendLiveLog(&live, strings.ToUpper(run.TaskType)+" wird vorbereitet")
|
||||
s.updateActive(live)
|
||||
mounted, err := s.mounts.Prepare(ctx, repo)
|
||||
if err != nil {
|
||||
return failedError(run, err)
|
||||
failedRun := failedError(run, err)
|
||||
appendLiveLog(&live, "Vorbereitung fehlgeschlagen: "+err.Error())
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
defer s.mounts.Cleanup(context.Background(), mounted)
|
||||
appendLiveLog(&live, "Repository bereit; Restic "+run.TaskType+" läuft")
|
||||
s.updateActive(live)
|
||||
if run.TaskType == "check" {
|
||||
err = s.restic.Check(ctx, mounted.Repository)
|
||||
} else {
|
||||
err = s.restic.Prune(ctx, mounted.Repository)
|
||||
}
|
||||
if err != nil {
|
||||
return failedError(run, err)
|
||||
failedRun := failedError(run, err)
|
||||
appendLiveLog(&live, strings.ToUpper(run.TaskType)+" fehlgeschlagen: "+err.Error())
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
run.Status, run.Message = "success", run.TaskType+" completed"
|
||||
appendLiveLog(&live, strings.ToUpper(run.TaskType)+" abgeschlossen")
|
||||
run.LiveLog = live.LiveLog
|
||||
return run
|
||||
}
|
||||
|
||||
@@ -248,6 +261,9 @@ func (s *Service) executeBackup(ctx context.Context, run model.Run) (result mode
|
||||
if !ok {
|
||||
return failed(run, "validation", "repository no longer exists")
|
||||
}
|
||||
live := run
|
||||
appendLiveLog(&live, "Backup wird vorbereitet")
|
||||
s.updateActive(live)
|
||||
mounted, err := s.mounts.Prepare(ctx, repo)
|
||||
if err != nil {
|
||||
return failedError(run, err)
|
||||
@@ -277,10 +293,36 @@ func (s *Service) executeBackup(ctx context.Context, run model.Run) (result mode
|
||||
return failed(run, "source", "source unavailable: "+path)
|
||||
}
|
||||
}
|
||||
summary, err := s.restic.Backup(ctx, mounted.Repository, job, prepared.Sources)
|
||||
appendLiveLog(&live, fmt.Sprintf("Restic-Backup gestartet: %d Quelle(n)", len(prepared.Sources)))
|
||||
s.updateActive(live)
|
||||
lastLoggedPercent := -5
|
||||
started := time.Now()
|
||||
summary, err := s.restic.Backup(ctx, mounted.Repository, job, prepared.Sources, func(progress restic.Progress) {
|
||||
live.ProgressPercent = progress.PercentDone * 100
|
||||
live.ProgressBytes, live.ProgressTotal = progress.BytesDone, progress.TotalBytes
|
||||
live.ProgressFiles, live.ProgressFileTotal = progress.FilesDone, progress.TotalFiles
|
||||
live.SecondsRemaining = progress.SecondsRemaining
|
||||
if elapsed := time.Since(started).Seconds(); elapsed > 0 {
|
||||
live.BytesPerSecond = float64(progress.BytesDone) / elapsed
|
||||
}
|
||||
if len(progress.CurrentFiles) > 0 {
|
||||
live.CurrentFile = progress.CurrentFiles[len(progress.CurrentFiles)-1]
|
||||
}
|
||||
percent := int(live.ProgressPercent)
|
||||
if percent >= lastLoggedPercent+5 {
|
||||
appendLiveLog(&live, fmt.Sprintf("Fortschritt: %.1f%%, %d/%d Dateien, %s/%s", live.ProgressPercent, live.ProgressFiles, live.ProgressFileTotal, formatBytes(live.ProgressBytes), formatBytes(live.ProgressTotal)))
|
||||
lastLoggedPercent = percent - percent%5
|
||||
}
|
||||
s.updateActive(live)
|
||||
})
|
||||
if err != nil {
|
||||
return failedError(run, err)
|
||||
appendLiveLog(&live, "Backup fehlgeschlagen: "+err.Error())
|
||||
failedRun := failedError(run, err)
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
appendLiveLog(&live, "Backup-Daten übertragen; Aufbewahrung wird angewendet")
|
||||
s.updateActive(live)
|
||||
if err := s.restic.Forget(ctx, mounted.Repository, job); err != nil {
|
||||
run.Status, run.Message, run.ErrorCode = "warning", "backup succeeded but retention failed: "+err.Error(), "repository"
|
||||
} else {
|
||||
@@ -288,10 +330,36 @@ func (s *Service) executeBackup(ctx context.Context, run model.Run) (result mode
|
||||
}
|
||||
run.SnapshotID, run.BytesAdded, run.FilesNew, run.FilesChanged = summary.SnapshotID, summary.DataAdded, summary.FilesNew, summary.FilesChanged
|
||||
run.BytesProcessed, run.FilesProcessed = summary.TotalBytesProcessed, summary.TotalFilesProcessed
|
||||
run.ProgressPercent, run.ProgressBytes, run.ProgressTotal = 100, live.ProgressBytes, live.ProgressTotal
|
||||
run.ProgressFiles, run.ProgressFileTotal = live.ProgressFiles, live.ProgressFileTotal
|
||||
run.BytesPerSecond, run.SecondsRemaining, run.CurrentFile = live.BytesPerSecond, 0, live.CurrentFile
|
||||
appendLiveLog(&live, "Backup abgeschlossen")
|
||||
run.LiveLog = live.LiveLog
|
||||
return run
|
||||
}
|
||||
|
||||
func (s *Service) updateActive(run model.Run) {
|
||||
s.queue.UpdateActive(run.ID, func(active *model.Run) {
|
||||
active.ProgressPercent, active.ProgressBytes, active.ProgressTotal = run.ProgressPercent, run.ProgressBytes, run.ProgressTotal
|
||||
active.ProgressFiles, active.ProgressFileTotal = run.ProgressFiles, run.ProgressFileTotal
|
||||
active.BytesPerSecond, active.SecondsRemaining = run.BytesPerSecond, run.SecondsRemaining
|
||||
active.CurrentFile = run.CurrentFile
|
||||
active.LiveLog = append([]string(nil), run.LiveLog...)
|
||||
})
|
||||
}
|
||||
|
||||
func appendLiveLog(run *model.Run, message string) {
|
||||
line := time.Now().Format("15:04:05") + " " + strings.ReplaceAll(strings.TrimSpace(message), "\n", " ")
|
||||
run.LiveLog = append(run.LiveLog, line)
|
||||
if len(run.LiveLog) > 100 {
|
||||
run.LiveLog = run.LiveLog[len(run.LiveLog)-100:]
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Service) executeRestore(ctx context.Context, run model.Run) model.Run {
|
||||
live := run
|
||||
appendLiveLog(&live, "Restore wird vorbereitet")
|
||||
s.updateActive(live)
|
||||
parts := strings.Split(run.JobID, "\x00")
|
||||
if len(parts) < 4 {
|
||||
return failed(run, "internal", "invalid restore payload")
|
||||
@@ -302,17 +370,30 @@ func (s *Service) executeRestore(ctx context.Context, run model.Run) model.Run {
|
||||
}
|
||||
task := model.RestoreTask{SchemaVersion: model.SchemaVersion, ID: run.ID, RepositoryID: parts[0], SnapshotID: parts[1], Target: parts[2], InPlace: parts[3] == "true", Confirmed: true, Includes: parts[4:]}
|
||||
if err := os.MkdirAll(task.Target, 0750); err != nil {
|
||||
return failedError(run, err)
|
||||
failedRun := failedError(run, err)
|
||||
appendLiveLog(&live, "Restore-Ziel konnte nicht vorbereitet werden: "+err.Error())
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
mounted, err := s.mounts.Prepare(ctx, repo)
|
||||
if err != nil {
|
||||
return failedError(run, err)
|
||||
failedRun := failedError(run, err)
|
||||
appendLiveLog(&live, "Repository konnte nicht vorbereitet werden: "+err.Error())
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
defer s.mounts.Cleanup(context.Background(), mounted)
|
||||
appendLiveLog(&live, "Restic stellt den Snapshot wieder her")
|
||||
s.updateActive(live)
|
||||
if err := s.restic.Restore(ctx, mounted.Repository, task); err != nil {
|
||||
return failedError(run, err)
|
||||
failedRun := failedError(run, err)
|
||||
appendLiveLog(&live, "Restore fehlgeschlagen: "+err.Error())
|
||||
failedRun.LiveLog = live.LiveLog
|
||||
return failedRun
|
||||
}
|
||||
run.Status, run.Message, run.JobID = "success", "restore completed", ""
|
||||
appendLiveLog(&live, "Restore abgeschlossen")
|
||||
run.LiveLog = live.LiveLog
|
||||
return run
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user