Count each file, byte and upload once in backup statistics #254
+2
-2
@@ -335,10 +335,10 @@ CreateSnapshot(opts)
|
|||||||
│ │
|
│ │
|
||||||
│ └─► Accumulate statistics
|
│ └─► Accumulate statistics
|
||||||
│
|
│
|
||||||
├─► SnapshotManager.UpdateSnapshotStatsExtended()
|
|
||||||
│
|
|
||||||
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
|
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
|
||||||
│
|
│
|
||||||
|
├─► SnapshotManager.UpdateSnapshotStatsExtended()
|
||||||
|
│
|
||||||
├─► SnapshotManager.ExportSnapshotMetadata()
|
├─► SnapshotManager.ExportSnapshotMetadata()
|
||||||
│ │
|
│ │
|
||||||
│ ├─► Copy database to temp file
|
│ ├─► Copy database to temp file
|
||||||
|
|||||||
@@ -22,6 +22,19 @@ the tag exists and is exercised; what is left is merging `next` to
|
|||||||
|
|
||||||
# Completed Steps
|
# 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
|
- 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
|
rules ([issue #224](https://git.eeqj.de/sneak/vaultik/issues/224)). The
|
||||||
startup banner went to stdout, so a `completion` script or a
|
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
|
- `started_at` (INTEGER) - Start timestamp
|
||||||
- `completed_at` (INTEGER) - Completion timestamp (NULL if in progress)
|
- `completed_at` (INTEGER) - Completion timestamp (NULL if in progress)
|
||||||
- `file_count` (INTEGER) - Number of files in snapshot
|
- `file_count` (INTEGER) - Number of files in snapshot
|
||||||
- `chunk_count` (INTEGER) - Number of unique chunks
|
- `chunk_count` (INTEGER) - Number of chunks this snapshot stored that were not stored before
|
||||||
- `blob_count` (INTEGER) - Number of blobs referenced
|
- `blob_count` (INTEGER) - Number of blobs this snapshot created
|
||||||
- `total_size` (INTEGER) - Total size of all files
|
- `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
|
- `blob_uncompressed_size` (INTEGER) - Total uncompressed size of all referenced blobs
|
||||||
- `compression_ratio` (REAL) - Compression ratio achieved
|
- `compression_ratio` (REAL) - Compression ratio achieved
|
||||||
- `compression_level` (INTEGER) - Compression level used for this snapshot
|
- `compression_level` (INTEGER) - Compression level used for this snapshot
|
||||||
|
|||||||
@@ -99,8 +99,8 @@ type Snapshot struct {
|
|||||||
StartedAt time.Time
|
StartedAt time.Time
|
||||||
CompletedAt *time.Time // nil if still in progress
|
CompletedAt *time.Time // nil if still in progress
|
||||||
FileCount int64
|
FileCount int64
|
||||||
ChunkCount int64
|
ChunkCount int64 // Chunks this snapshot stored that were not stored before
|
||||||
BlobCount int64
|
BlobCount int64 // Blobs this snapshot created
|
||||||
TotalSize int64 // Total size of all referenced files
|
TotalSize int64 // Total size of all referenced files
|
||||||
|
|
||||||
// BlobSize is the total size of all referenced blobs (compressed and
|
// BlobSize is the total size of all referenced blobs (compressed and
|
||||||
|
|||||||
@@ -127,6 +127,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
|||||||
snapshotID string,
|
snapshotID string,
|
||||||
blobUncompressedSize int64,
|
blobUncompressedSize int64,
|
||||||
compressionLevel int,
|
compressionLevel int,
|
||||||
|
uploadBytes int64,
|
||||||
uploadDurationMs int64,
|
uploadDurationMs int64,
|
||||||
) error {
|
) error {
|
||||||
compressionRatio, err := r.extendedCompressionRatio(
|
compressionRatio, err := r.extendedCompressionRatio(
|
||||||
@@ -141,7 +142,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
|||||||
SET blob_uncompressed_size = ?,
|
SET blob_uncompressed_size = ?,
|
||||||
compression_ratio = ?,
|
compression_ratio = ?,
|
||||||
compression_level = ?,
|
compression_level = ?,
|
||||||
upload_bytes = blob_size,
|
upload_bytes = ?,
|
||||||
upload_duration_ms = ?
|
upload_duration_ms = ?
|
||||||
WHERE id = ?
|
WHERE id = ?
|
||||||
`
|
`
|
||||||
@@ -149,11 +150,11 @@ func (r *SnapshotRepository) UpdateExtendedStats(
|
|||||||
if tx != nil {
|
if tx != nil {
|
||||||
_, err = tx.ExecContext(ctx, query,
|
_, err = tx.ExecContext(ctx, query,
|
||||||
blobUncompressedSize, compressionRatio, compressionLevel,
|
blobUncompressedSize, compressionRatio, compressionLevel,
|
||||||
uploadDurationMs, snapshotID)
|
uploadBytes, uploadDurationMs, snapshotID)
|
||||||
} else {
|
} else {
|
||||||
_, err = r.db.ExecWithLog(ctx, query,
|
_, err = r.db.ExecWithLog(ctx, query,
|
||||||
blobUncompressedSize, compressionRatio, compressionLevel,
|
blobUncompressedSize, compressionRatio, compressionLevel,
|
||||||
uploadDurationMs, snapshotID)
|
uploadBytes, uploadDurationMs, snapshotID)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -543,6 +544,30 @@ func (r *SnapshotRepository) GetSnapshotTotalCompressedSize(
|
|||||||
return totalSize, nil
|
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
|
// GetSnapshotUncompressedChunkSize returns the sum of plaintext sizes of all unique
|
||||||
// chunks referenced by a snapshot (via snapshot_files → file_chunks → chunks).
|
// chunks referenced by a snapshot (via snapshot_files → file_chunks → chunks).
|
||||||
func (r *SnapshotRepository) GetSnapshotUncompressedChunkSize(
|
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) {
|
func TestSnapshotRepositoryListRecent(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
|
|||||||
@@ -158,19 +158,3 @@ type UploadStats struct {
|
|||||||
MinDurationMs int64
|
MinDurationMs int64
|
||||||
MaxDurationMs 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
|
BlobsCreated atomic.Int64
|
||||||
BlobsUploaded atomic.Int64
|
BlobsUploaded atomic.Int64
|
||||||
BytesUploaded atomic.Int64
|
BytesUploaded atomic.Int64
|
||||||
UploadDurationMs atomic.Int64 // Total milliseconds spent uploading
|
|
||||||
CurrentFile atomic.Value // stores string
|
CurrentFile atomic.Value // stores string
|
||||||
TotalSize atomic.Int64 // Total size to process (set after scan phase)
|
TotalSize atomic.Int64 // Total size to process (set after scan phase)
|
||||||
TotalFiles atomic.Int64 // Total files to process in phase 2
|
TotalFiles atomic.Int64 // Total files to process in phase 2
|
||||||
@@ -231,9 +230,6 @@ func (pr *ProgressReporter) ReportUploadComplete(
|
|||||||
// Clear current upload
|
// Clear current upload
|
||||||
pr.stats.CurrentUpload.Store((*UploadInfo)(nil))
|
pr.stats.CurrentUpload.Store((*UploadInfo)(nil))
|
||||||
|
|
||||||
// Add to total upload duration
|
|
||||||
pr.stats.UploadDurationMs.Add(duration.Milliseconds())
|
|
||||||
|
|
||||||
// Calculate speed
|
// Calculate speed
|
||||||
if duration < time.Millisecond {
|
if duration < time.Millisecond {
|
||||||
duration = time.Millisecond
|
duration = time.Millisecond
|
||||||
|
|||||||
@@ -92,9 +92,6 @@ type Scanner struct {
|
|||||||
|
|
||||||
// Mutex for coordinating blob creation
|
// Mutex for coordinating blob creation
|
||||||
packerMu sync.Mutex // Blocks chunk production during 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
|
// Periodic status output intervals and thresholds for the scan and
|
||||||
@@ -134,7 +131,9 @@ type ScannerConfig struct {
|
|||||||
SkipErrors bool
|
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 {
|
type ScanResult struct {
|
||||||
FilesScanned int
|
FilesScanned int
|
||||||
FilesSkipped int
|
FilesSkipped int
|
||||||
@@ -144,6 +143,9 @@ type ScanResult struct {
|
|||||||
BytesDeleted int64
|
BytesDeleted int64
|
||||||
ChunksCreated int
|
ChunksCreated int
|
||||||
BlobsCreated int
|
BlobsCreated int
|
||||||
|
BlobsUploaded int
|
||||||
|
BytesUploaded int64
|
||||||
|
UploadDuration time.Duration
|
||||||
StartTime time.Time
|
StartTime time.Time
|
||||||
EndTime time.Time
|
EndTime time.Time
|
||||||
}
|
}
|
||||||
@@ -211,7 +213,6 @@ func (s *Scanner) Scan(
|
|||||||
s.snapshotID = snapshotID
|
s.snapshotID = snapshotID
|
||||||
// Store source path for file records (used during restore)
|
// Store source path for file records (used during restore)
|
||||||
s.currentSourcePath = path
|
s.currentSourcePath = path
|
||||||
s.scanCtx = ctx
|
|
||||||
result := &ScanResult{
|
result := &ScanResult{
|
||||||
StartTime: time.Now().UTC(),
|
StartTime: time.Now().UTC(),
|
||||||
}
|
}
|
||||||
@@ -219,7 +220,9 @@ func (s *Scanner) Scan(
|
|||||||
// Set blob handler for concurrent upload
|
// Set blob handler for concurrent upload
|
||||||
if s.storage != nil {
|
if s.storage != nil {
|
||||||
log.Debug("Setting blob handler for storage uploads")
|
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 {
|
} else {
|
||||||
log.Debug("No storage configured, blobs will not be uploaded")
|
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)")
|
log.Info("Phase 2/3: Skipping (no files need processing, metadata-only snapshot)")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Finalize result with blob statistics
|
result.EndTime = time.Now().UTC()
|
||||||
s.finalizeScanResult(ctx, result)
|
|
||||||
|
|
||||||
return result, nil
|
return result, nil
|
||||||
}
|
}
|
||||||
@@ -431,27 +433,6 @@ func (s *Scanner) summarizeScanPhase(
|
|||||||
s.ui.Completef("%s.", msg)
|
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
|
// loadKnownFiles loads the known files at and beneath path from the
|
||||||
// database into a map for fast lookup. Every loaded file the scan does
|
// 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
|
// 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
|
// 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()
|
startTime := time.Now().UTC()
|
||||||
finishedBlob := blobWithReader.FinishedBlob
|
finishedBlob := blobWithReader.FinishedBlob
|
||||||
|
|
||||||
|
result.BlobsCreated++
|
||||||
|
|
||||||
if s.progress != nil {
|
if s.progress != nil {
|
||||||
s.progress.ReportUploadStart(finishedBlob.Hash, finishedBlob.Compressed)
|
s.progress.ReportUploadStart(finishedBlob.Hash, finishedBlob.Compressed)
|
||||||
s.progress.GetStats().BlobsCreated.Add(1)
|
s.progress.GetStats().BlobsCreated.Add(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
ctx := s.scanCtx
|
|
||||||
if ctx == nil {
|
|
||||||
ctx = context.Background()
|
|
||||||
}
|
|
||||||
|
|
||||||
blobPath := fmt.Sprintf("blobs/%s/%s/%s",
|
blobPath := fmt.Sprintf("blobs/%s/%s/%s",
|
||||||
finishedBlob.Hash[:2], finishedBlob.Hash[2:4], finishedBlob.Hash)
|
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 {
|
if err != nil {
|
||||||
s.cleanupBlobTempFile(blobWithReader)
|
s.cleanupBlobTempFile(blobWithReader)
|
||||||
|
|
||||||
@@ -1563,6 +1544,7 @@ func (s *Scanner) uploadBlobIfNeeded(
|
|||||||
blobPath string,
|
blobPath string,
|
||||||
blobWithReader *blob.WithReader,
|
blobWithReader *blob.WithReader,
|
||||||
startTime time.Time,
|
startTime time.Time,
|
||||||
|
result *ScanResult,
|
||||||
) (bool, error) {
|
) (bool, error) {
|
||||||
finishedBlob := blobWithReader.FinishedBlob
|
finishedBlob := blobWithReader.FinishedBlob
|
||||||
|
|
||||||
@@ -1598,6 +1580,10 @@ func (s *Scanner) uploadBlobIfNeeded(
|
|||||||
uploadDuration := time.Since(startTime)
|
uploadDuration := time.Since(startTime)
|
||||||
uploadSpeedBps := float64(finishedBlob.Compressed) / uploadDuration.Seconds()
|
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.Completef("Uploaded blob %s (%s) in %s at %s.",
|
||||||
s.ui.Hex(finishedBlob.Hash),
|
s.ui.Hex(finishedBlob.Hash),
|
||||||
s.ui.Size(finishedBlob.Compressed),
|
s.ui.Size(finishedBlob.Compressed),
|
||||||
@@ -1813,9 +1799,9 @@ func (s *Scanner) processFileStreaming(
|
|||||||
size: chunk.Size,
|
size: chunk.Size,
|
||||||
})
|
})
|
||||||
|
|
||||||
s.updateChunkStats(chunkExists, chunk.Size, result)
|
|
||||||
|
|
||||||
if !chunkExists {
|
if !chunkExists {
|
||||||
|
s.updateChunkStats(chunk.Size, result)
|
||||||
|
|
||||||
err := s.addChunkToPacker(ctx, chunk)
|
err := s.addChunkToPacker(ctx, chunk)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Mark as a packer error so --skip-errors cannot swallow it:
|
// Mark as a packer error so --skip-errors cannot swallow it:
|
||||||
@@ -1843,27 +1829,17 @@ func (s *Scanner) processFileStreaming(
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// updateChunkStats updates scan result and progress stats for a processed chunk
|
// updateChunkStats counts a chunk that was not already stored. The scan
|
||||||
func (s *Scanner) updateChunkStats(
|
// result's file counts, BytesScanned and BytesSkipped are not touched
|
||||||
chunkExists bool, chunkSize int64, result *ScanResult,
|
// here: the scan phase counts each file once.
|
||||||
) {
|
func (s *Scanner) updateChunkStats(chunkSize int64, result *ScanResult) {
|
||||||
if chunkExists {
|
|
||||||
result.FilesSkipped++
|
|
||||||
|
|
||||||
result.BytesSkipped += chunkSize
|
|
||||||
if s.progress != nil {
|
|
||||||
s.progress.GetStats().BytesSkipped.Add(chunkSize)
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
result.ChunksCreated++
|
result.ChunksCreated++
|
||||||
result.BytesScanned += chunkSize
|
|
||||||
|
|
||||||
if s.progress != nil {
|
if s.progress != nil {
|
||||||
s.progress.GetStats().ChunksCreated.Add(1)
|
s.progress.GetStats().ChunksCreated.Add(1)
|
||||||
s.progress.GetStats().BytesProcessed.Add(chunkSize)
|
s.progress.GetStats().BytesProcessed.Add(chunkSize)
|
||||||
s.progress.UpdateChunkingActivity()
|
s.progress.UpdateChunkingActivity()
|
||||||
}
|
}
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// addChunkToPacker adds a chunk to the blob packer, finalizing the current
|
// addChunkToPacker adds a chunk to the blob packer, finalizing the current
|
||||||
|
|||||||
@@ -154,26 +154,6 @@ func (sm *SnapshotManager) CreateSnapshotWithName(
|
|||||||
return snapshotID, nil
|
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.
|
// UpdateSnapshotStatsExtended updates snapshot statistics with extended metrics.
|
||||||
// This includes compression level, uncompressed blob size, and upload duration.
|
// This includes compression level, uncompressed blob size, and upload duration.
|
||||||
func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
||||||
@@ -185,8 +165,8 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
|||||||
int64(stats.FilesScanned),
|
int64(stats.FilesScanned),
|
||||||
int64(stats.ChunksCreated),
|
int64(stats.ChunksCreated),
|
||||||
int64(stats.BlobsCreated),
|
int64(stats.BlobsCreated),
|
||||||
stats.BytesScanned,
|
stats.TotalSize,
|
||||||
stats.BytesUploaded,
|
stats.BlobSize,
|
||||||
)
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -196,6 +176,7 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
|
|||||||
return sm.repos.Snapshots.UpdateExtendedStats(ctx, tx, snapshotID,
|
return sm.repos.Snapshots.UpdateExtendedStats(ctx, tx, snapshotID,
|
||||||
stats.BlobUncompressedSize,
|
stats.BlobUncompressedSize,
|
||||||
stats.CompressionLevel,
|
stats.CompressionLevel,
|
||||||
|
stats.BytesUploaded,
|
||||||
stats.UploadDurationMs,
|
stats.UploadDurationMs,
|
||||||
)
|
)
|
||||||
})
|
})
|
||||||
@@ -890,7 +871,7 @@ func (sm *SnapshotManager) getFileSize(path string) int64 {
|
|||||||
// BackupStats contains statistics from a backup operation
|
// BackupStats contains statistics from a backup operation
|
||||||
type BackupStats struct {
|
type BackupStats struct {
|
||||||
FilesScanned int
|
FilesScanned int
|
||||||
BytesScanned int64
|
TotalSize int64 // Total size of all files examined
|
||||||
ChunksCreated int
|
ChunksCreated int
|
||||||
BlobsCreated int
|
BlobsCreated int
|
||||||
BytesUploaded int64
|
BytesUploaded int64
|
||||||
@@ -900,6 +881,7 @@ type BackupStats struct {
|
|||||||
type ExtendedBackupStats struct {
|
type ExtendedBackupStats struct {
|
||||||
BackupStats
|
BackupStats
|
||||||
|
|
||||||
|
BlobSize int64 // Total compressed size of all referenced blobs
|
||||||
BlobUncompressedSize int64 // Total uncompressed size of all referenced blobs
|
BlobUncompressedSize int64 // Total uncompressed size of all referenced blobs
|
||||||
CompressionLevel int // Compression level used for this snapshot
|
CompressionLevel int // Compression level used for this snapshot
|
||||||
UploadDurationMs int64 // Total milliseconds spent uploading to S3
|
UploadDurationMs int64 // Total milliseconds spent uploading to S3
|
||||||
|
|||||||
@@ -189,6 +189,11 @@ type snapshotStats struct {
|
|||||||
totalBytesUploaded int64
|
totalBytesUploaded int64
|
||||||
totalBlobsUploaded int
|
totalBlobsUploaded int
|
||||||
uploadDuration time.Duration
|
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
|
// createNamedSnapshot creates a single named snapshot
|
||||||
@@ -228,8 +233,6 @@ func (v *Vaultik) createNamedSnapshot(
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
v.collectUploadStats(scanner, stats)
|
|
||||||
|
|
||||||
err = v.finalizeSnapshotMetadata(snapshotID, stats)
|
err = v.finalizeSnapshotMetadata(snapshotID, stats)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -314,6 +317,9 @@ func (v *Vaultik) scanAllDirectories(
|
|||||||
stats.totalBytesSkipped += result.BytesSkipped
|
stats.totalBytesSkipped += result.BytesSkipped
|
||||||
stats.totalFilesDeleted += result.FilesDeleted
|
stats.totalFilesDeleted += result.FilesDeleted
|
||||||
stats.totalBytesDeleted += result.BytesDeleted
|
stats.totalBytesDeleted += result.BytesDeleted
|
||||||
|
stats.totalBlobsUploaded += result.BlobsUploaded
|
||||||
|
stats.totalBytesUploaded += result.BytesUploaded
|
||||||
|
stats.uploadDuration += result.UploadDuration
|
||||||
|
|
||||||
log.Info("Directory scan complete",
|
log.Info("Directory scan complete",
|
||||||
"path", dir,
|
"path", dir,
|
||||||
@@ -329,18 +335,6 @@ func (v *Vaultik) scanAllDirectories(
|
|||||||
return stats, nil
|
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
|
// finalizeSnapshotMetadata updates stats, exports metadata, and only then
|
||||||
// marks the snapshot complete. Recording completion last is deliberate: an
|
// marks the snapshot complete. Recording completion last is deliberate: an
|
||||||
// export interrupted by a crash leaves the snapshot incomplete rather than
|
// 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(
|
func (v *Vaultik) finalizeSnapshotMetadata(
|
||||||
snapshotID string, stats *snapshotStats,
|
snapshotID string, stats *snapshotStats,
|
||||||
) error {
|
) 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{
|
extStats := snapshot.ExtendedBackupStats{
|
||||||
BackupStats: snapshot.BackupStats{
|
BackupStats: snapshot.BackupStats{
|
||||||
FilesScanned: stats.totalFiles,
|
FilesScanned: stats.totalFiles,
|
||||||
BytesScanned: stats.totalBytes,
|
TotalSize: stats.totalBytes + stats.totalBytesSkipped,
|
||||||
ChunksCreated: stats.totalChunks,
|
ChunksCreated: stats.totalChunks,
|
||||||
BlobsCreated: stats.totalBlobs,
|
BlobsCreated: stats.totalBlobs,
|
||||||
BytesUploaded: stats.totalBytesUploaded,
|
BytesUploaded: stats.totalBytesUploaded,
|
||||||
},
|
},
|
||||||
BlobUncompressedSize: 0,
|
BlobSize: stats.blobSize,
|
||||||
|
BlobUncompressedSize: stats.blobUncompressedSize,
|
||||||
CompressionLevel: v.Config.CompressionLevel,
|
CompressionLevel: v.Config.CompressionLevel,
|
||||||
UploadDurationMs: stats.uploadDuration.Milliseconds(),
|
UploadDurationMs: stats.uploadDuration.Milliseconds(),
|
||||||
}
|
}
|
||||||
|
|
||||||
err := v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
|
err = v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("updating snapshot stats: %w", err)
|
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(
|
err = v.SnapshotManager.ExportSnapshotMetadata(
|
||||||
v.ctx, v.Config.IndexPath, snapshotID)
|
v.ctx, v.Config.IndexPath, snapshotID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -409,12 +411,10 @@ func (v *Vaultik) printSnapshotSummary(
|
|||||||
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
|
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
|
||||||
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
|
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
|
||||||
|
|
||||||
// Get total blob sizes from database
|
|
||||||
compressedSize, uncompressedSize := v.getSnapshotBlobSizes(snapshotID)
|
|
||||||
|
|
||||||
var compressionRatio float64
|
var compressionRatio float64
|
||||||
if uncompressedSize > 0 {
|
if stats.blobUncompressedSize > 0 {
|
||||||
compressionRatio = float64(compressedSize) / float64(uncompressedSize)
|
compressionRatio = float64(stats.blobSize) /
|
||||||
|
float64(stats.blobUncompressedSize)
|
||||||
} else {
|
} else {
|
||||||
compressionRatio = 1.0
|
compressionRatio = 1.0
|
||||||
}
|
}
|
||||||
@@ -442,8 +442,8 @@ func (v *Vaultik) printSnapshotSummary(
|
|||||||
|
|
||||||
if stats.totalBlobsUploaded > 0 {
|
if stats.totalBlobsUploaded > 0 {
|
||||||
v.UI.Detailf("Storage: %s compressed from %s (%.2fx ratio).",
|
v.UI.Detailf("Storage: %s compressed from %s (%.2fx ratio).",
|
||||||
v.UI.Size(compressedSize),
|
v.UI.Size(stats.blobSize),
|
||||||
v.UI.Size(uncompressedSize),
|
v.UI.Size(stats.blobUncompressedSize),
|
||||||
compressionRatio)
|
compressionRatio)
|
||||||
v.UI.Detailf("Upload: %d blobs, %s in %s (%s).",
|
v.UI.Detailf("Upload: %d blobs, %s in %s (%s).",
|
||||||
stats.totalBlobsUploaded,
|
stats.totalBlobsUploaded,
|
||||||
@@ -455,27 +455,6 @@ func (v *Vaultik) printSnapshotSummary(
|
|||||||
v.UI.Detailf("Snapshot create duration: %s.", v.UI.Duration(snapshotDuration))
|
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.
|
// SnapshotPurgeOptions contains options for the snapshot purge command.
|
||||||
type SnapshotPurgeOptions struct {
|
type SnapshotPurgeOptions struct {
|
||||||
KeepLatest bool // Keep only the most recent snapshot per name
|
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