Count each file, byte and upload once in backup statistics #254
+2
-2
@@ -335,10 +335,10 @@ CreateSnapshot(opts)
|
||||
│ │
|
||||
│ └─► Accumulate statistics
|
||||
│
|
||||
├─► SnapshotManager.UpdateSnapshotStatsExtended()
|
||||
│
|
||||
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
|
||||
│
|
||||
├─► SnapshotManager.UpdateSnapshotStatsExtended()
|
||||
│
|
||||
├─► SnapshotManager.ExportSnapshotMetadata()
|
||||
│ │
|
||||
│ ├─► Copy database to temp file
|
||||
|
||||
@@ -22,6 +22,19 @@ the tag exists and is exercised; what is left is merging `next` to
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-10-06: Made the backup summary and the `snapshots` row count each
|
||||
file, byte and upload once
|
||||
([issue #225](https://git.eeqj.de/sneak/vaultik/issues/225)). The
|
||||
scanner added a file's bytes again for each new chunk and counted a
|
||||
file as unchanged for each chunk already stored, so a first backup
|
||||
reported twice its size and "backed up" could go negative. Upload
|
||||
figures came from the progress reporter, which `--cron` turns off, and
|
||||
`blob_count` counted earlier paths' blobs again for each later path.
|
||||
The scanner now counts uploads itself; `blob_size`,
|
||||
`blob_uncompressed_size` and `compression_ratio` describe the blobs
|
||||
the snapshot references, and `docs/DATAMODEL.md` now says
|
||||
`chunk_count` and `blob_count` count what the run added.
|
||||
|
||||
- 2026-10-06: Made command output follow the README's stdout and stderr
|
||||
rules ([issue #224](https://git.eeqj.de/sneak/vaultik/issues/224)). The
|
||||
startup banner went to stdout, so a `completion` script or a
|
||||
|
||||
+3
-3
@@ -117,10 +117,10 @@ Tracks backup snapshots.
|
||||
- `started_at` (INTEGER) - Start timestamp
|
||||
- `completed_at` (INTEGER) - Completion timestamp (NULL if in progress)
|
||||
- `file_count` (INTEGER) - Number of files in snapshot
|
||||
- `chunk_count` (INTEGER) - Number of unique chunks
|
||||
- `blob_count` (INTEGER) - Number of blobs referenced
|
||||
- `chunk_count` (INTEGER) - Number of chunks this snapshot stored that were not stored before
|
||||
- `blob_count` (INTEGER) - Number of blobs this snapshot created
|
||||
- `total_size` (INTEGER) - Total size of all files
|
||||
- `blob_size` (INTEGER) - Total size of all blobs (compressed)
|
||||
- `blob_size` (INTEGER) - Total compressed size of all referenced blobs
|
||||
- `blob_uncompressed_size` (INTEGER) - Total uncompressed size of all referenced blobs
|
||||
- `compression_ratio` (REAL) - Compression ratio achieved
|
||||
- `compression_level` (INTEGER) - Compression level used for this snapshot
|
||||
|
||||
@@ -99,8 +99,8 @@ type Snapshot struct {
|
||||
StartedAt time.Time
|
||||
CompletedAt *time.Time // nil if still in progress
|
||||
FileCount int64
|
||||
ChunkCount int64
|
||||
BlobCount int64
|
||||
ChunkCount int64 // Chunks this snapshot stored that were not stored before
|
||||
BlobCount int64 // Blobs this snapshot created
|
||||
TotalSize int64 // Total size of all referenced files
|
||||
|
||||
// BlobSize is the total size of all referenced blobs (compressed and
|
||||
|
||||
@@ -127,6 +127,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
||||
snapshotID string,
|
||||
blobUncompressedSize int64,
|
||||
compressionLevel int,
|
||||
uploadBytes int64,
|
||||
uploadDurationMs int64,
|
||||
) error {
|
||||
compressionRatio, err := r.extendedCompressionRatio(
|
||||
@@ -141,7 +142,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
||||
SET blob_uncompressed_size = ?,
|
||||
compression_ratio = ?,
|
||||
compression_level = ?,
|
||||
upload_bytes = blob_size,
|
||||
upload_bytes = ?,
|
||||
upload_duration_ms = ?
|
||||
WHERE id = ?
|
||||
`
|
||||
@@ -149,11 +150,11 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
||||
if tx != nil {
|
||||
_, err = tx.ExecContext(ctx, query,
|
||||
blobUncompressedSize, compressionRatio, compressionLevel,
|
||||
uploadDurationMs, snapshotID)
|
||||
uploadBytes, uploadDurationMs, snapshotID)
|
||||
} else {
|
||||
_, err = r.db.ExecWithLog(ctx, query,
|
||||
blobUncompressedSize, compressionRatio, compressionLevel,
|
||||
uploadDurationMs, snapshotID)
|
||||
uploadBytes, uploadDurationMs, snapshotID)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
@@ -543,6 +544,30 @@ func (r *SnapshotRepository) GetSnapshotTotalCompressedSize(
|
||||
return totalSize, nil
|
||||
}
|
||||
|
||||
// GetSnapshotBlobSizes returns the total compressed and uncompressed sizes
|
||||
// of all blobs referenced by a snapshot.
|
||||
func (r *SnapshotRepository) GetSnapshotBlobSizes(
|
||||
ctx context.Context, snapshotID string,
|
||||
) (int64, int64, error) {
|
||||
query := `
|
||||
SELECT COALESCE(SUM(b.compressed_size), 0),
|
||||
COALESCE(SUM(b.uncompressed_size), 0)
|
||||
FROM snapshot_blobs sb
|
||||
JOIN blobs b ON sb.blob_hash = b.blob_hash
|
||||
WHERE sb.snapshot_id = ?
|
||||
`
|
||||
|
||||
var compressed, uncompressed int64
|
||||
|
||||
err := r.db.conn.QueryRowContext(ctx, query, snapshotID).Scan(
|
||||
&compressed, &uncompressed)
|
||||
if err != nil {
|
||||
return 0, 0, fmt.Errorf("querying snapshot blob sizes: %w", err)
|
||||
}
|
||||
|
||||
return compressed, uncompressed, nil
|
||||
}
|
||||
|
||||
// GetSnapshotUncompressedChunkSize returns the sum of plaintext sizes of all unique
|
||||
// chunks referenced by a snapshot (via snapshot_files → file_chunks → chunks).
|
||||
func (r *SnapshotRepository) GetSnapshotUncompressedChunkSize(
|
||||
|
||||
@@ -145,6 +145,65 @@ func TestSnapshotRepositoryUpdateCounts(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// GetSnapshotBlobSizes totals the blobs the snapshot references, and only
|
||||
// those.
|
||||
func TestSnapshotRepositoryGetSnapshotBlobSizes(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
db, cleanup := setupTestDB(t)
|
||||
defer cleanup()
|
||||
|
||||
ctx := context.Background()
|
||||
repos := database.NewRepositories(db)
|
||||
|
||||
snapshot := &database.Snapshot{
|
||||
ID: "2024-01-03T12:00:00Z",
|
||||
Hostname: testHostname,
|
||||
VaultikVersion: testVersion,
|
||||
StartedAt: time.Now().Truncate(time.Second),
|
||||
}
|
||||
|
||||
err := repos.Snapshots.Create(ctx, nil, snapshot)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to create snapshot: %v", err)
|
||||
}
|
||||
|
||||
blobs := []*database.Blob{
|
||||
{Hash: "referenced-1", CompressedSize: 10, UncompressedSize: 100},
|
||||
{Hash: "referenced-2", CompressedSize: 20, UncompressedSize: 200},
|
||||
{Hash: "unreferenced", CompressedSize: 40, UncompressedSize: 400},
|
||||
}
|
||||
|
||||
for _, blob := range blobs {
|
||||
blob.ID = types.NewBlobID()
|
||||
blob.CreatedTS = time.Now().Truncate(time.Second)
|
||||
|
||||
err = repos.Blobs.Create(ctx, nil, blob)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to create blob %s: %v", blob.Hash, err)
|
||||
}
|
||||
}
|
||||
|
||||
for _, blob := range blobs[:2] {
|
||||
err = repos.Snapshots.AddBlob(ctx, nil, snapshot.ID.String(),
|
||||
blob.ID, blob.Hash)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to add blob %s to snapshot: %v", blob.Hash, err)
|
||||
}
|
||||
}
|
||||
|
||||
compressed, uncompressed, err := repos.Snapshots.GetSnapshotBlobSizes(
|
||||
ctx, snapshot.ID.String())
|
||||
if err != nil {
|
||||
t.Fatalf("failed to get snapshot blob sizes: %v", err)
|
||||
}
|
||||
|
||||
if compressed != 30 || uncompressed != 300 {
|
||||
t.Errorf("blob sizes: got %d and %d, want 30 and 300",
|
||||
compressed, uncompressed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSnapshotRepositoryListRecent(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -158,19 +158,3 @@ type UploadStats struct {
|
||||
MinDurationMs int64
|
||||
MaxDurationMs int64
|
||||
}
|
||||
|
||||
// GetCountBySnapshot returns the count of uploads for a specific snapshot
|
||||
func (r *UploadRepository) GetCountBySnapshot(
|
||||
ctx context.Context, snapshotID string,
|
||||
) (int64, error) {
|
||||
query := `SELECT COUNT(*) FROM uploads WHERE snapshot_id = ?`
|
||||
|
||||
var count int64
|
||||
|
||||
err := r.conn.QueryRowContext(ctx, query, snapshotID).Scan(&count)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
return count, nil
|
||||
}
|
||||
|
||||
@@ -66,7 +66,6 @@ type ProgressStats struct {
|
||||
BlobsCreated atomic.Int64
|
||||
BlobsUploaded atomic.Int64
|
||||
BytesUploaded atomic.Int64
|
||||
UploadDurationMs atomic.Int64 // Total milliseconds spent uploading
|
||||
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
|
||||
@@ -231,9 +230,6 @@ func (pr *ProgressReporter) ReportUploadComplete(
|
||||
// Clear current upload
|
||||
pr.stats.CurrentUpload.Store((*UploadInfo)(nil))
|
||||
|
||||
// Add to total upload duration
|
||||
pr.stats.UploadDurationMs.Add(duration.Milliseconds())
|
||||
|
||||
// Calculate speed
|
||||
if duration < time.Millisecond {
|
||||
duration = time.Millisecond
|
||||
|
||||
@@ -92,9 +92,6 @@ type Scanner struct {
|
||||
|
||||
// Mutex for coordinating blob creation
|
||||
packerMu sync.Mutex // Blocks chunk production during blob creation
|
||||
|
||||
// Context for cancellation
|
||||
scanCtx context.Context //nolint:containedctx // set per-Scan for packer callbacks
|
||||
}
|
||||
|
||||
// Periodic status output intervals and thresholds for the scan and
|
||||
@@ -134,18 +131,23 @@ type ScannerConfig struct {
|
||||
SkipErrors bool
|
||||
}
|
||||
|
||||
// ScanResult contains the results of a scan operation
|
||||
// ScanResult contains the results of a scan operation. Files and bytes
|
||||
// are counted per file: BytesScanned is the size of the new and changed
|
||||
// files, BytesSkipped that of the unchanged ones.
|
||||
type ScanResult struct {
|
||||
FilesScanned int
|
||||
FilesSkipped int
|
||||
FilesDeleted int
|
||||
BytesScanned int64
|
||||
BytesSkipped int64
|
||||
BytesDeleted int64
|
||||
ChunksCreated int
|
||||
BlobsCreated int
|
||||
StartTime time.Time
|
||||
EndTime time.Time
|
||||
FilesScanned int
|
||||
FilesSkipped int
|
||||
FilesDeleted int
|
||||
BytesScanned int64
|
||||
BytesSkipped int64
|
||||
BytesDeleted int64
|
||||
ChunksCreated int
|
||||
BlobsCreated int
|
||||
BlobsUploaded int
|
||||
BytesUploaded int64
|
||||
UploadDuration time.Duration
|
||||
StartTime time.Time
|
||||
EndTime time.Time
|
||||
}
|
||||
|
||||
// NewScanner creates a new scanner instance
|
||||
@@ -211,7 +213,6 @@ func (s *Scanner) Scan(
|
||||
s.snapshotID = snapshotID
|
||||
// Store source path for file records (used during restore)
|
||||
s.currentSourcePath = path
|
||||
s.scanCtx = ctx
|
||||
result := &ScanResult{
|
||||
StartTime: time.Now().UTC(),
|
||||
}
|
||||
@@ -219,7 +220,9 @@ func (s *Scanner) Scan(
|
||||
// Set blob handler for concurrent upload
|
||||
if s.storage != nil {
|
||||
log.Debug("Setting blob handler for storage uploads")
|
||||
s.packer.SetBlobHandler(s.handleBlobReady)
|
||||
s.packer.SetBlobHandler(func(blobWithReader *blob.WithReader) error {
|
||||
return s.handleBlobReady(ctx, blobWithReader, result)
|
||||
})
|
||||
} else {
|
||||
log.Debug("No storage configured, blobs will not be uploaded")
|
||||
}
|
||||
@@ -288,8 +291,7 @@ func (s *Scanner) Scan(
|
||||
log.Info("Phase 2/3: Skipping (no files need processing, metadata-only snapshot)")
|
||||
}
|
||||
|
||||
// Finalize result with blob statistics
|
||||
s.finalizeScanResult(ctx, result)
|
||||
result.EndTime = time.Now().UTC()
|
||||
|
||||
return result, nil
|
||||
}
|
||||
@@ -431,27 +433,6 @@ func (s *Scanner) summarizeScanPhase(
|
||||
s.ui.Completef("%s.", msg)
|
||||
}
|
||||
|
||||
// finalizeScanResult populates final blob statistics in the scan result
|
||||
// by querying the packer and database for blob/upload counts
|
||||
func (s *Scanner) finalizeScanResult(ctx context.Context, result *ScanResult) {
|
||||
blobs := s.packer.GetFinishedBlobs()
|
||||
result.BlobsCreated += len(blobs)
|
||||
|
||||
// Query database for actual blob count created during this snapshot
|
||||
// The database is authoritative, especially for concurrent blob uploads
|
||||
// We count uploads rather than all snapshot_blobs to get only NEW blobs
|
||||
if s.snapshotID != "" {
|
||||
uploadCount, err := s.repos.Uploads.GetCountBySnapshot(ctx, s.snapshotID)
|
||||
if err != nil {
|
||||
log.Warn("Failed to query upload count from database", "error", err)
|
||||
} else {
|
||||
result.BlobsCreated = int(uploadCount)
|
||||
}
|
||||
}
|
||||
|
||||
result.EndTime = time.Now().UTC()
|
||||
}
|
||||
|
||||
// loadKnownFiles loads the known files at and beneath path from the
|
||||
// database into a map for fast lookup. Every loaded file the scan does
|
||||
// not find is counted as deleted. This avoids per-file database queries
|
||||
@@ -1511,24 +1492,24 @@ func (s *Scanner) finalizeProcessPhase(ctx context.Context, result *ScanResult)
|
||||
}
|
||||
|
||||
// handleBlobReady is called by the packer when a blob is finalized
|
||||
func (s *Scanner) handleBlobReady(blobWithReader *blob.WithReader) error {
|
||||
func (s *Scanner) handleBlobReady(
|
||||
ctx context.Context, blobWithReader *blob.WithReader, result *ScanResult,
|
||||
) error {
|
||||
startTime := time.Now().UTC()
|
||||
finishedBlob := blobWithReader.FinishedBlob
|
||||
|
||||
result.BlobsCreated++
|
||||
|
||||
if s.progress != nil {
|
||||
s.progress.ReportUploadStart(finishedBlob.Hash, finishedBlob.Compressed)
|
||||
s.progress.GetStats().BlobsCreated.Add(1)
|
||||
}
|
||||
|
||||
ctx := s.scanCtx
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
}
|
||||
|
||||
blobPath := fmt.Sprintf("blobs/%s/%s/%s",
|
||||
finishedBlob.Hash[:2], finishedBlob.Hash[2:4], finishedBlob.Hash)
|
||||
|
||||
blobExists, err := s.uploadBlobIfNeeded(ctx, blobPath, blobWithReader, startTime)
|
||||
blobExists, err := s.uploadBlobIfNeeded(
|
||||
ctx, blobPath, blobWithReader, startTime, result)
|
||||
if err != nil {
|
||||
s.cleanupBlobTempFile(blobWithReader)
|
||||
|
||||
@@ -1563,6 +1544,7 @@ func (s *Scanner) uploadBlobIfNeeded(
|
||||
blobPath string,
|
||||
blobWithReader *blob.WithReader,
|
||||
startTime time.Time,
|
||||
result *ScanResult,
|
||||
) (bool, error) {
|
||||
finishedBlob := blobWithReader.FinishedBlob
|
||||
|
||||
@@ -1598,6 +1580,10 @@ func (s *Scanner) uploadBlobIfNeeded(
|
||||
uploadDuration := time.Since(startTime)
|
||||
uploadSpeedBps := float64(finishedBlob.Compressed) / uploadDuration.Seconds()
|
||||
|
||||
result.BlobsUploaded++
|
||||
result.BytesUploaded += finishedBlob.Compressed
|
||||
result.UploadDuration += uploadDuration
|
||||
|
||||
s.ui.Completef("Uploaded blob %s (%s) in %s at %s.",
|
||||
s.ui.Hex(finishedBlob.Hash),
|
||||
s.ui.Size(finishedBlob.Compressed),
|
||||
@@ -1813,9 +1799,9 @@ func (s *Scanner) processFileStreaming(
|
||||
size: chunk.Size,
|
||||
})
|
||||
|
||||
s.updateChunkStats(chunkExists, chunk.Size, result)
|
||||
|
||||
if !chunkExists {
|
||||
s.updateChunkStats(chunk.Size, result)
|
||||
|
||||
err := s.addChunkToPacker(ctx, chunk)
|
||||
if err != nil {
|
||||
// Mark as a packer error so --skip-errors cannot swallow it:
|
||||
@@ -1843,26 +1829,16 @@ func (s *Scanner) processFileStreaming(
|
||||
return nil
|
||||
}
|
||||
|
||||
// updateChunkStats updates scan result and progress stats for a processed chunk
|
||||
func (s *Scanner) updateChunkStats(
|
||||
chunkExists bool, chunkSize int64, result *ScanResult,
|
||||
) {
|
||||
if chunkExists {
|
||||
result.FilesSkipped++
|
||||
// updateChunkStats counts a chunk that was not already stored. The scan
|
||||
// result's file counts, BytesScanned and BytesSkipped are not touched
|
||||
// here: the scan phase counts each file once.
|
||||
func (s *Scanner) updateChunkStats(chunkSize int64, result *ScanResult) {
|
||||
result.ChunksCreated++
|
||||
|
||||
result.BytesSkipped += chunkSize
|
||||
if s.progress != nil {
|
||||
s.progress.GetStats().BytesSkipped.Add(chunkSize)
|
||||
}
|
||||
} else {
|
||||
result.ChunksCreated++
|
||||
result.BytesScanned += chunkSize
|
||||
|
||||
if s.progress != nil {
|
||||
s.progress.GetStats().ChunksCreated.Add(1)
|
||||
s.progress.GetStats().BytesProcessed.Add(chunkSize)
|
||||
s.progress.UpdateChunkingActivity()
|
||||
}
|
||||
if s.progress != nil {
|
||||
s.progress.GetStats().ChunksCreated.Add(1)
|
||||
s.progress.GetStats().BytesProcessed.Add(chunkSize)
|
||||
s.progress.UpdateChunkingActivity()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -154,26 +154,6 @@ func (sm *SnapshotManager) CreateSnapshotWithName(
|
||||
return snapshotID, nil
|
||||
}
|
||||
|
||||
// UpdateSnapshotStats updates the statistics for a snapshot during backup
|
||||
func (sm *SnapshotManager) UpdateSnapshotStats(
|
||||
ctx context.Context, snapshotID string, stats BackupStats,
|
||||
) error {
|
||||
err := sm.repos.WithTx(ctx, func(ctx context.Context, tx *sql.Tx) error {
|
||||
return sm.repos.Snapshots.UpdateCounts(ctx, tx, snapshotID,
|
||||
int64(stats.FilesScanned),
|
||||
int64(stats.ChunksCreated),
|
||||
int64(stats.BlobsCreated),
|
||||
stats.BytesScanned,
|
||||
stats.BytesUploaded,
|
||||
)
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("updating snapshot stats: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// UpdateSnapshotStatsExtended updates snapshot statistics with extended metrics.
|
||||
// This includes compression level, uncompressed blob size, and upload duration.
|
||||
func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
||||
@@ -185,8 +165,8 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
||||
int64(stats.FilesScanned),
|
||||
int64(stats.ChunksCreated),
|
||||
int64(stats.BlobsCreated),
|
||||
stats.BytesScanned,
|
||||
stats.BytesUploaded,
|
||||
stats.TotalSize,
|
||||
stats.BlobSize,
|
||||
)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -196,6 +176,7 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
||||
return sm.repos.Snapshots.UpdateExtendedStats(ctx, tx, snapshotID,
|
||||
stats.BlobUncompressedSize,
|
||||
stats.CompressionLevel,
|
||||
stats.BytesUploaded,
|
||||
stats.UploadDurationMs,
|
||||
)
|
||||
})
|
||||
@@ -890,7 +871,7 @@ func (sm *SnapshotManager) getFileSize(path string) int64 {
|
||||
// BackupStats contains statistics from a backup operation
|
||||
type BackupStats struct {
|
||||
FilesScanned int
|
||||
BytesScanned int64
|
||||
TotalSize int64 // Total size of all files examined
|
||||
ChunksCreated int
|
||||
BlobsCreated int
|
||||
BytesUploaded int64
|
||||
@@ -900,6 +881,7 @@ type BackupStats struct {
|
||||
type ExtendedBackupStats struct {
|
||||
BackupStats
|
||||
|
||||
BlobSize int64 // Total compressed size of all referenced blobs
|
||||
BlobUncompressedSize int64 // Total uncompressed size of all referenced blobs
|
||||
CompressionLevel int // Compression level used for this snapshot
|
||||
UploadDurationMs int64 // Total milliseconds spent uploading to S3
|
||||
|
||||
@@ -189,6 +189,11 @@ type snapshotStats struct {
|
||||
totalBytesUploaded int64
|
||||
totalBlobsUploaded int
|
||||
uploadDuration time.Duration
|
||||
|
||||
// The sizes of all blobs the snapshot references, set by
|
||||
// finalizeSnapshotMetadata once snapshot_blobs is populated.
|
||||
blobSize int64
|
||||
blobUncompressedSize int64
|
||||
}
|
||||
|
||||
// createNamedSnapshot creates a single named snapshot
|
||||
@@ -228,8 +233,6 @@ func (v *Vaultik) createNamedSnapshot(
|
||||
return err
|
||||
}
|
||||
|
||||
v.collectUploadStats(scanner, stats)
|
||||
|
||||
err = v.finalizeSnapshotMetadata(snapshotID, stats)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -314,6 +317,9 @@ func (v *Vaultik) scanAllDirectories(
|
||||
stats.totalBytesSkipped += result.BytesSkipped
|
||||
stats.totalFilesDeleted += result.FilesDeleted
|
||||
stats.totalBytesDeleted += result.BytesDeleted
|
||||
stats.totalBlobsUploaded += result.BlobsUploaded
|
||||
stats.totalBytesUploaded += result.BytesUploaded
|
||||
stats.uploadDuration += result.UploadDuration
|
||||
|
||||
log.Info("Directory scan complete",
|
||||
"path", dir,
|
||||
@@ -329,18 +335,6 @@ func (v *Vaultik) scanAllDirectories(
|
||||
return stats, nil
|
||||
}
|
||||
|
||||
// collectUploadStats gathers upload statistics from the scanner's
|
||||
// progress reporter.
|
||||
func (v *Vaultik) collectUploadStats(scanner *snapshot.Scanner, stats *snapshotStats) {
|
||||
if s := scanner.GetProgress(); s != nil {
|
||||
progressStats := s.GetStats()
|
||||
stats.totalBytesUploaded = progressStats.BytesUploaded.Load()
|
||||
stats.totalBlobsUploaded = int(progressStats.BlobsUploaded.Load())
|
||||
stats.uploadDuration = time.Duration(
|
||||
progressStats.UploadDurationMs.Load()) * time.Millisecond
|
||||
}
|
||||
}
|
||||
|
||||
// finalizeSnapshotMetadata updates stats, exports metadata, and only then
|
||||
// marks the snapshot complete. Recording completion last is deliberate: an
|
||||
// export interrupted by a crash leaves the snapshot incomplete rather than
|
||||
@@ -350,31 +344,39 @@ func (v *Vaultik) collectUploadStats(scanner *snapshot.Scanner, stats *snapshotS
|
||||
func (v *Vaultik) finalizeSnapshotMetadata(
|
||||
snapshotID string, stats *snapshotStats,
|
||||
) error {
|
||||
// snapshot_blobs must be populated before the blob sizes below, which
|
||||
// total the snapshot's blobs, and before the export, which builds the
|
||||
// manifest and the trimmed metadata database from it.
|
||||
err := v.SnapshotManager.PopulateSnapshotBlobs(v.ctx, snapshotID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("populating snapshot blobs: %w", err)
|
||||
}
|
||||
|
||||
stats.blobSize, stats.blobUncompressedSize, err =
|
||||
v.Repositories.Snapshots.GetSnapshotBlobSizes(v.ctx, snapshotID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("getting snapshot blob sizes: %w", err)
|
||||
}
|
||||
|
||||
extStats := snapshot.ExtendedBackupStats{
|
||||
BackupStats: snapshot.BackupStats{
|
||||
FilesScanned: stats.totalFiles,
|
||||
BytesScanned: stats.totalBytes,
|
||||
TotalSize: stats.totalBytes + stats.totalBytesSkipped,
|
||||
ChunksCreated: stats.totalChunks,
|
||||
BlobsCreated: stats.totalBlobs,
|
||||
BytesUploaded: stats.totalBytesUploaded,
|
||||
},
|
||||
BlobUncompressedSize: 0,
|
||||
BlobSize: stats.blobSize,
|
||||
BlobUncompressedSize: stats.blobUncompressedSize,
|
||||
CompressionLevel: v.Config.CompressionLevel,
|
||||
UploadDurationMs: stats.uploadDuration.Milliseconds(),
|
||||
}
|
||||
|
||||
err := v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
|
||||
err = v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
|
||||
if err != nil {
|
||||
return fmt.Errorf("updating snapshot stats: %w", err)
|
||||
}
|
||||
|
||||
// snapshot_blobs must be populated before the export, which builds the
|
||||
// manifest and the trimmed metadata database from it.
|
||||
err = v.SnapshotManager.PopulateSnapshotBlobs(v.ctx, snapshotID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("populating snapshot blobs: %w", err)
|
||||
}
|
||||
|
||||
err = v.SnapshotManager.ExportSnapshotMetadata(
|
||||
v.ctx, v.Config.IndexPath, snapshotID)
|
||||
if err != nil {
|
||||
@@ -409,12 +411,10 @@ func (v *Vaultik) printSnapshotSummary(
|
||||
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
|
||||
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
|
||||
|
||||
// Get total blob sizes from database
|
||||
compressedSize, uncompressedSize := v.getSnapshotBlobSizes(snapshotID)
|
||||
|
||||
var compressionRatio float64
|
||||
if uncompressedSize > 0 {
|
||||
compressionRatio = float64(compressedSize) / float64(uncompressedSize)
|
||||
if stats.blobUncompressedSize > 0 {
|
||||
compressionRatio = float64(stats.blobSize) /
|
||||
float64(stats.blobUncompressedSize)
|
||||
} else {
|
||||
compressionRatio = 1.0
|
||||
}
|
||||
@@ -442,8 +442,8 @@ func (v *Vaultik) printSnapshotSummary(
|
||||
|
||||
if stats.totalBlobsUploaded > 0 {
|
||||
v.UI.Detailf("Storage: %s compressed from %s (%.2fx ratio).",
|
||||
v.UI.Size(compressedSize),
|
||||
v.UI.Size(uncompressedSize),
|
||||
v.UI.Size(stats.blobSize),
|
||||
v.UI.Size(stats.blobUncompressedSize),
|
||||
compressionRatio)
|
||||
v.UI.Detailf("Upload: %d blobs, %s in %s (%s).",
|
||||
stats.totalBlobsUploaded,
|
||||
@@ -455,27 +455,6 @@ func (v *Vaultik) printSnapshotSummary(
|
||||
v.UI.Detailf("Snapshot create duration: %s.", v.UI.Duration(snapshotDuration))
|
||||
}
|
||||
|
||||
// getSnapshotBlobSizes returns total compressed and uncompressed blob
|
||||
// sizes for a snapshot.
|
||||
func (v *Vaultik) getSnapshotBlobSizes(snapshotID string) (int64, int64) {
|
||||
var compressed, uncompressed int64
|
||||
|
||||
blobHashes, err := v.Repositories.Snapshots.GetBlobHashes(v.ctx, snapshotID)
|
||||
if err != nil {
|
||||
return 0, 0
|
||||
}
|
||||
|
||||
for _, hash := range blobHashes {
|
||||
blob, err := v.Repositories.Blobs.GetByHash(v.ctx, hash)
|
||||
if err == nil && blob != nil {
|
||||
compressed += blob.CompressedSize
|
||||
uncompressed += blob.UncompressedSize
|
||||
}
|
||||
}
|
||||
|
||||
return compressed, uncompressed
|
||||
}
|
||||
|
||||
// SnapshotPurgeOptions contains options for the snapshot purge command.
|
||||
type SnapshotPurgeOptions struct {
|
||||
KeepLatest bool // Keep only the most recent snapshot per name
|
||||
|
||||
@@ -0,0 +1,285 @@
|
||||
package vaultik_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/spf13/afero"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"sneak.berlin/go/vaultik/internal/config"
|
||||
"sneak.berlin/go/vaultik/internal/database"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
"sneak.berlin/go/vaultik/internal/storage/faultstore"
|
||||
"sneak.berlin/go/vaultik/internal/ui"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
|
||||
// These tests cover https://git.eeqj.de/sneak/vaultik/issues/225: the
|
||||
// summary printed after a backup, and the statistics stored in the
|
||||
// snapshots table, count each file, byte and upload once, and a --cron
|
||||
// run records its uploads.
|
||||
|
||||
// summaryUploadDelay slows every blob upload, so a run's upload time is
|
||||
// at least this long per blob even on a local store.
|
||||
const summaryUploadDelay = 20 * time.Millisecond
|
||||
|
||||
// summaryEnv is a backup setup whose user-facing output is kept in out.
|
||||
type summaryEnv struct {
|
||||
v *vaultik.Vaultik
|
||||
db *database.DB
|
||||
repos *database.Repositories
|
||||
out *bytes.Buffer
|
||||
|
||||
// aPath is a.bin, whose content copy.bin repeats; aSize is its size
|
||||
// and totalSize the size of all three source files.
|
||||
aPath string
|
||||
aSize int64
|
||||
totalSize int64
|
||||
}
|
||||
|
||||
// newSummaryEnv writes src/one/a.bin, src/one/small.txt and
|
||||
// src/two/copy.bin, a copy of a.bin. Every chunk of copy.bin is therefore
|
||||
// already stored by the time the backup reaches it.
|
||||
//
|
||||
// The snapshot names "first" and "second" back up src; "split" backs up
|
||||
// src/one and src/two as two paths.
|
||||
func newSummaryEnv(t *testing.T) *summaryEnv {
|
||||
t.Helper()
|
||||
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
srcDir := filepath.Join(tempDir, "src")
|
||||
dirOne := filepath.Join(srcDir, "one")
|
||||
dirTwo := filepath.Join(srcDir, "two")
|
||||
dbPath := filepath.Join(tempDir, "index.sqlite")
|
||||
ctx := context.Background()
|
||||
|
||||
aContent := bytesPattern("a-", int(3*faultChunkSize))
|
||||
smallContent := []byte("hello vaultik")
|
||||
files := map[string][]byte{
|
||||
filepath.Join(dirOne, "a.bin"): aContent,
|
||||
filepath.Join(dirOne, "small.txt"): smallContent,
|
||||
filepath.Join(dirTwo, "copy.bin"): aContent,
|
||||
}
|
||||
|
||||
for path, content := range files {
|
||||
require.NoError(t, fs.MkdirAll(filepath.Dir(path), 0o755))
|
||||
require.NoError(t, afero.WriteFile(fs, path, content, 0o644))
|
||||
}
|
||||
|
||||
cfg := faultTestConfig()
|
||||
cfg.IndexPath = dbPath
|
||||
cfg.ChunkSize = config.Size(faultChunkSize)
|
||||
cfg.Snapshots = map[string]config.SnapshotConfig{
|
||||
"first": {Paths: []string{srcDir}},
|
||||
"second": {Paths: []string{srcDir}},
|
||||
"split": {Paths: []string{dirOne, dirTwo}},
|
||||
}
|
||||
|
||||
inner, err := storage.NewFileStorer(filepath.Join(tempDir, "remote"))
|
||||
require.NoError(t, err)
|
||||
|
||||
store := faultstore.New(inner)
|
||||
store.OnPut = func(key string) faultstore.PutAction {
|
||||
if strings.HasPrefix(key, "blobs/") {
|
||||
time.Sleep(summaryUploadDelay)
|
||||
}
|
||||
|
||||
return faultstore.PutNormal
|
||||
}
|
||||
|
||||
db, err := database.New(ctx, dbPath)
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = db.Close() })
|
||||
|
||||
repos := database.NewRepositories(db)
|
||||
out := &bytes.Buffer{}
|
||||
v := newBackupVaultik(ctx, cfg, store, repos, db, fs)
|
||||
v.UI = ui.NewWithColor(out, false)
|
||||
|
||||
return &summaryEnv{
|
||||
v: v,
|
||||
db: db,
|
||||
repos: repos,
|
||||
out: out,
|
||||
aPath: filepath.Join(dirOne, "a.bin"),
|
||||
aSize: int64(len(aContent)),
|
||||
totalSize: int64(2*len(aContent) + len(smallContent)),
|
||||
}
|
||||
}
|
||||
|
||||
// backUp runs a backup of the named snapshot and returns its output.
|
||||
func (e *summaryEnv) backUp(t *testing.T, name string, cron bool) string {
|
||||
t.Helper()
|
||||
|
||||
e.out.Reset()
|
||||
require.NoError(t, e.v.CreateSnapshot(&vaultik.SnapshotCreateOptions{
|
||||
Cron: cron,
|
||||
Snapshots: []string{name},
|
||||
}))
|
||||
|
||||
return e.out.String()
|
||||
}
|
||||
|
||||
// snapshot returns the local snapshots row of the snapshot named name.
|
||||
func (e *summaryEnv) snapshot(t *testing.T, name string) *database.Snapshot {
|
||||
t.Helper()
|
||||
|
||||
ctx := context.Background()
|
||||
snap, err := e.repos.Snapshots.GetByID(ctx,
|
||||
localSnapshotID(ctx, t, e.repos, name))
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, snap)
|
||||
|
||||
return snap
|
||||
}
|
||||
|
||||
// uploads returns how many blobs the snapshot uploaded and their
|
||||
// total size, as recorded in the uploads table.
|
||||
func (e *summaryEnv) uploads(t *testing.T, snapshotID string) (int64, int64) {
|
||||
t.Helper()
|
||||
|
||||
var count, size int64
|
||||
|
||||
err := e.db.Conn().QueryRowContext(context.Background(), `
|
||||
SELECT COUNT(*), COALESCE(SUM(size), 0)
|
||||
FROM uploads WHERE snapshot_id = ?`, snapshotID).Scan(&count, &size)
|
||||
require.NoError(t, err)
|
||||
|
||||
return count, size
|
||||
}
|
||||
|
||||
// referencedBlobSizes returns the compressed and uncompressed sizes of
|
||||
// all blobs the snapshot references.
|
||||
func (e *summaryEnv) referencedBlobSizes(
|
||||
t *testing.T, snapshotID string,
|
||||
) (int64, int64) {
|
||||
t.Helper()
|
||||
|
||||
var compressed, uncompressed int64
|
||||
|
||||
err := e.db.Conn().QueryRowContext(context.Background(), `
|
||||
SELECT COALESCE(SUM(b.compressed_size), 0),
|
||||
COALESCE(SUM(b.uncompressed_size), 0)
|
||||
FROM snapshot_blobs sb JOIN blobs b ON b.blob_hash = sb.blob_hash
|
||||
WHERE sb.snapshot_id = ?`, snapshotID).Scan(&compressed, &uncompressed)
|
||||
require.NoError(t, err)
|
||||
|
||||
return compressed, uncompressed
|
||||
}
|
||||
|
||||
// filesLine returns the summary's line of file counts.
|
||||
func filesLine(examined, backedUp, unchanged int) string {
|
||||
return fmt.Sprintf("Files: %d examined, %d backed up, %d unchanged.",
|
||||
examined, backedUp, unchanged)
|
||||
}
|
||||
|
||||
// dataLine returns the summary's line of byte counts.
|
||||
func (e *summaryEnv) dataLine(total, backedUp int64) string {
|
||||
return fmt.Sprintf("Data: %s total (%s backed up).",
|
||||
e.v.UI.Size(total), e.v.UI.Size(backedUp))
|
||||
}
|
||||
|
||||
// A first backup stores copy.bin's chunks while backing up a.bin, so
|
||||
// copy.bin's chunks are deduplicated within the run. Each file and byte
|
||||
// is still counted once.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryFirstRun(t *testing.T) {
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
summary := env.backUp(t, "first", false)
|
||||
|
||||
assert.Contains(t, summary, filesLine(3, 3, 0))
|
||||
assert.Contains(t, summary, env.dataLine(env.totalSize, env.totalSize))
|
||||
|
||||
snap := env.snapshot(t, "first")
|
||||
uploadCount, uploadBytes := env.uploads(t, snap.ID.String())
|
||||
require.Positive(t, uploadCount)
|
||||
|
||||
assert.Contains(t, summary, fmt.Sprintf("Upload: %d blobs, %s in ",
|
||||
uploadCount, env.v.UI.Size(uploadBytes)))
|
||||
|
||||
assert.Equal(t, int64(3), snap.FileCount)
|
||||
assert.Equal(t, env.totalSize, snap.TotalSize)
|
||||
assert.Equal(t, uploadCount, snap.BlobCount)
|
||||
assert.Equal(t, uploadBytes, snap.UploadBytes)
|
||||
}
|
||||
|
||||
// An incremental backup where a.bin's mtime changed but its content did
|
||||
// not: a.bin is backed up again and every one of its chunks is already
|
||||
// stored.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) {
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
env.backUp(t, "first", false)
|
||||
|
||||
later := time.Now().Add(time.Hour)
|
||||
require.NoError(t, os.Chtimes(env.aPath, later, later))
|
||||
|
||||
summary := env.backUp(t, "second", false)
|
||||
|
||||
assert.Contains(t, summary, filesLine(3, 1, 2))
|
||||
assert.Contains(t, summary, env.dataLine(env.totalSize, env.aSize))
|
||||
assert.NotContains(t, summary, "Upload:")
|
||||
|
||||
snap := env.snapshot(t, "second")
|
||||
compressed, uncompressed := env.referencedBlobSizes(t, snap.ID.String())
|
||||
require.Positive(t, compressed)
|
||||
|
||||
assert.Equal(t, env.totalSize, snap.TotalSize)
|
||||
assert.Zero(t, snap.ChunkCount)
|
||||
assert.Zero(t, snap.BlobCount)
|
||||
assert.Zero(t, snap.UploadBytes)
|
||||
assert.Equal(t, compressed, snap.BlobSize,
|
||||
"blob_size must total the blobs the snapshot references")
|
||||
assert.Equal(t, uncompressed, snap.BlobUncompressedSize)
|
||||
}
|
||||
|
||||
// Under --cron the progress reporter is off; the upload figures must
|
||||
// still reach the summary and the snapshots row. The snapshot has two
|
||||
// paths, each backed up by its own scan.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryCronRunRecordsUploads(t *testing.T) {
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
summary := env.backUp(t, "split", true)
|
||||
|
||||
snap := env.snapshot(t, "split")
|
||||
uploadCount, uploadBytes := env.uploads(t, snap.ID.String())
|
||||
require.Positive(t, uploadCount)
|
||||
|
||||
assert.Contains(t, summary, filesLine(3, 3, 0))
|
||||
assert.Contains(t, summary, env.dataLine(env.totalSize, env.totalSize))
|
||||
assert.Contains(t, summary, fmt.Sprintf("Upload: %d blobs, %s in ",
|
||||
uploadCount, env.v.UI.Size(uploadBytes)))
|
||||
|
||||
assert.Equal(t, env.totalSize, snap.TotalSize)
|
||||
assert.Equal(t, uploadCount, snap.BlobCount,
|
||||
"blob_count must count each blob once, however many paths the "+
|
||||
"snapshot has")
|
||||
assert.Equal(t, uploadBytes, snap.UploadBytes)
|
||||
assert.GreaterOrEqual(t, snap.UploadDurationMs,
|
||||
uploadCount*summaryUploadDelay.Milliseconds())
|
||||
|
||||
compressed, uncompressed := env.referencedBlobSizes(t, snap.ID.String())
|
||||
require.Positive(t, uncompressed)
|
||||
|
||||
assert.Equal(t, compressed, snap.BlobSize)
|
||||
assert.Equal(t, uncompressed, snap.BlobUncompressedSize)
|
||||
assert.InDelta(t, float64(compressed)/float64(uncompressed),
|
||||
snap.CompressionRatio, 1e-9)
|
||||
}
|
||||
Reference in New Issue
Block a user