Count each file, byte and upload once in backup statistics (closes #225)
check / check (push) Waiting to run
check / check (push) Waiting to run
The scanner added a changed file's bytes again for each new chunk and counted a file as unchanged for each chunk already stored. It now counts files and bytes once, in the scan phase, and counts its own uploads, so a --cron run, which has no progress reporter, records them. The blob count no longer adds earlier paths' blobs again. The snapshots row now stores the size of all files in total_size and the referenced blobs' sizes in blob_size, blob_uncompressed_size and compression_ratio, as docs/DATAMODEL.md says. Those sizes come from one query, and a failed query fails the snapshot. DATAMODEL.md now says chunk_count and blob_count count what the run added. Removed UpdateSnapshotStats and GetCountBySnapshot, which nothing calls any more. Model: opus-5-5
This commit was merged in pull request #254.
This commit is contained in:
+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,7 +131,9 @@ 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
|
||||
@@ -144,6 +143,9 @@ type ScanResult struct {
|
||||
BytesDeleted int64
|
||||
ChunksCreated int
|
||||
BlobsCreated int
|
||||
BlobsUploaded int
|
||||
BytesUploaded int64
|
||||
UploadDuration time.Duration
|
||||
StartTime time.Time
|
||||
EndTime time.Time
|
||||
}
|
||||
@@ -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,20 +1829,11 @@ 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++
|
||||
|
||||
result.BytesSkipped += chunkSize
|
||||
if s.progress != nil {
|
||||
s.progress.GetStats().BytesSkipped.Add(chunkSize)
|
||||
}
|
||||
} else {
|
||||
// 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.BytesScanned += chunkSize
|
||||
|
||||
if s.progress != nil {
|
||||
s.progress.GetStats().ChunksCreated.Add(1)
|
||||
@@ -1864,7 +1841,6 @@ func (s *Scanner) updateChunkStats(
|
||||
s.progress.UpdateChunkingActivity()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// addChunkToPacker adds a chunk to the blob packer, finalizing the current
|
||||
// blob if needed
|
||||
|
||||
@@ -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