Compare commits

..
1 Commits
Author SHA1 Message Date
sneak 1e9644edbd Write only the document to stdout from snapshot remove --json (closes #251)
check / check (push) Canceled after 0s
When the destination store cannot be reached, `snapshot remove` still
removes the snapshot from the local index and warns. The warning went
through the UI, which writes to stdout, so under `--json` it landed
ahead of the document and `| jq` failed on a command that exited 0.
Under `--json` the UI warning is now skipped. The logger's warning,
which goes to stderr in every mode, now also says to run `vaultik
prune` once the destination store is reachable.

The new test runs the command through `Entry` against a missing
destination directory, decodes stdout as exactly one JSON document, and
finds the warning and the `vaultik prune` follow-up on stderr.

Model: opus-5-5
2026-10-07 00:27:59 +00:00
15 changed files with 192 additions and 501 deletions
+2 -2
View File
@@ -335,10 +335,10 @@ CreateSnapshot(opts)
│ │ │ │
│ └─► Accumulate statistics │ └─► Accumulate statistics
│ │
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
│
├─► SnapshotManager.UpdateSnapshotStatsExtended() ├─► SnapshotManager.UpdateSnapshotStatsExtended()
│ │
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
│
├─► SnapshotManager.ExportSnapshotMetadata() ├─► SnapshotManager.ExportSnapshotMetadata()
│ │ │ │
│ ├─► Copy database to temp file │ ├─► Copy database to temp file
+2 -3
View File
@@ -352,9 +352,8 @@ may hold snapshots this host doesn't know about), which is what
prune` invocation to run as a follow-up. Local row cleanup (files, prune` invocation to run as a follow-up. Local row cleanup (files,
chunks, blobs the snapshot was the last referrer for) runs chunks, blobs the snapshot was the last referrer for) runs
automatically. If the destination store is unreachable, the local-DB automatically. If the destination store is unreachable, the local-DB
removal still completes and a warning is emitted; run `vaultik snapshot removal still completes and a warning is emitted; rerun `vaultik prune`
remove <snapshot-id>` again once the store is reachable to remove the once the store is reachable to finish remote cleanup. To wipe everything
snapshot's metadata from it (`vaultik prune` does not). To wipe everything
on the destination in one go, use `vaultik remote nuke --force`. on the destination in one go, use `vaultik remote nuke --force`.
* `--local-only`: Skip remote cleanup; only touch the local index * `--local-only`: Skip remote cleanup; only touch the local index
* `--dry-run`: Show what would be deleted without deleting * `--dry-run`: Show what would be deleted without deleting
+2 -17
View File
@@ -28,23 +28,8 @@ the tag exists and is exercised; what is left is merging `next` to
warning that the snapshot's metadata was left on the destination store warning that the snapshot's metadata was left on the destination store
went to stdout ahead of the document, breaking `| jq` on a command that went to stdout ahead of the document, breaking `| jq` on a command that
exited 0. Under `--json` the warning now reaches stderr only, through exited 0. Under `--json` the warning now reaches stderr only, through
the logger. The warning, the README and the command's help said the logger, whose record now also says to run `vaultik prune` once the
`vaultik prune` would finish the cleanup, but `prune` never removes destination store is reachable.
snapshot metadata; they now say to run `vaultik snapshot remove` for the
snapshot again once the destination store is reachable.
- 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
+3 -3
View File
@@ -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 chunks this snapshot stored that were not stored before - `chunk_count` (INTEGER) - Number of unique chunks
- `blob_count` (INTEGER) - Number of blobs this snapshot created - `blob_count` (INTEGER) - Number of blobs referenced
- `total_size` (INTEGER) - Total size of all files - `total_size` (INTEGER) - Total size of all files
- `blob_size` (INTEGER) - Total compressed size of all referenced blobs - `blob_size` (INTEGER) - Total size of all blobs (compressed)
- `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
+8 -10
View File
@@ -88,26 +88,24 @@ func TestEntryJSONFailureIsReportedOnStderr(t *testing.T) {
} }
// TestEntrySnapshotRemoveJSONWarningIsOnStderr runs `snapshot remove // TestEntrySnapshotRemoveJSONWarningIsOnStderr runs `snapshot remove
// --json` on a snapshot in the local index, against a destination // --json` against a destination directory that does not exist. The
// directory that does not exist. The command removes the snapshot from // command still exits 0 after removing the snapshot from the local
// the local index and still exits 0. Its stdout must hold the document // index. Its stdout must hold the document alone, with the warning
// alone, with the warning about the destination store, and the command // about the destination store, and the `vaultik prune` follow-up, on
// to run again once it is reachable, on stderr. // stderr.
// //
//nolint:paralleltest // replaces os.Args, os.Stdout, os.Stderr and the xdg globals //nolint:paralleltest // replaces os.Args, os.Stdout, os.Stderr and the xdg globals
func TestEntrySnapshotRemoveJSONWarningIsOnStderr(t *testing.T) { func TestEntrySnapshotRemoveJSONWarningIsOnStderr(t *testing.T) {
configPath, indexPath := writeMissingDestinationConfig(t) configPath, _ := writeMissingDestinationConfig(t)
seedStaleSnapshotRecord(t, indexPath)
code, stdout, stderr := runEntry(t, flagConfig, configPath, code, stdout, stderr := runEntry(t, flagConfig, configPath,
cmdSnapshot, cmdRemove, stalePruneSnapshotID, flagJSON) cmdSnapshot, cmdRemove, someSnapshotID, flagJSON)
require.Equal(t, 0, code) require.Equal(t, 0, code)
requireExactlyOneJSONDocument(t, stdout) requireExactlyOneJSONDocument(t, stdout)
assert.Contains(t, stderr, assert.Contains(t, stderr,
"Could not remove snapshot metadata from remote storage") "Could not remove snapshot metadata from remote storage")
assert.Contains(t, stderr, assert.Contains(t, stderr, "run 'vaultik prune'")
"run 'vaultik snapshot remove "+stalePruneSnapshotID+"' again")
} }
// writeMissingDestinationConfig builds a config whose destination // writeMissingDestinationConfig builds a config whose destination
+2 -3
View File
@@ -258,9 +258,8 @@ Use --local-only to skip the remote half (e.g. when you want to forget a
snapshot locally without touching the destination store). snapshot locally without touching the destination store).
If the remote is unreachable, the local-database removal still completes If the remote is unreachable, the local-database removal still completes
and a warning is emitted; run 'vaultik snapshot remove <snapshot-id>' again and a warning is emitted; rerun 'vaultik prune' once the destination store
once the destination store is reachable to remove the snapshot's metadata is reachable to finish remote cleanup.
from it ('vaultik prune' does not).
To wipe the entire destination store and start over, use 'vaultik remote To wipe the entire destination store and start over, use 'vaultik remote
nuke --force' — it is the single supported entry point for that.`, nuke --force' — it is the single supported entry point for that.`,
+2 -2
View File
@@ -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 // Chunks this snapshot stored that were not stored before ChunkCount int64
BlobCount int64 // Blobs this snapshot created BlobCount int64
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
+3 -28
View File
@@ -127,7 +127,6 @@ 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(
@@ -142,7 +141,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
SET blob_uncompressed_size = ?, SET blob_uncompressed_size = ?,
compression_ratio = ?, compression_ratio = ?,
compression_level = ?, compression_level = ?,
upload_bytes = ?, upload_bytes = blob_size,
upload_duration_ms = ? upload_duration_ms = ?
WHERE id = ? WHERE id = ?
` `
@@ -150,11 +149,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,
uploadBytes, uploadDurationMs, snapshotID) uploadDurationMs, snapshotID)
} else { } else {
_, err = r.db.ExecWithLog(ctx, query, _, err = r.db.ExecWithLog(ctx, query,
blobUncompressedSize, compressionRatio, compressionLevel, blobUncompressedSize, compressionRatio, compressionLevel,
uploadBytes, uploadDurationMs, snapshotID) uploadDurationMs, snapshotID)
} }
if err != nil { if err != nil {
@@ -544,30 +543,6 @@ 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(
-59
View File
@@ -145,65 +145,6 @@ 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()
+16
View File
@@ -158,3 +158,19 @@ 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
}
+4
View File
@@ -66,6 +66,7 @@ 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
@@ -230,6 +231,9 @@ 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
+52 -28
View File
@@ -92,6 +92,9 @@ 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
@@ -131,9 +134,7 @@ type ScannerConfig struct {
SkipErrors bool SkipErrors bool
} }
// ScanResult contains the results of a scan operation. Files and bytes // ScanResult contains the results of a scan operation
// 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
@@ -143,9 +144,6 @@ 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
} }
@@ -213,6 +211,7 @@ 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(),
} }
@@ -220,9 +219,7 @@ 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(func(blobWithReader *blob.WithReader) error { s.packer.SetBlobHandler(s.handleBlobReady)
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")
} }
@@ -291,7 +288,8 @@ 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)")
} }
result.EndTime = time.Now().UTC() // Finalize result with blob statistics
s.finalizeScanResult(ctx, result)
return result, nil return result, nil
} }
@@ -433,6 +431,27 @@ 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
@@ -1492,24 +1511,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( func (s *Scanner) handleBlobReady(blobWithReader *blob.WithReader) error {
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( blobExists, err := s.uploadBlobIfNeeded(ctx, blobPath, blobWithReader, startTime)
ctx, blobPath, blobWithReader, startTime, result)
if err != nil { if err != nil {
s.cleanupBlobTempFile(blobWithReader) s.cleanupBlobTempFile(blobWithReader)
@@ -1544,7 +1563,6 @@ 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
@@ -1580,10 +1598,6 @@ 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),
@@ -1799,9 +1813,9 @@ func (s *Scanner) processFileStreaming(
size: chunk.Size, size: chunk.Size,
}) })
if !chunkExists { s.updateChunkStats(chunkExists, chunk.Size, result)
s.updateChunkStats(chunk.Size, result)
if !chunkExists {
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:
@@ -1829,11 +1843,20 @@ func (s *Scanner) processFileStreaming(
return nil return nil
} }
// updateChunkStats counts a chunk that was not already stored. The scan // updateChunkStats updates scan result and progress stats for a processed chunk
// result's file counts, BytesScanned and BytesSkipped are not touched func (s *Scanner) updateChunkStats(
// here: the scan phase counts each file once. chunkExists bool, chunkSize int64, result *ScanResult,
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)
@@ -1841,6 +1864,7 @@ func (s *Scanner) updateChunkStats(chunkSize int64, result *ScanResult) {
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
// blob if needed // blob if needed
+23 -5
View File
@@ -154,6 +154,26 @@ 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(
@@ -165,8 +185,8 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
int64(stats.FilesScanned), int64(stats.FilesScanned),
int64(stats.ChunksCreated), int64(stats.ChunksCreated),
int64(stats.BlobsCreated), int64(stats.BlobsCreated),
stats.TotalSize, stats.BytesScanned,
stats.BlobSize, stats.BytesUploaded,
) )
if err != nil { if err != nil {
return err return err
@@ -176,7 +196,6 @@ 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,
) )
}) })
@@ -871,7 +890,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
TotalSize int64 // Total size of all files examined BytesScanned int64
ChunksCreated int ChunksCreated int
BlobsCreated int BlobsCreated int
BytesUploaded int64 BytesUploaded int64
@@ -881,7 +900,6 @@ 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
+58 -41
View File
@@ -189,11 +189,6 @@ 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
@@ -233,6 +228,8 @@ 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
@@ -317,9 +314,6 @@ 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,
@@ -335,6 +329,18 @@ 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
@@ -344,39 +350,31 @@ func (v *Vaultik) scanAllDirectories(
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,
TotalSize: stats.totalBytes + stats.totalBytesSkipped, BytesScanned: stats.totalBytes,
ChunksCreated: stats.totalChunks, ChunksCreated: stats.totalChunks,
BlobsCreated: stats.totalBlobs, BlobsCreated: stats.totalBlobs,
BytesUploaded: stats.totalBytesUploaded, BytesUploaded: stats.totalBytesUploaded,
}, },
BlobSize: stats.blobSize, BlobUncompressedSize: 0,
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 {
@@ -411,10 +409,12 @@ 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 stats.blobUncompressedSize > 0 { if uncompressedSize > 0 {
compressionRatio = float64(stats.blobSize) / compressionRatio = float64(compressedSize) / float64(uncompressedSize)
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(stats.blobSize), v.UI.Size(compressedSize),
v.UI.Size(stats.blobUncompressedSize), v.UI.Size(uncompressedSize),
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,6 +455,27 @@ 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
@@ -1229,10 +1250,8 @@ func (v *Vaultik) confirmRemoveSnapshot(snapshotID string, opts *RemoveOptions)
// removeSnapshotRemote strips the snapshot's metadata from the // removeSnapshotRemote strips the snapshot's metadata from the
// destination store, warning and proceeding on failure: the local-DB // destination store, warning and proceeding on failure: the local-DB
// removal has already happened, so the user is told the remote half // removal has already happened, so the user is told the remote half
// didn't finish and to run `vaultik snapshot remove` for the snapshot // didn't finish and can retry with `vaultik prune` once the destination
// again once the destination store is reachable (`vaultik prune` never // store is reachable. Returns true when the remote removal succeeded.
// removes snapshot metadata). Returns true when the remote removal
// succeeded.
func (v *Vaultik) removeSnapshotRemote(snapshotID string, opts *RemoveOptions) bool { func (v *Vaultik) removeSnapshotRemote(snapshotID string, opts *RemoveOptions) bool {
log.Info("Removing snapshot metadata from remote storage", log.Info("Removing snapshot metadata from remote storage",
"snapshot_id", snapshotID) "snapshot_id", snapshotID)
@@ -1241,18 +1260,16 @@ func (v *Vaultik) removeSnapshotRemote(snapshotID string, opts *RemoveOptions) b
err := v.deleteRemoteSnapshotByKey(remoteKey) err := v.deleteRemoteSnapshotByKey(remoteKey)
if err != nil { if err != nil {
removeCommand := "vaultik snapshot remove " + snapshotID
log.Warn("Could not remove snapshot metadata from remote storage; "+ log.Warn("Could not remove snapshot metadata from remote storage; "+
"run '"+removeCommand+"' again once the remote is reachable", "run '"+pruneCommandHint+"' once the remote is reachable "+
"error", err) "to finish cleanup", "error", err)
// The UI writes to stdout, which under --json holds only the // The UI writes to stdout, which under --json holds only the
// document; the log record above is the warning on stderr. // document; the log record above is the warning on stderr.
if v.UI != nil && !opts.JSON { if v.UI != nil && !opts.JSON {
v.UI.Warningf("Could not remove snapshot metadata from remote: "+ v.UI.Warningf("Could not remove snapshot metadata from remote: "+
"%v. Run '%s' again once the remote is reachable.", "%v. Run '%s' once the remote is reachable to finish cleanup.",
err, removeCommand) err, pruneCommandHint)
} }
return false return false
-285
View File
@@ -1,285 +0,0 @@
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)
}