diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 309186f..0541199 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -335,10 +335,10 @@ CreateSnapshot(opts) │ │ │ └─► Accumulate statistics │ - ├─► SnapshotManager.UpdateSnapshotStatsExtended() - │ ├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs │ + ├─► SnapshotManager.UpdateSnapshotStatsExtended() + │ ├─► SnapshotManager.ExportSnapshotMetadata() │ │ │ ├─► Copy database to temp file diff --git a/TODO.md b/TODO.md index b2e6cb7..ec17c26 100644 --- a/TODO.md +++ b/TODO.md @@ -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 diff --git a/docs/DATAMODEL.md b/docs/DATAMODEL.md index 4fe62a5..a1967a3 100644 --- a/docs/DATAMODEL.md +++ b/docs/DATAMODEL.md @@ -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 diff --git a/internal/database/models.go b/internal/database/models.go index ed27a8c..9c23c59 100644 --- a/internal/database/models.go +++ b/internal/database/models.go @@ -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 diff --git a/internal/database/snapshots.go b/internal/database/snapshots.go index 4479985..09d3b3a 100644 --- a/internal/database/snapshots.go +++ b/internal/database/snapshots.go @@ -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( diff --git a/internal/database/snapshots_test.go b/internal/database/snapshots_test.go index 9d07c8e..5ddbd14 100644 --- a/internal/database/snapshots_test.go +++ b/internal/database/snapshots_test.go @@ -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() diff --git a/internal/database/uploads.go b/internal/database/uploads.go index 312373e..7939089 100644 --- a/internal/database/uploads.go +++ b/internal/database/uploads.go @@ -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 -} diff --git a/internal/snapshot/progress.go b/internal/snapshot/progress.go index 5cd151a..9aba0f5 100644 --- a/internal/snapshot/progress.go +++ b/internal/snapshot/progress.go @@ -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 diff --git a/internal/snapshot/scanner.go b/internal/snapshot/scanner.go index e74921a..6676f6e 100644 --- a/internal/snapshot/scanner.go +++ b/internal/snapshot/scanner.go @@ -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() } } diff --git a/internal/snapshot/snapshot.go b/internal/snapshot/snapshot.go index 82a12a7..2848864 100644 --- a/internal/snapshot/snapshot.go +++ b/internal/snapshot/snapshot.go @@ -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 diff --git a/internal/vaultik/snapshot.go b/internal/vaultik/snapshot.go index 7c948f7..82a116e 100644 --- a/internal/vaultik/snapshot.go +++ b/internal/vaultik/snapshot.go @@ -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 diff --git a/internal/vaultik/snapshot_summary_test.go b/internal/vaultik/snapshot_summary_test.go new file mode 100644 index 0000000..a3562f3 --- /dev/null +++ b/internal/vaultik/snapshot_summary_test.go @@ -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) +}