A machine restoring after the original is gone has no local index and cannot know a snapshot's human ID; snapshot list shows such snapshots only by their remote key, but restore and verify accepted only the human ID, so recovery could not be done as documented. Restore and verify now also accept a remote key, or an unambiguous leading part of it as snapshot list prints it, resolved against the store's metadata listing. Human IDs are never pure hex, which tells the two forms apart. Deep verify reads the single snapshot in the downloaded per-snapshot database. A new README section walks the recovery end to end; a test backs up, then lists, restores and deep-verifies with an empty index, another hostname and no age_recipients. model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)
720 lines
19 KiB
Go
720 lines
19 KiB
Go
package vaultik
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"database/sql"
|
|
"encoding/hex"
|
|
"errors"
|
|
"fmt"
|
|
"hash"
|
|
"io"
|
|
"os"
|
|
"time"
|
|
|
|
"github.com/klauspost/compress/zstd"
|
|
|
|
// Blank import registers the pure-Go sqlite driver for database/sql.
|
|
_ "modernc.org/sqlite"
|
|
"sneak.berlin/go/vaultik/internal/log"
|
|
"sneak.berlin/go/vaultik/internal/snapshot"
|
|
)
|
|
|
|
// Sentinel errors for snapshot verification failures.
|
|
var (
|
|
errVerificationFailed = errors.New("verification failed")
|
|
errSecretKeyRequired = errors.New(
|
|
"VAULTIK_AGE_SECRET_KEY not set; required for deep verification")
|
|
errChunksOutOfOrder = errors.New("chunks out of order")
|
|
errChunkHashMismatch = errors.New("chunk hash mismatch")
|
|
errTrailingBlobData = errors.New(
|
|
"blob has unexpected trailing bytes not covered by chunk list")
|
|
errManifestExtraBlob = errors.New("manifest contains blob not in database")
|
|
errBlobSizeMismatch = errors.New("blob size mismatch")
|
|
)
|
|
|
|
// verifyStatusFailed is the JSON status value for a failed verification.
|
|
const verifyStatusFailed = "failed"
|
|
|
|
// VerifyOptions contains options for the verify command
|
|
type VerifyOptions struct {
|
|
Deep bool
|
|
JSON bool
|
|
}
|
|
|
|
// VerifyResult contains the result of a snapshot verification
|
|
//
|
|
//nolint:tagliatelle // snake_case is the established JSON output format
|
|
type VerifyResult struct {
|
|
SnapshotID string `json:"snapshot_id"`
|
|
Status string `json:"status"` // "ok" or "failed"
|
|
Mode string `json:"mode"` // "shallow" or "deep"
|
|
BlobCount int `json:"blob_count"`
|
|
TotalSize int64 `json:"total_size"`
|
|
Verified int `json:"verified"`
|
|
Missing int `json:"missing"`
|
|
MissingSize int64 `json:"missing_size,omitempty"`
|
|
ErrorMessage string `json:"error,omitempty"`
|
|
}
|
|
|
|
// deepVerifyFailure records a failure in the result and returns it appropriately
|
|
func (v *Vaultik) deepVerifyFailure(
|
|
result *VerifyResult, opts *VerifyOptions, msg string, err error,
|
|
) error {
|
|
result.Status = verifyStatusFailed
|
|
|
|
result.ErrorMessage = msg
|
|
if opts.JSON {
|
|
return v.outputVerifyJSON(result)
|
|
}
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return fmt.Errorf("%w: %s", errVerificationFailed, msg)
|
|
}
|
|
|
|
// RunDeepVerify executes deep verification operation
|
|
func (v *Vaultik) RunDeepVerify(snapshotID string, opts *VerifyOptions) error {
|
|
result := &VerifyResult{
|
|
SnapshotID: snapshotID,
|
|
Mode: "deep",
|
|
}
|
|
|
|
if !v.CanDecrypt() {
|
|
return v.deepVerifyFailure(result, opts,
|
|
errSecretKeyRequired.Error(), errSecretKeyRequired)
|
|
}
|
|
|
|
log.Info("Starting snapshot verification", "snapshot_id", snapshotID, "mode", "deep")
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("Deep verification of snapshot: %s\n\n", snapshotID)
|
|
}
|
|
|
|
manifest, tempDB, dbBlobs, err := v.loadVerificationData(snapshotID, opts, result)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
defer func() {
|
|
if tempDB != nil {
|
|
_ = tempDB.Close()
|
|
}
|
|
}()
|
|
|
|
result.BlobCount = len(dbBlobs)
|
|
|
|
var totalSize int64
|
|
for _, blob := range dbBlobs {
|
|
totalSize += blob.CompressedSize
|
|
}
|
|
|
|
result.TotalSize = totalSize
|
|
|
|
err = v.runVerificationSteps(manifest, dbBlobs, tempDB, opts, result, totalSize)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
result.Status = "ok"
|
|
result.Verified = len(dbBlobs)
|
|
|
|
if opts.JSON {
|
|
return v.outputVerifyJSON(result)
|
|
}
|
|
|
|
log.Info("✓ Verification completed successfully",
|
|
"snapshot_id", snapshotID, "mode", "deep", "blobs_verified", len(dbBlobs))
|
|
v.stdoutf("\n✓ Verification completed successfully\n")
|
|
v.stdoutf(" Snapshot: %s\n", snapshotID)
|
|
v.stdoutf(" Blobs verified: %d\n", len(dbBlobs))
|
|
v.stdoutf(" Total size: %s\n", ubytes(totalSize))
|
|
|
|
return nil
|
|
}
|
|
|
|
// loadVerificationData downloads manifest, database, and blob list for verification
|
|
func (v *Vaultik) loadVerificationData(
|
|
snapshotID string, opts *VerifyOptions, result *VerifyResult,
|
|
) (*snapshot.Manifest, *tempDB, []snapshot.BlobInfo, error) {
|
|
// Resolve the identifier to the snapshot's remote key. A human ID is
|
|
// hashed; a remote key (or its abbreviation, as printed for a
|
|
// remote-only snapshot) is used as-is, so a host with no local index
|
|
// can verify a snapshot it can only see on the store.
|
|
remoteKey, err := v.resolveSnapshotRemoteKey(snapshotID)
|
|
if err != nil {
|
|
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
|
fmt.Sprintf("resolving snapshot identifier: %v", err), err)
|
|
}
|
|
|
|
// Download manifest. downloadManifestByKey is the single reader for
|
|
// remote manifests; see its doc comment.
|
|
log.Info("Downloading manifest", "remote_key", remoteKey)
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("Downloading manifest...\n")
|
|
}
|
|
|
|
manifest, err := v.downloadManifestByKey(remoteKey)
|
|
if err != nil {
|
|
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
|
fmt.Sprintf("failed to download manifest: %v", err),
|
|
fmt.Errorf("failed to download manifest: %w", err))
|
|
}
|
|
|
|
log.Info("Manifest loaded",
|
|
"manifest_blob_count", manifest.BlobCount,
|
|
"manifest_total_size", ubytes(manifest.TotalCompressedSize))
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("Manifest loaded: %d blobs (%s)\n",
|
|
manifest.BlobCount, ubytes(manifest.TotalCompressedSize))
|
|
v.stdoutf("Downloading and decrypting database...\n")
|
|
}
|
|
|
|
// Download and decrypt database
|
|
dbPath := fmt.Sprintf("metadata/%s/db.zst.age", remoteKey)
|
|
log.Info("Downloading encrypted database", "path", dbPath)
|
|
|
|
dbReader, err := v.Storage.Get(v.ctx, dbPath)
|
|
if err != nil {
|
|
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
|
fmt.Sprintf("failed to download database: %v", err),
|
|
fmt.Errorf("failed to download database: %w", err))
|
|
}
|
|
|
|
defer func() { _ = dbReader.Close() }()
|
|
|
|
tdb, err := v.decryptAndLoadDatabase(dbReader)
|
|
if err != nil {
|
|
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
|
fmt.Sprintf("failed to decrypt database: %v", err),
|
|
fmt.Errorf("failed to decrypt database: %w", err))
|
|
}
|
|
|
|
dbBlobs, err := v.getBlobsFromDatabase(tdb.DB)
|
|
if err != nil {
|
|
_ = tdb.Close()
|
|
|
|
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
|
fmt.Sprintf("failed to get blobs from database: %v", err),
|
|
fmt.Errorf("failed to get blobs from database: %w", err))
|
|
}
|
|
|
|
var dbTotalSize int64
|
|
for _, b := range dbBlobs {
|
|
dbTotalSize += b.CompressedSize
|
|
}
|
|
|
|
log.Info("Database loaded",
|
|
"db_blob_count", len(dbBlobs),
|
|
"db_total_size", ubytes(dbTotalSize))
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("Database loaded: %d blobs (%s)\n",
|
|
len(dbBlobs), ubytes(dbTotalSize))
|
|
}
|
|
|
|
return manifest, tdb, dbBlobs, nil
|
|
}
|
|
|
|
// runVerificationSteps executes manifest verification, blob existence
|
|
// check, and deep content verification.
|
|
func (v *Vaultik) runVerificationSteps(
|
|
manifest *snapshot.Manifest,
|
|
dbBlobs []snapshot.BlobInfo,
|
|
tdb *tempDB,
|
|
opts *VerifyOptions,
|
|
result *VerifyResult,
|
|
totalSize int64,
|
|
) error {
|
|
if !opts.JSON {
|
|
v.stdoutf("Verifying manifest against database...\n")
|
|
}
|
|
|
|
err := v.verifyManifestAgainstDatabase(manifest, dbBlobs)
|
|
if err != nil {
|
|
return v.deepVerifyFailure(result, opts, err.Error(), err)
|
|
}
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("Manifest verified.\n")
|
|
v.stdoutf("Checking blob existence in remote storage...\n")
|
|
}
|
|
|
|
err = v.verifyBlobExistenceFromDB(dbBlobs)
|
|
if err != nil {
|
|
return v.deepVerifyFailure(result, opts, err.Error(), err)
|
|
}
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf("All blobs exist.\n")
|
|
v.stdoutf("Downloading and verifying blob contents (%d blobs, %s)...\n",
|
|
len(dbBlobs), ubytes(totalSize))
|
|
}
|
|
|
|
err = v.performDeepVerificationFromDB(dbBlobs, tdb.DB, opts)
|
|
if err != nil {
|
|
return v.deepVerifyFailure(result, opts, err.Error(), err)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// tempDB wraps sql.DB with cleanup
|
|
type tempDB struct {
|
|
*sql.DB
|
|
|
|
tempPath string
|
|
}
|
|
|
|
func (t *tempDB) Close() error {
|
|
err := t.DB.Close()
|
|
_ = os.Remove(t.tempPath)
|
|
|
|
return err
|
|
}
|
|
|
|
// decryptAndLoadDatabase decrypts and loads the binary SQLite database
|
|
// from the encrypted stream.
|
|
func (v *Vaultik) decryptAndLoadDatabase(reader io.ReadCloser) (*tempDB, error) {
|
|
// Get decryptor
|
|
decryptor, err := v.GetDecryptor()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to get decryptor: %w", err)
|
|
}
|
|
|
|
// Decrypt the stream
|
|
decryptedReader, err := decryptor.DecryptStream(reader)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to decrypt database: %w", err)
|
|
}
|
|
|
|
// Decompress the binary database
|
|
decompressor, err := zstd.NewReader(decryptedReader)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to create decompressor: %w", err)
|
|
}
|
|
defer decompressor.Close()
|
|
|
|
// Create temporary file for the database
|
|
tempFile, err := os.CreateTemp("", "vaultik-verify-*.db")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to create temp file: %w", err)
|
|
}
|
|
|
|
tempPath := tempFile.Name()
|
|
|
|
// Stream decompress directly to file
|
|
log.Info("Decompressing database...")
|
|
|
|
written, err := io.Copy(tempFile, decompressor)
|
|
if err != nil {
|
|
_ = tempFile.Close()
|
|
_ = os.Remove(tempPath)
|
|
|
|
return nil, fmt.Errorf("failed to decompress database: %w", err)
|
|
}
|
|
|
|
_ = tempFile.Close()
|
|
|
|
log.Info("Database decompressed", "size", ubytes(written))
|
|
|
|
// Open the database
|
|
db, err := sql.Open("sqlite", tempPath)
|
|
if err != nil {
|
|
_ = os.Remove(tempPath)
|
|
|
|
return nil, fmt.Errorf("failed to open database: %w", err)
|
|
}
|
|
|
|
return &tempDB{
|
|
DB: db,
|
|
tempPath: tempPath,
|
|
}, nil
|
|
}
|
|
|
|
// verifyBlob downloads and verifies a single blob
|
|
func (v *Vaultik) verifyBlob(blobInfo snapshot.BlobInfo, db *sql.DB) error {
|
|
// Download blob using shared fetch method
|
|
reader, _, err := v.FetchBlob(v.ctx, blobInfo.Hash, blobInfo.CompressedSize)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to download: %w", err)
|
|
}
|
|
|
|
defer func() { _ = reader.Close() }()
|
|
|
|
// Get decryptor
|
|
decryptor, err := v.GetDecryptor()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to get decryptor: %w", err)
|
|
}
|
|
|
|
// Decrypt blob
|
|
decryptedReader, err := decryptor.DecryptStream(reader)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to decrypt: %w", err)
|
|
}
|
|
|
|
// Decompress blob
|
|
decompressor, err := zstd.NewReader(decryptedReader)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to decompress: %w", err)
|
|
}
|
|
defer decompressor.Close()
|
|
|
|
// A blob's hash — its remote name — is the double SHA256 of its
|
|
// decompressed plaintext (see blobgen.Writer.Sum256), not of the
|
|
// encrypted bytes. Hash the plaintext as chunk verification streams
|
|
// it, then compare on completion.
|
|
plaintextHasher := sha256.New()
|
|
hashedStream := io.TeeReader(decompressor, plaintextHasher)
|
|
|
|
chunkCount, err := v.verifyBlobChunks(db, blobInfo.Hash, hashedStream)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = v.verifyBlobFinalIntegrity(hashedStream, plaintextHasher, blobInfo.Hash)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
log.Info("Blob verified",
|
|
"hash", blobInfo.Hash[:16]+"...",
|
|
"chunks", chunkCount,
|
|
"size", ubytes(blobInfo.CompressedSize),
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
// verifyBlobChunks queries blob chunks from the database and verifies
|
|
// each chunk's hash against the decompressed blob stream.
|
|
func (v *Vaultik) verifyBlobChunks(
|
|
db *sql.DB, blobHash string, decompressor io.Reader,
|
|
) (int, error) {
|
|
query := `
|
|
SELECT bc.chunk_hash, bc.offset, bc.length
|
|
FROM blob_chunks bc
|
|
JOIN blobs b ON bc.blob_id = b.id
|
|
WHERE b.blob_hash = ?
|
|
ORDER BY bc.offset
|
|
`
|
|
|
|
rows, err := db.QueryContext(v.ctx, query, blobHash)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to query blob chunks: %w", err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var lastOffset int64 = -1
|
|
|
|
chunkCount := 0
|
|
totalRead := int64(0)
|
|
|
|
// Verify each chunk in the blob
|
|
for rows.Next() {
|
|
var (
|
|
chunkHash string
|
|
offset, length int64
|
|
)
|
|
|
|
err := rows.Scan(&chunkHash, &offset, &length)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to scan chunk row: %w", err)
|
|
}
|
|
|
|
// Verify chunk ordering
|
|
if offset <= lastOffset {
|
|
return 0, fmt.Errorf("%w: offset %d after %d",
|
|
errChunksOutOfOrder, offset, lastOffset)
|
|
}
|
|
|
|
lastOffset = offset
|
|
|
|
// Read chunk data from decompressed stream
|
|
if offset > totalRead {
|
|
// Skip to the correct offset
|
|
skipBytes := offset - totalRead
|
|
|
|
_, err = io.CopyN(io.Discard, decompressor, skipBytes)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to skip to offset %d: %w", offset, err)
|
|
}
|
|
|
|
totalRead = offset
|
|
}
|
|
|
|
// Read chunk data
|
|
chunkData := make([]byte, length)
|
|
|
|
_, err = io.ReadFull(decompressor, chunkData)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to read chunk at offset %d: %w", offset, err)
|
|
}
|
|
|
|
totalRead += length
|
|
|
|
// Verify chunk hash
|
|
hasher := sha256.New()
|
|
hasher.Write(chunkData)
|
|
calculatedHash := hex.EncodeToString(hasher.Sum(nil))
|
|
|
|
if calculatedHash != chunkHash {
|
|
return 0, fmt.Errorf("%w at offset %d: calculated %s, expected %s",
|
|
errChunkHashMismatch, offset, calculatedHash, chunkHash)
|
|
}
|
|
|
|
chunkCount++
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("error iterating blob chunks: %w", err)
|
|
}
|
|
|
|
return chunkCount, nil
|
|
}
|
|
|
|
// verifyBlobFinalIntegrity checks that no trailing data exists in the
|
|
// decompressed stream and that the blob hash matches the expected value.
|
|
func (v *Vaultik) verifyBlobFinalIntegrity(
|
|
plaintext io.Reader, plaintextHasher hash.Hash, expectedHash string,
|
|
) error {
|
|
// Verify no remaining data in blob - if the chunk list is accurate,
|
|
// the blob should be fully consumed.
|
|
remaining, err := io.Copy(io.Discard, plaintext)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to check for remaining blob data: %w", err)
|
|
}
|
|
|
|
if remaining > 0 {
|
|
return fmt.Errorf("%w: %d bytes", errTrailingBlobData, remaining)
|
|
}
|
|
|
|
// The blob hash is the double SHA256 of its plaintext content.
|
|
firstHash := plaintextHasher.Sum(nil)
|
|
secondHash := sha256.Sum256(firstHash)
|
|
calculatedBlobHash := hex.EncodeToString(secondHash[:])
|
|
|
|
if calculatedBlobHash != expectedHash {
|
|
return fmt.Errorf("%w: calculated %s, expected %s",
|
|
errBlobHashMismatch, calculatedBlobHash, expectedHash)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// getBlobsFromDatabase gets all blobs for the snapshot from the database.
|
|
//
|
|
// The exported per-snapshot database holds exactly one snapshot's data
|
|
// (see cleanSnapshotDB), so every row in snapshot_blobs belongs to it.
|
|
// We select them directly rather than filtering by the human snapshot ID,
|
|
// which a host restoring from the store alone does not have.
|
|
func (v *Vaultik) getBlobsFromDatabase(db *sql.DB) ([]snapshot.BlobInfo, error) {
|
|
query := `
|
|
SELECT b.blob_hash, b.compressed_size
|
|
FROM snapshot_blobs sb
|
|
JOIN blobs b ON sb.blob_hash = b.blob_hash
|
|
ORDER BY b.blob_hash
|
|
`
|
|
|
|
rows, err := db.QueryContext(v.ctx, query)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to query snapshot blobs: %w", err)
|
|
}
|
|
|
|
defer func() { _ = rows.Close() }()
|
|
|
|
var blobs []snapshot.BlobInfo
|
|
|
|
for rows.Next() {
|
|
var (
|
|
hash string
|
|
size int64
|
|
)
|
|
|
|
err := rows.Scan(&hash, &size)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to scan blob row: %w", err)
|
|
}
|
|
|
|
blobs = append(blobs, snapshot.BlobInfo{
|
|
Hash: hash,
|
|
CompressedSize: size,
|
|
})
|
|
}
|
|
|
|
err = rows.Err()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("error iterating blobs: %w", err)
|
|
}
|
|
|
|
return blobs, nil
|
|
}
|
|
|
|
// verifyManifestAgainstDatabase verifies the manifest matches the
|
|
// authoritative database.
|
|
func (v *Vaultik) verifyManifestAgainstDatabase(
|
|
manifest *snapshot.Manifest, dbBlobs []snapshot.BlobInfo,
|
|
) error {
|
|
log.Info("Verifying manifest against database")
|
|
|
|
// Build map of database blobs
|
|
dbBlobMap := make(map[string]int64)
|
|
for _, blob := range dbBlobs {
|
|
dbBlobMap[blob.Hash] = blob.CompressedSize
|
|
}
|
|
|
|
// Build map of manifest blobs
|
|
manifestBlobMap := make(map[string]int64)
|
|
for _, blob := range manifest.Blobs {
|
|
manifestBlobMap[blob.Hash] = blob.CompressedSize
|
|
}
|
|
|
|
// Check counts match
|
|
if len(dbBlobMap) != len(manifestBlobMap) {
|
|
log.Warn("Manifest blob count mismatch",
|
|
"database_blobs", len(dbBlobMap),
|
|
"manifest_blobs", len(manifestBlobMap),
|
|
)
|
|
// This is a warning, not an error - database is authoritative
|
|
}
|
|
|
|
// Check each manifest blob exists in database with correct size
|
|
for hash, manifestSize := range manifestBlobMap {
|
|
dbSize, exists := dbBlobMap[hash]
|
|
if !exists {
|
|
return fmt.Errorf("%w: %s", errManifestExtraBlob, hash)
|
|
}
|
|
|
|
if dbSize != manifestSize {
|
|
return fmt.Errorf(
|
|
"%w: blob %s: database has %d bytes, manifest has %d bytes",
|
|
errBlobSizeMismatch, hash, dbSize, manifestSize)
|
|
}
|
|
}
|
|
|
|
log.Info("✓ Manifest verified against database",
|
|
"manifest_blobs", len(manifestBlobMap),
|
|
"database_blobs", len(dbBlobMap),
|
|
)
|
|
|
|
return nil
|
|
}
|
|
|
|
// verifyBlobExistenceFromDB checks that all blobs from database exist in S3
|
|
func (v *Vaultik) verifyBlobExistenceFromDB(blobs []snapshot.BlobInfo) error {
|
|
log.Info("Verifying blob existence in S3", "blob_count", len(blobs))
|
|
|
|
for i, blob := range blobs {
|
|
// Construct blob path
|
|
blobPath := fmt.Sprintf("blobs/%s/%s/%s", blob.Hash[:2], blob.Hash[2:4], blob.Hash)
|
|
|
|
// Check blob exists
|
|
stat, err := v.Storage.Stat(v.ctx, blobPath)
|
|
if err != nil {
|
|
return fmt.Errorf("blob %s missing from storage: %w", blob.Hash, err)
|
|
}
|
|
|
|
// Verify size matches
|
|
if stat.Size != blob.CompressedSize {
|
|
return fmt.Errorf(
|
|
"%w: blob %s: S3 has %d bytes, database has %d bytes",
|
|
errBlobSizeMismatch, blob.Hash, stat.Size, blob.CompressedSize)
|
|
}
|
|
|
|
// Progress update every 100 blobs
|
|
if (i+1)%progressLogEvery == 0 || i == len(blobs)-1 {
|
|
log.Info("Blob existence check progress",
|
|
"checked", i+1,
|
|
"total", len(blobs),
|
|
"percent", fmt.Sprintf("%.1f%%",
|
|
float64(i+1)/float64(len(blobs))*percentScale),
|
|
)
|
|
}
|
|
}
|
|
|
|
log.Info("✓ All blobs exist in storage")
|
|
|
|
return nil
|
|
}
|
|
|
|
// performDeepVerificationFromDB downloads and verifies the content of
|
|
// each blob using the database as source.
|
|
func (v *Vaultik) performDeepVerificationFromDB(
|
|
blobs []snapshot.BlobInfo, db *sql.DB, opts *VerifyOptions,
|
|
) error {
|
|
// Calculate total bytes for ETA
|
|
var totalBytesExpected int64
|
|
for _, b := range blobs {
|
|
totalBytesExpected += b.CompressedSize
|
|
}
|
|
|
|
log.Info("Starting deep verification - downloading and verifying all blobs",
|
|
"blob_count", len(blobs),
|
|
"total_size", ubytes(totalBytesExpected),
|
|
)
|
|
|
|
startTime := time.Now()
|
|
bytesProcessed := int64(0)
|
|
|
|
for i, blobInfo := range blobs {
|
|
// Verify individual blob
|
|
err := v.verifyBlob(blobInfo, db)
|
|
if err != nil {
|
|
return fmt.Errorf("blob %s verification failed: %w", blobInfo.Hash, err)
|
|
}
|
|
|
|
bytesProcessed += blobInfo.CompressedSize
|
|
elapsed := time.Since(startTime)
|
|
remaining := len(blobs) - (i + 1)
|
|
|
|
// Calculate ETA based on bytes processed
|
|
var eta time.Duration
|
|
|
|
if bytesProcessed > 0 {
|
|
bytesPerSec := float64(bytesProcessed) / elapsed.Seconds()
|
|
|
|
bytesRemaining := totalBytesExpected - bytesProcessed
|
|
if bytesPerSec > 0 {
|
|
eta = time.Duration(float64(bytesRemaining)/bytesPerSec) * time.Second
|
|
}
|
|
}
|
|
|
|
log.Info("Verification progress",
|
|
"blobs_done", i+1,
|
|
"blobs_total", len(blobs),
|
|
"blobs_remaining", remaining,
|
|
"bytes_done", bytesProcessed,
|
|
"bytes_done_human", ubytes(bytesProcessed),
|
|
"bytes_total", totalBytesExpected,
|
|
"bytes_total_human", ubytes(totalBytesExpected),
|
|
"elapsed", elapsed.Round(time.Second),
|
|
"eta", eta.Round(time.Second),
|
|
)
|
|
|
|
if !opts.JSON {
|
|
v.stdoutf(" Verified %d/%d blobs (%d remaining) - %s/%s - elapsed %s, eta %s\n",
|
|
i+1, len(blobs), remaining,
|
|
ubytes(bytesProcessed),
|
|
ubytes(totalBytesExpected),
|
|
elapsed.Round(time.Second),
|
|
eta.Round(time.Second))
|
|
}
|
|
}
|
|
|
|
totalElapsed := time.Since(startTime)
|
|
log.Info("✓ Deep verification completed successfully",
|
|
"blobs_verified", len(blobs),
|
|
"total_bytes", bytesProcessed,
|
|
"total_bytes_human", ubytes(bytesProcessed),
|
|
"duration", totalElapsed.Round(time.Second),
|
|
)
|
|
|
|
return nil
|
|
}
|