Keep the progress line of a multi-path snapshot within 100% (closes #271)
check / check (push) Waiting to run
check / check (push) Waiting to run
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 is measured from the bytes processed since the current path's scan phase ended, so a later path's scan phase does not lower it. 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
This commit is contained in:
@@ -56,23 +56,27 @@ const (
|
||||
|
||||
// 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 // Total size to process (set after scan phase)
|
||||
TotalFiles atomic.Int64 // Total files to process in phase 2
|
||||
ProcessStartTime atomic.Value // stores time.Time when processing starts
|
||||
StartTime time.Time
|
||||
mu sync.RWMutex
|
||||
lastDetailTime time.Time
|
||||
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 current path's processing started, and
|
||||
// BytesProcessed at that moment.
|
||||
processStartTime time.Time
|
||||
processStartBytes int64
|
||||
|
||||
// Upload tracking
|
||||
CurrentUpload atomic.Value // stores *UploadInfo
|
||||
@@ -148,10 +152,36 @@ func (pr *ProgressReporter) GetStats() *ProgressStats {
|
||||
return pr.stats
|
||||
}
|
||||
|
||||
// SetTotalSize sets the total size to process (after scan phase)
|
||||
func (pr *ProgressReporter) SetTotalSize(size int64) {
|
||||
pr.stats.TotalSize.Store(size)
|
||||
pr.stats.ProcessStartTime.Store(time.Now().UTC())
|
||||
// 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. The rate covers the current path's processing
|
||||
// only, so that the scan phase of a later path, which processes nothing,
|
||||
// does not lower it.
|
||||
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 current
|
||||
// path's processing started, 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
|
||||
@@ -357,19 +387,12 @@ func (pr *ProgressReporter) printSummaryStatus() {
|
||||
// Calculate ETA if we have total size and are processing
|
||||
etaStr := ""
|
||||
|
||||
if totalSize > 0 && bytesProcessed > 0 {
|
||||
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
|
||||
if ok && !processStart.IsZero() {
|
||||
processElapsed := time.Since(processStart)
|
||||
|
||||
rate := float64(bytesProcessed) / processElapsed.Seconds()
|
||||
if rate > 0 {
|
||||
remainingBytes := totalSize - bytesProcessed
|
||||
remainingSeconds := float64(remainingBytes) / rate
|
||||
eta := time.Duration(remainingSeconds * float64(time.Second))
|
||||
etaStr = " | ETA: " + formatDuration(eta)
|
||||
}
|
||||
}
|
||||
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()
|
||||
@@ -421,25 +444,18 @@ func (pr *ProgressReporter) printDetailedStatus() {
|
||||
log.Info("Elapsed time", "duration", formatDuration(elapsed))
|
||||
|
||||
// Calculate and show ETA if we have data
|
||||
if totalSize > 0 && bytesProcessed > 0 {
|
||||
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
|
||||
if ok && !processStart.IsZero() {
|
||||
processElapsed := time.Since(processStart)
|
||||
|
||||
processRate := float64(bytesProcessed) / processElapsed.Seconds()
|
||||
if 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))
|
||||
}
|
||||
}
|
||||
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",
|
||||
|
||||
Reference in New Issue
Block a user