Files
vaultik/internal/snapshot/progress.go
T
clawbot c06c4d201b
check / check (push) Waiting to run
Keep the progress line of a multi-path snapshot within 100% (closes #271)
One progress reporter spans every path of a snapshot. Its counts of
files and bytes processed run across all the paths, while the totals
they were divided by were reset to each path's own, 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. SetTotalSize is renamed AddTotalSize because it
now adds.

The issue says the ETA went negative; it was computed negative and
printed as `unknown`.

Model: opus-5-5
2026-10-08 07:29:11 +02:00

492 lines
14 KiB
Go

package snapshot
import (
"context"
"fmt"
"os"
"os/signal"
"sync"
"sync/atomic"
"syscall"
"time"
"github.com/dustin/go-humanize"
"sneak.berlin/go/vaultik/internal/log"
)
const (
// SummaryInterval defines how often one-line status updates are printed.
// These updates show current progress, ETA, and the file being processed.
SummaryInterval = 10 * time.Second
// DetailInterval defines how often multi-line detailed status reports are
// printed. These reports include comprehensive statistics about files,
// chunks, blobs, and uploads.
DetailInterval = 60 * time.Second
// UploadProgressInterval defines how often upload progress messages are logged.
UploadProgressInterval = 15 * time.Second
)
const (
// bitsPerByte converts byte counts to bit counts for speed display.
bitsPerByte = 8
// percentScale converts a ratio to a percentage.
percentScale = 100
// currentFileMaxLen is the display width used for current-file paths.
currentFileMaxLen = 40
// secondsPerMinute and minutesPerHour are used for duration formatting.
secondsPerMinute = 60
minutesPerHour = 60
// Bit-rate thresholds for human-readable upload speed formatting.
bitsPerGbit = 1e9
bitsPerMbit = 1e6
bitsPerKbit = 1e3
// ellipsis prefixes truncated paths and suffixes shortened hashes.
ellipsis = "..."
// hashPrefixLen is how many hex characters of a blob hash to show in logs.
hashPrefixLen = 8
)
// ProgressStats holds atomic counters for progress tracking
type ProgressStats struct {
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
FilesProcessed atomic.Int64 // Files actually processed in phase 2
FilesSkipped atomic.Int64 // Files skipped due to no changes
BytesScanned atomic.Int64 // Bytes from new/changed files only
BytesSkipped atomic.Int64 // Bytes from unchanged files
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
ChunksCreated atomic.Int64
BlobsCreated atomic.Int64
BlobsUploaded atomic.Int64
BytesUploaded atomic.Int64
CurrentFile atomic.Value // stores string
TotalSize atomic.Int64 // Size to process in the paths scanned so far
TotalFiles atomic.Int64 // Files to process in the paths scanned so far
StartTime time.Time
mu sync.RWMutex
lastDetailTime time.Time
// Guarded by mu: when the last scan phase ended, and
// BytesProcessed at that moment.
processStartTime time.Time
processStartBytes int64
// Upload tracking
CurrentUpload atomic.Value // stores *UploadInfo
lastChunkingTime time.Time // Track when we last showed chunking progress
}
// UploadInfo tracks current upload progress
type UploadInfo struct {
BlobHash string
Size int64
StartTime time.Time
LastLogTime time.Time
}
// ProgressReporter handles periodic progress reporting
type ProgressReporter struct {
stats *ProgressStats
ctx context.Context //nolint:containedctx // bound at construction
cancel context.CancelFunc
wg sync.WaitGroup
detailTicker *time.Ticker
summaryTicker *time.Ticker
sigChan chan os.Signal
}
// NewProgressReporter creates a new progress reporter
func NewProgressReporter() *ProgressReporter {
stats := &ProgressStats{
StartTime: time.Now().UTC(),
lastDetailTime: time.Now().UTC(),
}
stats.CurrentFile.Store("")
ctx, cancel := context.WithCancel(context.Background())
pr := &ProgressReporter{
stats: stats,
ctx: ctx,
cancel: cancel,
summaryTicker: time.NewTicker(SummaryInterval),
detailTicker: time.NewTicker(DetailInterval),
sigChan: make(chan os.Signal, 1),
}
// Register for SIGUSR1
signal.Notify(pr.sigChan, syscall.SIGUSR1)
return pr
}
// Start begins the progress reporting
func (pr *ProgressReporter) Start() {
pr.wg.Add(1)
go pr.run()
// Print initial multi-line status
pr.printDetailedStatus()
}
// Stop stops the progress reporting
func (pr *ProgressReporter) Stop() {
pr.cancel()
pr.summaryTicker.Stop()
pr.detailTicker.Stop()
signal.Stop(pr.sigChan)
close(pr.sigChan)
pr.wg.Wait()
}
// GetStats returns the progress stats for updating
func (pr *ProgressReporter) GetStats() *ProgressStats {
return pr.stats
}
// AddTotalSize adds the size one path of the snapshot has to process to
// the total, once that path's scan phase is done, and starts measuring
// the processing rate again. The processed counts run across every path,
// 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
func formatDuration(d time.Duration) string {
if d < 0 {
return "unknown"
}
if d < time.Minute {
return fmt.Sprintf("%ds", int(d.Seconds()))
}
if d < time.Hour {
return fmt.Sprintf("%dm%ds", int(d.Minutes()), int(d.Seconds())%secondsPerMinute)
}
return fmt.Sprintf("%dh%dm", int(d.Hours()), int(d.Minutes())%minutesPerHour)
}
func formatPercent(numerator, denominator int64) string {
if denominator == 0 {
return "0.0%"
}
return fmt.Sprintf("%.1f%%", float64(numerator)/float64(denominator)*percentScale)
}
func formatRatio(compressed, uncompressed int64) string {
if uncompressed == 0 {
return "1.00"
}
ratio := float64(compressed) / float64(uncompressed)
return fmt.Sprintf("%.2f", ratio)
}
func truncatePath(path string, maxLen int) string {
if len(path) <= maxLen {
return path
}
// Keep the last maxLen-len(ellipsis) characters and prepend the ellipsis.
return ellipsis + path[len(path)-(maxLen-len(ellipsis)):]
}
// safeUint64 converts a non-negative int64 counter to uint64 for display,
// clamping negative values to zero.
func safeUint64(n int64) uint64 {
if n < 0 {
return 0
}
return uint64(n)
}
// ReportUploadStart marks the beginning of a blob upload
func (pr *ProgressReporter) ReportUploadStart(blobHash string, size int64) {
info := &UploadInfo{
BlobHash: blobHash,
Size: size,
StartTime: time.Now().UTC(),
}
pr.stats.CurrentUpload.Store(info)
// Log the start of upload
log.Info("Starting blob upload",
"hash", blobHash[:hashPrefixLen]+ellipsis,
"size", humanize.Bytes(safeUint64(size)))
}
// ReportUploadComplete marks the completion of a blob upload
func (pr *ProgressReporter) ReportUploadComplete(
blobHash string, size int64, duration time.Duration,
) {
// Clear current upload
pr.stats.CurrentUpload.Store((*UploadInfo)(nil))
// Calculate speed
if duration < time.Millisecond {
duration = time.Millisecond
}
bytesPerSec := float64(size) / duration.Seconds()
bitsPerSec := bytesPerSec * bitsPerByte
// Format speed
var speedStr string
switch {
case bitsPerSec >= bitsPerGbit:
speedStr = fmt.Sprintf("%.1fGbit/sec", bitsPerSec/bitsPerGbit)
case bitsPerSec >= bitsPerMbit:
speedStr = fmt.Sprintf("%.0fMbit/sec", bitsPerSec/bitsPerMbit)
case bitsPerSec >= bitsPerKbit:
speedStr = fmt.Sprintf("%.0fKbit/sec", bitsPerSec/bitsPerKbit)
default:
speedStr = fmt.Sprintf("%.0fbit/sec", bitsPerSec)
}
log.Info("Blob upload completed",
"hash", blobHash[:hashPrefixLen]+ellipsis,
"size", humanize.Bytes(safeUint64(size)),
"duration", formatDuration(duration),
"speed", speedStr)
}
// UpdateChunkingActivity updates the last chunking time
func (pr *ProgressReporter) UpdateChunkingActivity() {
pr.stats.mu.Lock()
pr.stats.lastChunkingTime = time.Now().UTC()
pr.stats.mu.Unlock()
}
// ReportUploadProgress reports current upload progress with instantaneous speed
func (pr *ProgressReporter) ReportUploadProgress(
blobHash string, bytesUploaded, totalSize int64, instantSpeed float64,
) {
// Update the current upload info with progress
uploadInfo, ok := pr.stats.CurrentUpload.Load().(*UploadInfo)
if ok && uploadInfo != nil {
now := time.Now()
// Only log at the configured interval
if now.Sub(uploadInfo.LastLogTime) >= UploadProgressInterval {
// Format speed in bits/second using humanize
bitsPerSec := instantSpeed * bitsPerByte
speedStr := humanize.SI(bitsPerSec, "bit/sec")
percent := float64(bytesUploaded) / float64(totalSize) * percentScale
// Calculate ETA based on current speed
etaStr := "unknown"
if instantSpeed > 0 && bytesUploaded < totalSize {
remainingBytes := totalSize - bytesUploaded
remainingSeconds := float64(remainingBytes) / instantSpeed
eta := time.Duration(remainingSeconds * float64(time.Second))
etaStr = formatDuration(eta)
}
log.Info("Blob upload progress",
"hash", blobHash[:hashPrefixLen]+ellipsis,
"progress", fmt.Sprintf("%.1f%%", percent),
"uploaded", humanize.Bytes(safeUint64(bytesUploaded)),
"total", humanize.Bytes(safeUint64(totalSize)),
"speed", speedStr,
"eta", etaStr)
uploadInfo.LastLogTime = now
}
}
}
// run is the main progress reporting loop
func (pr *ProgressReporter) run() {
defer pr.wg.Done()
for {
select {
case <-pr.ctx.Done():
return
case <-pr.summaryTicker.C:
pr.printSummaryStatus()
case <-pr.detailTicker.C:
pr.printDetailedStatus()
case <-pr.sigChan:
// SIGUSR1 received, print detailed status
log.Info("SIGUSR1 received, printing detailed status")
pr.printDetailedStatus()
}
}
}
// printSummaryStatus prints a one-line status update
func (pr *ProgressReporter) printSummaryStatus() {
// Check if we're currently uploading
uploadInfo, ok := pr.stats.CurrentUpload.Load().(*UploadInfo)
if ok && uploadInfo != nil {
// Show upload progress instead
pr.printUploadProgress(uploadInfo)
return
}
// Only show chunking progress if we've done chunking recently
pr.stats.mu.RLock()
timeSinceLastChunk := time.Since(pr.stats.lastChunkingTime)
pr.stats.mu.RUnlock()
if timeSinceLastChunk > SummaryInterval*2 {
// No recent chunking activity, don't show progress
return
}
elapsed := time.Since(pr.stats.StartTime)
bytesScanned := pr.stats.BytesScanned.Load()
bytesSkipped := pr.stats.BytesSkipped.Load()
bytesProcessed := pr.stats.BytesProcessed.Load()
totalSize := pr.stats.TotalSize.Load()
currentFile, _ := pr.stats.CurrentFile.Load().(string)
// Calculate ETA if we have total size and are processing
etaStr := ""
processRate := pr.stats.processRate()
if totalSize > 0 && processRate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / processRate
eta := time.Duration(remainingSeconds * float64(time.Second))
etaStr = " | ETA: " + formatDuration(eta)
}
rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds()
// Show files processed / total files to process
filesProcessed := pr.stats.FilesProcessed.Load()
totalFiles := pr.stats.TotalFiles.Load()
status := fmt.Sprintf("Snapshot progress: %d/%d files, %s/%s (%.1f%%), %s/s%s",
filesProcessed,
totalFiles,
humanize.Bytes(safeUint64(bytesProcessed)),
humanize.Bytes(safeUint64(totalSize)),
float64(bytesProcessed)/float64(totalSize)*percentScale,
humanize.Bytes(uint64(rate)),
etaStr,
)
if currentFile != "" {
status += " | Current: " + truncatePath(currentFile, currentFileMaxLen)
}
log.Info(status)
}
// printDetailedStatus prints a multi-line detailed status
func (pr *ProgressReporter) printDetailedStatus() {
pr.stats.mu.Lock()
pr.stats.lastDetailTime = time.Now().UTC()
pr.stats.mu.Unlock()
elapsed := time.Since(pr.stats.StartTime)
filesScanned := pr.stats.FilesScanned.Load()
filesSkipped := pr.stats.FilesSkipped.Load()
bytesScanned := pr.stats.BytesScanned.Load()
bytesSkipped := pr.stats.BytesSkipped.Load()
bytesProcessed := pr.stats.BytesProcessed.Load()
totalSize := pr.stats.TotalSize.Load()
chunksCreated := pr.stats.ChunksCreated.Load()
blobsCreated := pr.stats.BlobsCreated.Load()
blobsUploaded := pr.stats.BlobsUploaded.Load()
bytesUploaded := pr.stats.BytesUploaded.Load()
currentFile, _ := pr.stats.CurrentFile.Load().(string)
totalBytes := bytesScanned + bytesSkipped
rate := float64(totalBytes) / elapsed.Seconds()
log.Notice("=== Snapshot Progress Report ===")
log.Info("Elapsed time", "duration", formatDuration(elapsed))
// Calculate and show ETA if we have data
processRate := pr.stats.processRate()
if totalSize > 0 && processRate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / processRate
eta := time.Duration(remainingSeconds * float64(time.Second))
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
log.Info("Overall progress",
"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",
"scanned", filesScanned,
"skipped", filesSkipped,
"total", filesScanned,
"skip_rate", formatPercent(filesSkipped, filesScanned))
log.Info("Data scanned",
"new", humanize.Bytes(safeUint64(bytesScanned)),
"skipped", humanize.Bytes(safeUint64(bytesSkipped)),
"total", humanize.Bytes(safeUint64(totalBytes)),
"scan_rate", humanize.Bytes(uint64(rate))+"/s")
log.Info("Chunks created", "count", chunksCreated)
log.Info("Blobs status",
"created", blobsCreated,
"uploaded", blobsUploaded,
"pending", blobsCreated-blobsUploaded)
log.Info("Total uploaded to remote",
"uploaded", humanize.Bytes(safeUint64(bytesUploaded)),
"compression_ratio", formatRatio(bytesUploaded, bytesScanned))
if currentFile != "" {
log.Info("Current file", "path", currentFile)
}
log.Notice("=============================")
}
// printUploadProgress prints upload progress
func (pr *ProgressReporter) printUploadProgress(_ *UploadInfo) {
// This function is called repeatedly during upload, not just at start
// Don't print anything here - the actual progress is shown by ReportUploadProgress
}