Keep the progress line of a multi-path snapshot within 100% #278
@@ -22,6 +22,17 @@ the tag exists and is exercised; what is left is merging `next` to
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
|
- 2026-10-08: Kept the progress line of a snapshot with more than one
|
||||||
|
path within 100%
|
||||||
|
([issue #271](https://git.eeqj.de/sneak/vaultik/issues/271)). The
|
||||||
|
bytes and files processed counted every path, but the totals they
|
||||||
|
were divided by held only the current path's, so a second path
|
||||||
|
smaller than the first showed more than 100% and an ETA of `unknown`.
|
||||||
|
The totals now add up over the paths scanned so far. The rate behind
|
||||||
|
the ETA restarts when each path's scan phase ends, so that phase does
|
||||||
|
not lower the rate while the path is processed; until then the rate
|
||||||
|
is the previous path's.
|
||||||
|
|
||||||
- 2026-10-08: Made a second `snapshot create` of one name succeed when
|
- 2026-10-08: Made a second `snapshot create` of one name succeed when
|
||||||
it starts in the same second as the first
|
it starts in the same second as the first
|
||||||
([issue #270](https://git.eeqj.de/sneak/vaultik/issues/270)). The
|
([issue #270](https://git.eeqj.de/sneak/vaultik/issues/270)). The
|
||||||
|
|||||||
@@ -56,23 +56,27 @@ const (
|
|||||||
|
|
||||||
// ProgressStats holds atomic counters for progress tracking
|
// ProgressStats holds atomic counters for progress tracking
|
||||||
type ProgressStats struct {
|
type ProgressStats struct {
|
||||||
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
|
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
|
||||||
FilesProcessed atomic.Int64 // Files actually processed in phase 2
|
FilesProcessed atomic.Int64 // Files actually processed in phase 2
|
||||||
FilesSkipped atomic.Int64 // Files skipped due to no changes
|
FilesSkipped atomic.Int64 // Files skipped due to no changes
|
||||||
BytesScanned atomic.Int64 // Bytes from new/changed files only
|
BytesScanned atomic.Int64 // Bytes from new/changed files only
|
||||||
BytesSkipped atomic.Int64 // Bytes from unchanged files
|
BytesSkipped atomic.Int64 // Bytes from unchanged files
|
||||||
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
|
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
|
||||||
ChunksCreated atomic.Int64
|
ChunksCreated atomic.Int64
|
||||||
BlobsCreated atomic.Int64
|
BlobsCreated atomic.Int64
|
||||||
BlobsUploaded atomic.Int64
|
BlobsUploaded atomic.Int64
|
||||||
BytesUploaded atomic.Int64
|
BytesUploaded atomic.Int64
|
||||||
CurrentFile atomic.Value // stores string
|
CurrentFile atomic.Value // stores string
|
||||||
TotalSize atomic.Int64 // Total size to process (set after scan phase)
|
TotalSize atomic.Int64 // Size to process in the paths scanned so far
|
||||||
TotalFiles atomic.Int64 // Total files to process in phase 2
|
TotalFiles atomic.Int64 // Files to process in the paths scanned so far
|
||||||
ProcessStartTime atomic.Value // stores time.Time when processing starts
|
StartTime time.Time
|
||||||
StartTime time.Time
|
mu sync.RWMutex
|
||||||
mu sync.RWMutex
|
lastDetailTime time.Time
|
||||||
lastDetailTime time.Time
|
|
||||||
|
// Guarded by mu: when the last scan phase ended, and
|
||||||
|
// BytesProcessed at that moment.
|
||||||
|
processStartTime time.Time
|
||||||
|
processStartBytes int64
|
||||||
|
|
||||||
// Upload tracking
|
// Upload tracking
|
||||||
CurrentUpload atomic.Value // stores *UploadInfo
|
CurrentUpload atomic.Value // stores *UploadInfo
|
||||||
@@ -148,10 +152,36 @@ func (pr *ProgressReporter) GetStats() *ProgressStats {
|
|||||||
return pr.stats
|
return pr.stats
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetTotalSize sets the total size to process (after scan phase)
|
// AddTotalSize adds the size one path of the snapshot has to process to
|
||||||
func (pr *ProgressReporter) SetTotalSize(size int64) {
|
// the total, once that path's scan phase is done, and starts measuring
|
||||||
pr.stats.TotalSize.Store(size)
|
// the processing rate again. The processed counts run across every path,
|
||||||
pr.stats.ProcessStartTime.Store(time.Now().UTC())
|
// so the total does too. Restarting the rate keeps a path's scan phase,
|
||||||
|
// which processes nothing, out of the rate while that path is processed.
|
||||||
|
// During a later path's scan phase the rate is still the previous
|
||||||
|
// path's, and falls as that scan goes on.
|
||||||
|
func (pr *ProgressReporter) AddTotalSize(size int64) {
|
||||||
|
pr.stats.TotalSize.Add(size)
|
||||||
|
|
||||||
|
pr.stats.mu.Lock()
|
||||||
|
defer pr.stats.mu.Unlock()
|
||||||
|
|
||||||
|
pr.stats.processStartTime = time.Now().UTC()
|
||||||
|
pr.stats.processStartBytes = pr.stats.BytesProcessed.Load()
|
||||||
|
}
|
||||||
|
|
||||||
|
// processRate returns the bytes processed per second since the last scan
|
||||||
|
// phase ended, or 0 before the first path's scan phase is done.
|
||||||
|
func (s *ProgressStats) processRate() float64 {
|
||||||
|
s.mu.RLock()
|
||||||
|
defer s.mu.RUnlock()
|
||||||
|
|
||||||
|
if s.processStartTime.IsZero() {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
|
||||||
|
processed := s.BytesProcessed.Load() - s.processStartBytes
|
||||||
|
|
||||||
|
return float64(processed) / time.Since(s.processStartTime).Seconds()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Helper functions
|
// Helper functions
|
||||||
@@ -357,19 +387,12 @@ func (pr *ProgressReporter) printSummaryStatus() {
|
|||||||
// Calculate ETA if we have total size and are processing
|
// Calculate ETA if we have total size and are processing
|
||||||
etaStr := ""
|
etaStr := ""
|
||||||
|
|
||||||
if totalSize > 0 && bytesProcessed > 0 {
|
processRate := pr.stats.processRate()
|
||||||
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
|
if totalSize > 0 && processRate > 0 {
|
||||||
if ok && !processStart.IsZero() {
|
remainingBytes := totalSize - bytesProcessed
|
||||||
processElapsed := time.Since(processStart)
|
remainingSeconds := float64(remainingBytes) / processRate
|
||||||
|
eta := time.Duration(remainingSeconds * float64(time.Second))
|
||||||
rate := float64(bytesProcessed) / processElapsed.Seconds()
|
etaStr = " | ETA: " + formatDuration(eta)
|
||||||
if rate > 0 {
|
|
||||||
remainingBytes := totalSize - bytesProcessed
|
|
||||||
remainingSeconds := float64(remainingBytes) / rate
|
|
||||||
eta := time.Duration(remainingSeconds * float64(time.Second))
|
|
||||||
etaStr = " | ETA: " + formatDuration(eta)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds()
|
rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds()
|
||||||
@@ -421,25 +444,18 @@ func (pr *ProgressReporter) printDetailedStatus() {
|
|||||||
log.Info("Elapsed time", "duration", formatDuration(elapsed))
|
log.Info("Elapsed time", "duration", formatDuration(elapsed))
|
||||||
|
|
||||||
// Calculate and show ETA if we have data
|
// Calculate and show ETA if we have data
|
||||||
if totalSize > 0 && bytesProcessed > 0 {
|
processRate := pr.stats.processRate()
|
||||||
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
|
if totalSize > 0 && processRate > 0 {
|
||||||
if ok && !processStart.IsZero() {
|
remainingBytes := totalSize - bytesProcessed
|
||||||
processElapsed := time.Since(processStart)
|
remainingSeconds := float64(remainingBytes) / processRate
|
||||||
|
eta := time.Duration(remainingSeconds * float64(time.Second))
|
||||||
processRate := float64(bytesProcessed) / processElapsed.Seconds()
|
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
|
||||||
if processRate > 0 {
|
log.Info("Overall progress",
|
||||||
remainingBytes := totalSize - bytesProcessed
|
"percent", fmt.Sprintf("%.1f%%", percentComplete),
|
||||||
remainingSeconds := float64(remainingBytes) / processRate
|
"processed", humanize.Bytes(safeUint64(bytesProcessed)),
|
||||||
eta := time.Duration(remainingSeconds * float64(time.Second))
|
"total", humanize.Bytes(safeUint64(totalSize)),
|
||||||
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
|
"rate", humanize.Bytes(uint64(processRate))+"/s",
|
||||||
log.Info("Overall progress",
|
"eta", formatDuration(eta))
|
||||||
"percent", fmt.Sprintf("%.1f%%", percentComplete),
|
|
||||||
"processed", humanize.Bytes(safeUint64(bytesProcessed)),
|
|
||||||
"total", humanize.Bytes(safeUint64(totalSize)),
|
|
||||||
"rate", humanize.Bytes(uint64(processRate))+"/s",
|
|
||||||
"eta", formatDuration(eta))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("Files processed",
|
log.Info("Files processed",
|
||||||
|
|||||||
@@ -0,0 +1,54 @@
|
|||||||
|
//nolint:testpackage // exercises the unexported processRate helper
|
||||||
|
package snapshot
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestProcessRateNotLoweredByLaterScanPhase feeds the reporter two paths
|
||||||
|
// that each take 10 seconds to process, the second after a 30-second scan
|
||||||
|
// phase. That scan phase processes nothing, so the rate while the second
|
||||||
|
// path is processed must be that path's own.
|
||||||
|
func TestProcessRateNotLoweredByLaterScanPhase(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
const (
|
||||||
|
pathSize = 1000
|
||||||
|
processingTime = 10 * time.Second
|
||||||
|
scanTime = 30 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// Never started, but Stop releases its tickers and signal handler.
|
||||||
|
progress := NewProgressReporter()
|
||||||
|
defer progress.Stop()
|
||||||
|
|
||||||
|
stats := progress.GetStats()
|
||||||
|
|
||||||
|
// Moves the processing start time back instead of sleeping.
|
||||||
|
elapse := func(d time.Duration) {
|
||||||
|
stats.mu.Lock()
|
||||||
|
defer stats.mu.Unlock()
|
||||||
|
|
||||||
|
stats.processStartTime = stats.processStartTime.Add(-d)
|
||||||
|
}
|
||||||
|
|
||||||
|
progress.AddTotalSize(pathSize)
|
||||||
|
stats.BytesProcessed.Add(pathSize)
|
||||||
|
elapse(processingTime)
|
||||||
|
|
||||||
|
elapse(scanTime)
|
||||||
|
progress.AddTotalSize(pathSize)
|
||||||
|
stats.BytesProcessed.Add(pathSize)
|
||||||
|
elapse(processingTime)
|
||||||
|
|
||||||
|
want := pathSize / processingTime.Seconds()
|
||||||
|
got := stats.processRate()
|
||||||
|
|
||||||
|
// The test's own run time adds to the elapsed time, so got is a hair
|
||||||
|
// under want.
|
||||||
|
if got > want || got < want*0.99 {
|
||||||
|
t.Errorf("rate while the second path is processed is %.1f bytes/s, "+
|
||||||
|
"want %.1f", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,94 @@
|
|||||||
|
package snapshot_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/spf13/afero"
|
||||||
|
"sneak.berlin/go/vaultik/internal/database"
|
||||||
|
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TestProgressPercentWithSecondPathSmaller backs up two paths with one
|
||||||
|
// scanner, as a snapshot with two paths does, the second path smaller
|
||||||
|
// than the first. The progress line divides the bytes processed by the
|
||||||
|
// total size, so both must count both paths to stay within 100%.
|
||||||
|
func TestProgressPercentWithSecondPathSmaller(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
fs := afero.NewMemMapFs()
|
||||||
|
files := map[string]string{
|
||||||
|
"/large/one.txt": strings.Repeat("1", 4000),
|
||||||
|
"/large/two.txt": strings.Repeat("2", 4000),
|
||||||
|
"/small/three.txt": strings.Repeat("3", 1000),
|
||||||
|
}
|
||||||
|
|
||||||
|
for path, content := range files {
|
||||||
|
err := fs.MkdirAll(filepath.Dir(path), 0755)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("mkdir: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
err = afero.WriteFile(fs, path, []byte(content), 0644)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("write %s: %v", path, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
db, err := database.NewTestDB()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("create test db: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
t.Cleanup(func() {
|
||||||
|
cerr := db.Close()
|
||||||
|
if cerr != nil {
|
||||||
|
t.Errorf("close db: %v", cerr)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
repos := database.NewRepositories(db)
|
||||||
|
|
||||||
|
scanner := snapshot.NewScanner(snapshot.ScannerConfig{
|
||||||
|
FS: fs,
|
||||||
|
ChunkSize: int64(1024 * 16),
|
||||||
|
Repositories: repos,
|
||||||
|
MaxBlobSize: int64(1024 * 1024),
|
||||||
|
CompressionLevel: 3,
|
||||||
|
AgeRecipients: []string{testAgePublicKey},
|
||||||
|
EnableProgress: true,
|
||||||
|
})
|
||||||
|
|
||||||
|
// Never started, but Stop releases its tickers and signal handler.
|
||||||
|
progress := scanner.GetProgress()
|
||||||
|
defer progress.Stop()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
snapshotID := "test-snapshot-progress"
|
||||||
|
createTestSnapshotRecord(ctx, t, repos, snapshotID)
|
||||||
|
|
||||||
|
for _, path := range []string{"/large", "/small"} {
|
||||||
|
_, err := scanner.Scan(ctx, path, snapshotID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("scanning %s: %v", path, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
stats := progress.GetStats()
|
||||||
|
|
||||||
|
// Directories count toward the total size but produce no chunks, so
|
||||||
|
// the percentage ends just under 100%.
|
||||||
|
percent := float64(stats.BytesProcessed.Load()) /
|
||||||
|
float64(stats.TotalSize.Load()) * 100
|
||||||
|
if percent > 100 {
|
||||||
|
t.Errorf("progress after both paths is %.1f%%, want at most 100%%",
|
||||||
|
percent)
|
||||||
|
}
|
||||||
|
|
||||||
|
if stats.FilesProcessed.Load() != stats.TotalFiles.Load() {
|
||||||
|
t.Errorf("progress after both paths is %d of %d files, want all",
|
||||||
|
stats.FilesProcessed.Load(), stats.TotalFiles.Load())
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -407,8 +407,8 @@ func (s *Scanner) summarizeScanPhase(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if s.progress != nil {
|
if s.progress != nil {
|
||||||
s.progress.SetTotalSize(totalSizeToProcess)
|
s.progress.AddTotalSize(totalSizeToProcess)
|
||||||
s.progress.GetStats().TotalFiles.Store(int64(len(filesToProcess)))
|
s.progress.GetStats().TotalFiles.Add(int64(len(filesToProcess)))
|
||||||
}
|
}
|
||||||
|
|
||||||
log.Info("Phase 1 complete",
|
log.Info("Phase 1 complete",
|
||||||
|
|||||||
Reference in New Issue
Block a user