2 Commits
Author SHA1 Message Date
sneak 80b44d1c16 Count each file, byte and upload once in backup statistics (closes #225)
check / check (push) Waiting to run
The scanner added a changed file's bytes again for each new chunk and
counted a file as unchanged for each chunk already stored. It now
counts files and bytes once, in the scan phase, and counts its own
uploads, so a --cron run, which has no progress reporter, records
them. The blob count no longer adds earlier paths' blobs again.

The snapshots row now stores the size of all files in total_size and
the referenced blobs' sizes in blob_size, blob_uncompressed_size and
compression_ratio, as docs/DATAMODEL.md says. Those sizes come from
one query, and a failed query fails the snapshot. DATAMODEL.md now
says chunk_count and blob_count count what the run added.

Removed UpdateSnapshotStats and GetCountBySnapshot, which nothing
calls any more.

Model: opus-5-5
2026-10-07 00:00:51 +00:00
clawbot 85d4ef118d Keep command output to the README's stdout and stderr rules (closes #224)
check / check (push) Successful in 16m49s
The startup banner moves from stdout to stderr, so a `completion`
script, a `config get` value and the hidden `__complete` command print
only their own output. `--quiet`, `--cron` and `--json` still suppress
it.

A failing `remote info`, `prune` or `snapshot remove` under `--json`
now reports its error on stderr. Their reporters returned early under
`--json`, so the failure reached neither stream.

`snapshot verify --quiet` writes no report. A failure is still
returned and printed on stderr, with the same exit status.

Judgement call: the banner's stream, posted on the issue for the owner.

Model: opus-5-5
2026-10-07 01:12:15 +02:00
23 changed files with 836 additions and 292 deletions
+2 -2
View File
@@ -335,10 +335,10 @@ CreateSnapshot(opts)
│ │
│ └─► Accumulate statistics
│
├─► SnapshotManager.UpdateSnapshotStatsExtended()
│
├─► SnapshotManager.PopulateSnapshotBlobs() // record referenced blobs
│
├─► SnapshotManager.UpdateSnapshotStatsExtended()
│
├─► SnapshotManager.ExportSnapshotMetadata()
│ │
│ ├─► Copy database to temp file
+20 -16
View File
@@ -200,8 +200,10 @@ local index or the destination store. `config`, `database delete`,
### stdout and stderr
Log output — everything from `--verbose` and `--debug`, and every
warning and error the logger emits — goes to **stderr**. stdout carries
the output you asked for: tables, and the documents produced by `--json`.
warning and error the logger emits — goes to **stderr**, and so does the
startup banner. stdout carries the output you asked for: tables, the
documents produced by `--json`, `config get` values, and completion
scripts.
This means `vaultik snapshot list --verbose > out.txt` captures the
listing and leaves the diagnostics on your terminal. To capture both,
@@ -645,8 +647,8 @@ Work planned after 1.0. Loosely ordered by priority.
## output style
Every command's user-facing output is governed by `internal/ui`, in one
of two ways. Color is enabled when stdout is a TTY and the `NO_COLOR`
environment variable is unset (https://no-color.org/).
of two ways. Color is enabled when the stream written to is a TTY and
the `NO_COLOR` environment variable is unset (https://no-color.org/).
* **Status, progress, warnings, and errors** go through the `internal/ui`
message methods below: marker-prefixed, colored on a TTY, and — except
@@ -656,24 +658,26 @@ environment variable is unset (https://no-color.org/).
`config init`, `config set`, and `database delete`.
* **The data a command exists to produce** is written plain, with no
marker and no color, because a marker would corrupt a table or a
parsed document. This covers the `version`, `info`, and `remote info`
reports, the `snapshot list` table, `config get` values, and every
`--json` document. `--quiet` silences the human reports and tables
(`version`, `info`, `remote info`, `snapshot list`) but never the
machine-consumed `config get` value or the `--json` documents, which a
script depends on. The `database delete` confirmation prompt is also
written this way and always shown: it is an interactive exchange the
operator must see.
parsed document. This covers the `version`, `info`, `remote info` and
`snapshot verify` reports, the `snapshot list` table, `config get`
values, and every `--json` document. `--quiet` silences the human
reports and tables (`version`, `info`, `remote info`, `snapshot
verify`, `snapshot list`) but never the machine-consumed `config get`
value or the `--json` documents, which a script depends on. The
`database delete` confirmation prompt is also written this way and
always shown: it is an interactive exchange the operator must see.
`internal/ui` writes to stdout; it is the output the user asked for.
Structured log records are a different thing and go through
`internal/log`, which writes to stderr (see "stdout and stderr" above).
`internal/ui` writes to stdout; it is the output the user asked for. The
exceptions are the startup banner and the error a failed command ends
with, which go to stderr. Structured log records are a different thing
and go through `internal/log`, which writes to stderr (see "stdout and
stderr" above).
Message classes:
| Class | Marker | Alignment | Use for |
|-------|--------|-----------|---------|
| Banner | none | column 0 | The startup line printed once per invocation |
| Banner | none | column 0 | The startup line printed once per invocation, on stderr |
| Begin | `》` (white) | column 0 | An operation is about to start (present-continuous verb) |
| Complete | `》` (green) | column 0 | An operation just finished (past-tense verb) |
| Info | `》` (white) | column 0 | Neutral status update |
+22
View File
@@ -22,6 +22,28 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps
- 2026-10-06: Made the backup summary and the `snapshots` row count each
file, byte and upload once
([issue #225](https://git.eeqj.de/sneak/vaultik/issues/225)). The
scanner added a file's bytes again for each new chunk and counted a
file as unchanged for each chunk already stored, so a first backup
reported twice its size and "backed up" could go negative. Upload
figures came from the progress reporter, which `--cron` turns off, and
`blob_count` counted earlier paths' blobs again for each later path.
The scanner now counts uploads itself; `blob_size`,
`blob_uncompressed_size` and `compression_ratio` describe the blobs
the snapshot references, and `docs/DATAMODEL.md` now says
`chunk_count` and `blob_count` count what the run added.
- 2026-10-06: Made command output follow the README's stdout and stderr
rules ([issue #224](https://git.eeqj.de/sneak/vaultik/issues/224)). The
startup banner went to stdout, so a `completion` script or a
`config get` value started with it; the banner now goes to stderr. A
failing `remote info`, `prune` or `snapshot remove` under `--json`
printed nothing on either stream, and now reports its error on stderr.
`snapshot verify --quiet` printed its whole report; it now prints
none, and a failure still reaches stderr with the same exit status.
- 2026-10-06: Made a backup without `--cron` of a snapshot with two or
more `paths` complete instead of panicking with `close of closed
channel` ([issue #253](https://git.eeqj.de/sneak/vaultik/issues/253)).
+3 -3
View File
@@ -117,10 +117,10 @@ Tracks backup snapshots.
- `started_at` (INTEGER) - Start timestamp
- `completed_at` (INTEGER) - Completion timestamp (NULL if in progress)
- `file_count` (INTEGER) - Number of files in snapshot
- `chunk_count` (INTEGER) - Number of unique chunks
- `blob_count` (INTEGER) - Number of blobs referenced
- `chunk_count` (INTEGER) - Number of chunks this snapshot stored that were not stored before
- `blob_count` (INTEGER) - Number of blobs this snapshot created
- `total_size` (INTEGER) - Total size of all files
- `blob_size` (INTEGER) - Total size of all blobs (compressed)
- `blob_size` (INTEGER) - Total compressed size of all referenced blobs
- `blob_uncompressed_size` (INTEGER) - Total uncompressed size of all referenced blobs
- `compression_ratio` (REAL) - Compression ratio achieved
- `compression_level` (INTEGER) - Compression level used for this snapshot
+13 -18
View File
@@ -200,10 +200,10 @@ func RunApp(ctx context.Context, app *fx.App) error {
}
// errReported marks a failure the operation has already shown the user
// (and deliberately withheld under --json). Entry turns it into a
// non-zero exit status without printing anything further, so the error
// line is not doubled. It flows up from RunOperation through cobra to
// Entry.
// (or, under `snapshot verify --json`, put in its document). Entry
// turns it into a non-zero exit status without printing anything
// further, so the error line is not doubled. It flows up from
// RunOperation through cobra to Entry.
var errReported = errors.New("operation failed")
// RunOperation runs op against the Vaultik instance inside the fx app
@@ -220,10 +220,10 @@ var errReported = errors.New("operation failed")
// interrupt OnStop cancels op and waits for the goroutine to return, so
// op's cleanup (removing decrypted scratch files) runs before the
// process exits; the wait is bounded by shutdownTimeout. report is
// called with a non-canceled failure so the caller can log it (and
// suppress it under --json) before it becomes errReported. A context
// cancellation is the interrupt path, not a failure: it is neither
// reported nor counted as one.
// called with a non-canceled failure so the caller can show it to the
// user before it becomes errReported. A context cancellation is the
// interrupt path, not a failure: it is neither reported nor counted as
// one.
func RunOperation(
ctx context.Context, opts AppOptions,
op func(v *vaultik.Vaultik) error, report func(err error),
@@ -293,13 +293,12 @@ func RunOperation(
// runVaultikApp runs the standard single-operation command lifecycle
// shared by the snapshot list/purge/remove and remote nuke subcommands:
// resolve the config, then run op against the Vaultik instance through
// RunOperation, reporting a failure prefixed with failMsg (suppressed
// while suppressErrors is true, e.g. under --json). mode says whether the
// command takes the PID lock. jsonOutput marks a command whose stdout is a
// JSON document: it quiets the UI but, unlike Quiet, leaves the stderr log
// level alone.
// RunOperation, reporting a failure prefixed with failMsg on stderr. mode
// says whether the command takes the PID lock. jsonOutput marks a command
// whose stdout is a JSON document: it quiets the UI but, unlike Quiet,
// leaves the stderr log level alone.
func runVaultikApp(
cmd *cobra.Command, mode lockMode, jsonOutput, suppressErrors bool,
cmd *cobra.Command, mode lockMode, jsonOutput bool,
failMsg string, op func(v *vaultik.Vaultik) error,
) error {
configPath, err := ResolveConfigPath()
@@ -319,10 +318,6 @@ func runVaultikApp(
},
Mode: mode,
}, op, func(err error) {
if suppressErrors {
return
}
log.Error(failMsg, "error", err)
ReportErrorf("%s: %v", failMsg, err)
})
+10 -9
View File
@@ -16,16 +16,18 @@ import (
const shortCommitLen = 12
// Entry is the main entry point for the CLI application.
// It prints the startup banner to stdout (unless a banner-suppressing
// It prints the startup banner to stderr (unless a banner-suppressing
// flag is present in os.Args — see bannerSuppressedInArgs), executes the
// root cobra command, and routes any returned error through the
// ui.Writer so the user sees a properly formatted "🛑 ERROR:" line.
// The banner goes to stderr because stdout carries only the output the
// user asked for, such as a completion script or a `config get` value.
//
// It returns the process exit code (0 on success, 1 on error) rather
// than calling os.Exit, so that main's deferred profile writers run
// before the process ends. See run in cmd/vaultik/main.go.
func Entry() int {
emitStartupBanner(os.Args[1:], os.Stdout)
emitStartupBanner(os.Args[1:], os.Stderr)
rootCmd := NewRootCommand()
rootCmd.SilenceErrors = true
@@ -33,8 +35,9 @@ func Entry() int {
err := rootCmd.Execute()
if err != nil {
// An operation that ran inside the fx app has already reported
// its own failure (and suppressed it under --json); errReported
// says so. Printing it again here would double the error line.
// its own failure (`snapshot verify --json` puts it in the
// document instead); errReported says so. Printing it again
// here would double the error line.
// Every other error — bad arguments, a config that would not
// load — reaches Entry unreported, so it is shown here.
if !errors.Is(err, errReported) {
@@ -49,9 +52,8 @@ func Entry() int {
// emitStartupBanner writes the startup banner to w unless args (the
// argument vector with the program name already stripped) contains a
// flag that suppresses it. Split out of Entry so that the decision — the
// only thing standing between a --json invocation and a parseable
// stdout — is reachable from a test without running the whole CLI.
// flag that suppresses it. Split out of Entry so that the decision is
// reachable from a test without running the whole CLI.
func emitStartupBanner(args []string, w io.Writer) {
if bannerSuppressedInArgs(args) {
return
@@ -86,8 +88,7 @@ func ReportErrorf(format string, args ...any) {
// --json is a subcommand flag rather than a persistent one, but so is
// --cron (it exists only on `snapshot create`), so this adds no new
// class of imprecision. The only cost of a false positive is a missing
// decorative banner; the cost of a false negative is a corrupt document
// on stdout, so the scan errs deliberately in that direction.
// decorative banner.
func bannerSuppressedInArgs(args []string) bool {
for _, a := range args {
if a == "--" {
+16 -43
View File
@@ -36,23 +36,14 @@ const (
// strips it before scanning, so it has to be present.
programName = "vaultik"
// someSnapshotID is any snapshot identifier: these tests never run
// the command, so it only has to occupy the positional argument.
// someSnapshotID only fills the positional argument; no test needs
// the snapshot to exist.
someSnapshotID = "host_2026-01-01T00:00:00Z"
)
// placeholderJSONDocument stands in for whatever document a --json
// command writes to stdout. `snapshot list --json` with no snapshots
// prints exactly this; the other --json commands print an object rather
// than an array, but this test is not about their shape. It is about
// what is on stdout *before* them, which is the same for all of them
// because Entry prints the banner before cobra has parsed anything and
// therefore before it can know which command is running.
const placeholderJSONDocument = "[]\n"
// jsonArgumentVectors are the argument vectors of every --json
// invocation the CLI accepts, with the program name stripped exactly as
// Entry strips it. Each one must leave stdout untouched by the banner.
// Entry strips it. Each one must suppress the banner.
//
//nolint:gochecknoglobals // read-only test fixture shared by two tests
var jsonArgumentVectors = map[string][]string{
@@ -74,39 +65,23 @@ var jsonArgumentVectors = map[string][]string{
},
}
// TestJSONInvocationStdoutIsExactlyOneDocument is the CLI-layer
// regression guard for issue #106: `vaultik snapshot list --json | jq`
// must work with no other flags.
//
// internal/vaultik's TestListSnapshots_JSONStdoutIsOnlyTheDocument
// guards the same contract one layer down, but it calls the library
// function directly and so cannot see Entry, which is where the
// contamination was: the startup banner is written to stdout before
// cobra parses anything, and the suppression scan did not know about
// --json. The two banner lines and the blank line landed ahead of the
// document and `jq` refused the result.
//
// The document is a constant here because this test is about the
// argument vectors, one per --json command; the one that runs a real
// command end to end is TestEntryJSONStdoutIsExactlyOneDocument below.
func TestJSONInvocationStdoutIsExactlyOneDocument(t *testing.T) {
// TestJSONInvocationSuppressesBanner checks that every --json
// invocation suppresses the startup banner, as the README says --json
// does along with --quiet and --cron. The scan is over the raw argument
// vector, so each position and spelling of --json is listed.
func TestJSONInvocationSuppressesBanner(t *testing.T) {
t.Parallel()
for name, argv := range jsonArgumentVectors {
t.Run(name, func(t *testing.T) {
t.Parallel()
var stdout bytes.Buffer
var banner bytes.Buffer
emitStartupBanner(argv, &stdout)
emitStartupBanner(argv, &banner)
require.Empty(t, stdout.String(),
"nothing may reach stdout ahead of a --json document")
_, err := stdout.WriteString(placeholderJSONDocument)
require.NoError(t, err)
requireExactlyOneJSONDocument(t, stdout.String())
assert.Empty(t, banner.String(),
"--json suppresses the banner")
})
}
}
@@ -127,11 +102,11 @@ func TestBannerStillPrintedWithoutSuppressingFlag(t *testing.T) {
t.Run(name, func(t *testing.T) {
t.Parallel()
var stdout bytes.Buffer
var banner bytes.Buffer
emitStartupBanner(argv, &stdout)
emitStartupBanner(argv, &banner)
assert.Contains(t, stdout.String(), "starting up at",
assert.Contains(t, banner.String(), "starting up at",
"the banner belongs on invocations that did not opt out")
})
}
@@ -247,9 +222,7 @@ func TestEntryJSONStdoutIsExactlyOneDocument(t *testing.T) {
// captureProcessStdout redirects the process's own stdout to a pipe for
// the duration of fn and returns what was written to it. The redirection
// has to be at the file-descriptor level rather than through an injected
// writer, because the banner and the JSON encoder reach os.Stdout
// independently and the point of the test is that both land in the same
// place.
// writer, because the commands Entry runs reach os.Stdout directly.
//
// Not parallel-safe: os.Stdout is process-global.
func captureProcessStdout(t *testing.T, fn func()) string {
+3 -3
View File
@@ -15,8 +15,8 @@ import (
// run, so a failing command must come back with a non-zero code rather
// than ending the process here.
//
// Stdout is captured only to keep the banner and command output off the
// test log; the assertion is on the returned code.
// Stdout and stderr are captured only to keep the banner and command
// output off the test log; the assertion is on the returned code.
//
//nolint:paralleltest // replaces os.Args and rootFlags
func TestEntryReturnsStatusCode(t *testing.T) {
@@ -50,7 +50,7 @@ func TestEntryReturnsStatusCode(t *testing.T) {
var code int
_ = captureProcessStdout(t, func() { code = Entry() })
_, _ = captureProcessStdoutAndStderr(t, func() { code = Entry() })
assert.Equal(t, testCase.want, code)
})
+152
View File
@@ -0,0 +1,152 @@
package cli //nolint:testpackage // shares hermeticConfig and the capture helpers
import (
"context"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"github.com/adrg/xdg"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/database"
)
// TestEntryCompletionStdoutIsTheScript runs `vaultik completion bash`,
// whose stdout the README tells the user to source. The script has to
// start on the first line.
//
//nolint:paralleltest // replaces os.Args, os.Stdout and os.Stderr
func TestEntryCompletionStdoutIsTheScript(t *testing.T) {
code, stdout, _ := runEntry(t, "completion", "bash")
require.Equal(t, 0, code)
firstLine, _, _ := strings.Cut(stdout, "\n")
assert.True(t, strings.HasPrefix(firstLine, "# bash completion"),
"the first line of stdout must be the script's, got %q", firstLine)
}
// TestEntryConfigGetStdoutIsTheValue runs `vaultik config get`, whose
// stdout a script reads as the value and nothing else.
//
//nolint:paralleltest // replaces os.Args, os.Stdout and os.Stderr
func TestEntryConfigGetStdoutIsTheValue(t *testing.T) {
configPath := filepath.Join(t.TempDir(), "config.yml")
require.NoError(t, os.WriteFile(configPath,
[]byte("hostname: test-host\n"), configFileMode))
code, stdout, _ := runEntry(t,
flagConfig, configPath, "config", "get", "hostname")
require.Equal(t, 0, code)
assert.Equal(t, "test-host\n", stdout)
}
// TestEntryJSONFailureIsReportedOnStderr runs each --json command that
// writes no document when it fails, against a destination it cannot
// use. The error must reach stderr, and stdout must stay empty.
//
//nolint:paralleltest // replaces os.Args, os.Stdout, os.Stderr and the xdg globals
func TestEntryJSONFailureIsReportedOnStderr(t *testing.T) {
for _, testCase := range []struct {
name string
args []string
wantOnStderr string
}{
{
name: "remote info",
args: []string{cmdRemote, cmdInfo, flagJSON},
wantOnStderr: "Failed to get remote info",
},
{
name: "prune",
args: []string{cmdPrune, flagJSON},
wantOnStderr: "Prune failed",
},
{
name: "snapshot remove",
args: []string{cmdSnapshot, cmdRemove, someSnapshotID, flagJSON},
wantOnStderr: "Failed to remove snapshot",
},
} {
t.Run(testCase.name, func(t *testing.T) {
configPath := writeUnusableDestinationConfig(t)
code, stdout, stderr := runEntry(t,
append([]string{flagConfig, configPath}, testCase.args...)...)
assert.Equal(t, 1, code)
assert.Empty(t, stdout,
"a failed --json command has no document to write")
assert.Contains(t, stderr, testCase.wantOnStderr,
"the failure must be reported on stderr")
})
}
}
// writeUnusableDestinationConfig builds a config whose destination
// directory does not exist, which fails `remote info`, and whose local
// index is bound to another destination, which fails `prune` and
// `snapshot remove` (a missing destination alone only makes `snapshot
// remove` warn). Returns the config path.
func writeUnusableDestinationConfig(t *testing.T) string {
t.Helper()
dir := t.TempDir()
configPath := filepath.Join(dir, "config.yml")
indexPath := filepath.Join(dir, "index.sqlite")
contents := fmt.Sprintf(hermeticConfig,
filepath.Join(dir, "source"),
filepath.Join(dir, "missing-store"),
indexPath)
require.NoError(t,
os.WriteFile(configPath, []byte(contents), configFileMode))
// The PID lock lives under xdg.DataHome, which xdg resolves at
// package init; point it at the temp dir so the test neither
// touches nor collides with the real one.
t.Setenv("XDG_DATA_HOME", filepath.Join(dir, "data"))
xdg.Reload()
t.Cleanup(xdg.Reload)
ctx := context.Background()
db, err := database.New(ctx, indexPath)
require.NoError(t, err)
defer func() { require.NoError(t, db.Close()) }()
require.NoError(t, database.NewRepositories(db).LocalMeta.Set(ctx,
database.LocalMetaKeyStorageURL, "file://"+filepath.Join(dir, "other")))
return configPath
}
// runEntry runs Entry with args after the program name and returns its
// exit code and what it wrote to stdout and stderr.
//
// Not parallel-safe: it replaces os.Args, os.Stdout and os.Stderr.
func runEntry(t *testing.T, args ...string) (int, string, string) {
t.Helper()
previousArgs := os.Args
t.Cleanup(func() {
os.Args = previousArgs
rootFlags = RootFlags{}
})
os.Args = append([]string{programName}, args...)
var code int
stdout, stderr := captureProcessStdoutAndStderr(t,
func() { code = Entry() })
return code, stdout, stderr
}
-4
View File
@@ -48,10 +48,6 @@ work (e.g. after a crashed backup or to reclaim storage).`,
}, func(v *vaultik.Vaultik) error {
return v.Prune(opts)
}, func(err error) {
if opts.JSON {
return
}
log.Error("Prune operation failed", "error", err)
ReportErrorf("Prune failed: %v", err)
})
+1 -5
View File
@@ -45,7 +45,7 @@ This is destructive and irreversible. Requires --force.`,
return errNukeNeedsForce
}
return runVaultikApp(cmd, mutating, false, false, "Remote nuke failed",
return runVaultikApp(cmd, mutating, false, "Remote nuke failed",
func(v *vaultik.Vaultik) error {
return v.NukeRemote(true)
})
@@ -92,10 +92,6 @@ func newRemoteInfoCommand() *cobra.Command {
}, func(v *vaultik.Vaultik) error {
return v.RemoteInfo(jsonOutput)
}, func(err error) {
if jsonOutput {
return
}
log.Error("Failed to get remote info", "error", err)
ReportErrorf("Failed to get remote info: %v", err)
})
+3 -3
View File
@@ -126,7 +126,7 @@ func newSnapshotListCommand() *cobra.Command {
Long: "Lists all snapshots with their ID, timestamp, and compressed size",
Args: cobra.NoArgs,
RunE: func(cmd *cobra.Command, _ []string) error {
return runVaultikApp(cmd, readOnly, false, false,
return runVaultikApp(cmd, readOnly, false,
"Failed to list snapshots",
func(v *vaultik.Vaultik) error {
return v.ListSnapshots(jsonOutput)
@@ -162,7 +162,7 @@ restrict the operation to specific snapshot names.`,
return errPurgeCriteriaBoth
}
return runVaultikApp(cmd, mutating, false, false,
return runVaultikApp(cmd, mutating, false,
"Failed to purge snapshots",
func(v *vaultik.Vaultik) error {
return v.PurgeSnapshotsWithOptions(opts)
@@ -265,7 +265,7 @@ To wipe the entire destination store and start over, use 'vaultik remote
nuke --force' — it is the single supported entry point for that.`,
Args: requireSnapshotIDArg,
RunE: func(cmd *cobra.Command, args []string) error {
return runVaultikApp(cmd, mutating, opts.JSON, opts.JSON,
return runVaultikApp(cmd, mutating, opts.JSON,
"Failed to remove snapshot",
func(v *vaultik.Vaultik) error {
_, err := v.RemoveSnapshot(args[0], opts)
+2 -2
View File
@@ -99,8 +99,8 @@ type Snapshot struct {
StartedAt time.Time
CompletedAt *time.Time // nil if still in progress
FileCount int64
ChunkCount int64
BlobCount int64
ChunkCount int64 // Chunks this snapshot stored that were not stored before
BlobCount int64 // Blobs this snapshot created
TotalSize int64 // Total size of all referenced files
// BlobSize is the total size of all referenced blobs (compressed and
+28 -3
View File
@@ -127,6 +127,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
snapshotID string,
blobUncompressedSize int64,
compressionLevel int,
uploadBytes int64,
uploadDurationMs int64,
) error {
compressionRatio, err := r.extendedCompressionRatio(
@@ -141,7 +142,7 @@ func (r *SnapshotRepository) UpdateExtendedStats(
SET blob_uncompressed_size = ?,
compression_ratio = ?,
compression_level = ?,
upload_bytes = blob_size,
upload_bytes = ?,
upload_duration_ms = ?
WHERE id = ?
`
@@ -149,11 +150,11 @@ func (r *SnapshotRepository) UpdateExtendedStats(
if tx != nil {
_, err = tx.ExecContext(ctx, query,
blobUncompressedSize, compressionRatio, compressionLevel,
uploadDurationMs, snapshotID)
uploadBytes, uploadDurationMs, snapshotID)
} else {
_, err = r.db.ExecWithLog(ctx, query,
blobUncompressedSize, compressionRatio, compressionLevel,
uploadDurationMs, snapshotID)
uploadBytes, uploadDurationMs, snapshotID)
}
if err != nil {
@@ -543,6 +544,30 @@ func (r *SnapshotRepository) GetSnapshotTotalCompressedSize(
return totalSize, nil
}
// GetSnapshotBlobSizes returns the total compressed and uncompressed sizes
// of all blobs referenced by a snapshot.
func (r *SnapshotRepository) GetSnapshotBlobSizes(
ctx context.Context, snapshotID string,
) (int64, int64, error) {
query := `
SELECT COALESCE(SUM(b.compressed_size), 0),
COALESCE(SUM(b.uncompressed_size), 0)
FROM snapshot_blobs sb
JOIN blobs b ON sb.blob_hash = b.blob_hash
WHERE sb.snapshot_id = ?
`
var compressed, uncompressed int64
err := r.db.conn.QueryRowContext(ctx, query, snapshotID).Scan(
&compressed, &uncompressed)
if err != nil {
return 0, 0, fmt.Errorf("querying snapshot blob sizes: %w", err)
}
return compressed, uncompressed, nil
}
// GetSnapshotUncompressedChunkSize returns the sum of plaintext sizes of all unique
// chunks referenced by a snapshot (via snapshot_files → file_chunks → chunks).
func (r *SnapshotRepository) GetSnapshotUncompressedChunkSize(
+59
View File
@@ -145,6 +145,65 @@ func TestSnapshotRepositoryUpdateCounts(t *testing.T) {
}
}
// GetSnapshotBlobSizes totals the blobs the snapshot references, and only
// those.
func TestSnapshotRepositoryGetSnapshotBlobSizes(t *testing.T) {
t.Parallel()
db, cleanup := setupTestDB(t)
defer cleanup()
ctx := context.Background()
repos := database.NewRepositories(db)
snapshot := &database.Snapshot{
ID: "2024-01-03T12:00:00Z",
Hostname: testHostname,
VaultikVersion: testVersion,
StartedAt: time.Now().Truncate(time.Second),
}
err := repos.Snapshots.Create(ctx, nil, snapshot)
if err != nil {
t.Fatalf("failed to create snapshot: %v", err)
}
blobs := []*database.Blob{
{Hash: "referenced-1", CompressedSize: 10, UncompressedSize: 100},
{Hash: "referenced-2", CompressedSize: 20, UncompressedSize: 200},
{Hash: "unreferenced", CompressedSize: 40, UncompressedSize: 400},
}
for _, blob := range blobs {
blob.ID = types.NewBlobID()
blob.CreatedTS = time.Now().Truncate(time.Second)
err = repos.Blobs.Create(ctx, nil, blob)
if err != nil {
t.Fatalf("failed to create blob %s: %v", blob.Hash, err)
}
}
for _, blob := range blobs[:2] {
err = repos.Snapshots.AddBlob(ctx, nil, snapshot.ID.String(),
blob.ID, blob.Hash)
if err != nil {
t.Fatalf("failed to add blob %s to snapshot: %v", blob.Hash, err)
}
}
compressed, uncompressed, err := repos.Snapshots.GetSnapshotBlobSizes(
ctx, snapshot.ID.String())
if err != nil {
t.Fatalf("failed to get snapshot blob sizes: %v", err)
}
if compressed != 30 || uncompressed != 300 {
t.Errorf("blob sizes: got %d and %d, want 30 and 300",
compressed, uncompressed)
}
}
func TestSnapshotRepositoryListRecent(t *testing.T) {
t.Parallel()
-16
View File
@@ -158,19 +158,3 @@ type UploadStats struct {
MinDurationMs int64
MaxDurationMs int64
}
// GetCountBySnapshot returns the count of uploads for a specific snapshot
func (r *UploadRepository) GetCountBySnapshot(
ctx context.Context, snapshotID string,
) (int64, error) {
query := `SELECT COUNT(*) FROM uploads WHERE snapshot_id = ?`
var count int64
err := r.conn.QueryRowContext(ctx, query, snapshotID).Scan(&count)
if err != nil {
return 0, err
}
return count, nil
}
-4
View File
@@ -66,7 +66,6 @@ type ProgressStats struct {
BlobsCreated atomic.Int64
BlobsUploaded atomic.Int64
BytesUploaded atomic.Int64
UploadDurationMs atomic.Int64 // Total milliseconds spent uploading
CurrentFile atomic.Value // stores string
TotalSize atomic.Int64 // Total size to process (set after scan phase)
TotalFiles atomic.Int64 // Total files to process in phase 2
@@ -231,9 +230,6 @@ func (pr *ProgressReporter) ReportUploadComplete(
// Clear current upload
pr.stats.CurrentUpload.Store((*UploadInfo)(nil))
// Add to total upload duration
pr.stats.UploadDurationMs.Add(duration.Milliseconds())
// Calculate speed
if duration < time.Millisecond {
duration = time.Millisecond
+28 -52
View File
@@ -92,9 +92,6 @@ type Scanner struct {
// Mutex for coordinating blob creation
packerMu sync.Mutex // Blocks chunk production during blob creation
// Context for cancellation
scanCtx context.Context //nolint:containedctx // set per-Scan for packer callbacks
}
// Periodic status output intervals and thresholds for the scan and
@@ -134,7 +131,9 @@ type ScannerConfig struct {
SkipErrors bool
}
// ScanResult contains the results of a scan operation
// ScanResult contains the results of a scan operation. Files and bytes
// are counted per file: BytesScanned is the size of the new and changed
// files, BytesSkipped that of the unchanged ones.
type ScanResult struct {
FilesScanned int
FilesSkipped int
@@ -144,6 +143,9 @@ type ScanResult struct {
BytesDeleted int64
ChunksCreated int
BlobsCreated int
BlobsUploaded int
BytesUploaded int64
UploadDuration time.Duration
StartTime time.Time
EndTime time.Time
}
@@ -211,7 +213,6 @@ func (s *Scanner) Scan(
s.snapshotID = snapshotID
// Store source path for file records (used during restore)
s.currentSourcePath = path
s.scanCtx = ctx
result := &ScanResult{
StartTime: time.Now().UTC(),
}
@@ -219,7 +220,9 @@ func (s *Scanner) Scan(
// Set blob handler for concurrent upload
if s.storage != nil {
log.Debug("Setting blob handler for storage uploads")
s.packer.SetBlobHandler(s.handleBlobReady)
s.packer.SetBlobHandler(func(blobWithReader *blob.WithReader) error {
return s.handleBlobReady(ctx, blobWithReader, result)
})
} else {
log.Debug("No storage configured, blobs will not be uploaded")
}
@@ -288,8 +291,7 @@ func (s *Scanner) Scan(
log.Info("Phase 2/3: Skipping (no files need processing, metadata-only snapshot)")
}
// Finalize result with blob statistics
s.finalizeScanResult(ctx, result)
result.EndTime = time.Now().UTC()
return result, nil
}
@@ -431,27 +433,6 @@ func (s *Scanner) summarizeScanPhase(
s.ui.Completef("%s.", msg)
}
// finalizeScanResult populates final blob statistics in the scan result
// by querying the packer and database for blob/upload counts
func (s *Scanner) finalizeScanResult(ctx context.Context, result *ScanResult) {
blobs := s.packer.GetFinishedBlobs()
result.BlobsCreated += len(blobs)
// Query database for actual blob count created during this snapshot
// The database is authoritative, especially for concurrent blob uploads
// We count uploads rather than all snapshot_blobs to get only NEW blobs
if s.snapshotID != "" {
uploadCount, err := s.repos.Uploads.GetCountBySnapshot(ctx, s.snapshotID)
if err != nil {
log.Warn("Failed to query upload count from database", "error", err)
} else {
result.BlobsCreated = int(uploadCount)
}
}
result.EndTime = time.Now().UTC()
}
// loadKnownFiles loads the known files at and beneath path from the
// database into a map for fast lookup. Every loaded file the scan does
// not find is counted as deleted. This avoids per-file database queries
@@ -1511,24 +1492,24 @@ func (s *Scanner) finalizeProcessPhase(ctx context.Context, result *ScanResult)
}
// handleBlobReady is called by the packer when a blob is finalized
func (s *Scanner) handleBlobReady(blobWithReader *blob.WithReader) error {
func (s *Scanner) handleBlobReady(
ctx context.Context, blobWithReader *blob.WithReader, result *ScanResult,
) error {
startTime := time.Now().UTC()
finishedBlob := blobWithReader.FinishedBlob
result.BlobsCreated++
if s.progress != nil {
s.progress.ReportUploadStart(finishedBlob.Hash, finishedBlob.Compressed)
s.progress.GetStats().BlobsCreated.Add(1)
}
ctx := s.scanCtx
if ctx == nil {
ctx = context.Background()
}
blobPath := fmt.Sprintf("blobs/%s/%s/%s",
finishedBlob.Hash[:2], finishedBlob.Hash[2:4], finishedBlob.Hash)
blobExists, err := s.uploadBlobIfNeeded(ctx, blobPath, blobWithReader, startTime)
blobExists, err := s.uploadBlobIfNeeded(
ctx, blobPath, blobWithReader, startTime, result)
if err != nil {
s.cleanupBlobTempFile(blobWithReader)
@@ -1563,6 +1544,7 @@ func (s *Scanner) uploadBlobIfNeeded(
blobPath string,
blobWithReader *blob.WithReader,
startTime time.Time,
result *ScanResult,
) (bool, error) {
finishedBlob := blobWithReader.FinishedBlob
@@ -1598,6 +1580,10 @@ func (s *Scanner) uploadBlobIfNeeded(
uploadDuration := time.Since(startTime)
uploadSpeedBps := float64(finishedBlob.Compressed) / uploadDuration.Seconds()
result.BlobsUploaded++
result.BytesUploaded += finishedBlob.Compressed
result.UploadDuration += uploadDuration
s.ui.Completef("Uploaded blob %s (%s) in %s at %s.",
s.ui.Hex(finishedBlob.Hash),
s.ui.Size(finishedBlob.Compressed),
@@ -1813,9 +1799,9 @@ func (s *Scanner) processFileStreaming(
size: chunk.Size,
})
s.updateChunkStats(chunkExists, chunk.Size, result)
if !chunkExists {
s.updateChunkStats(chunk.Size, result)
err := s.addChunkToPacker(ctx, chunk)
if err != nil {
// Mark as a packer error so --skip-errors cannot swallow it:
@@ -1843,20 +1829,11 @@ func (s *Scanner) processFileStreaming(
return nil
}
// updateChunkStats updates scan result and progress stats for a processed chunk
func (s *Scanner) updateChunkStats(
chunkExists bool, chunkSize int64, result *ScanResult,
) {
if chunkExists {
result.FilesSkipped++
result.BytesSkipped += chunkSize
if s.progress != nil {
s.progress.GetStats().BytesSkipped.Add(chunkSize)
}
} else {
// updateChunkStats counts a chunk that was not already stored. The scan
// result's file counts, BytesScanned and BytesSkipped are not touched
// here: the scan phase counts each file once.
func (s *Scanner) updateChunkStats(chunkSize int64, result *ScanResult) {
result.ChunksCreated++
result.BytesScanned += chunkSize
if s.progress != nil {
s.progress.GetStats().ChunksCreated.Add(1)
@@ -1864,7 +1841,6 @@ func (s *Scanner) updateChunkStats(
s.progress.UpdateChunkingActivity()
}
}
}
// addChunkToPacker adds a chunk to the blob packer, finalizing the current
// blob if needed
+5 -23
View File
@@ -154,26 +154,6 @@ func (sm *SnapshotManager) CreateSnapshotWithName(
return snapshotID, nil
}
// UpdateSnapshotStats updates the statistics for a snapshot during backup
func (sm *SnapshotManager) UpdateSnapshotStats(
ctx context.Context, snapshotID string, stats BackupStats,
) error {
err := sm.repos.WithTx(ctx, func(ctx context.Context, tx *sql.Tx) error {
return sm.repos.Snapshots.UpdateCounts(ctx, tx, snapshotID,
int64(stats.FilesScanned),
int64(stats.ChunksCreated),
int64(stats.BlobsCreated),
stats.BytesScanned,
stats.BytesUploaded,
)
})
if err != nil {
return fmt.Errorf("updating snapshot stats: %w", err)
}
return nil
}
// UpdateSnapshotStatsExtended updates snapshot statistics with extended metrics.
// This includes compression level, uncompressed blob size, and upload duration.
func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
@@ -185,8 +165,8 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
int64(stats.FilesScanned),
int64(stats.ChunksCreated),
int64(stats.BlobsCreated),
stats.BytesScanned,
stats.BytesUploaded,
stats.TotalSize,
stats.BlobSize,
)
if err != nil {
return err
@@ -196,6 +176,7 @@ func (sm *SnapshotManager) UpdateSnapshotStatsExtended(
return sm.repos.Snapshots.UpdateExtendedStats(ctx, tx, snapshotID,
stats.BlobUncompressedSize,
stats.CompressionLevel,
stats.BytesUploaded,
stats.UploadDurationMs,
)
})
@@ -890,7 +871,7 @@ func (sm *SnapshotManager) getFileSize(path string) int64 {
// BackupStats contains statistics from a backup operation
type BackupStats struct {
FilesScanned int
BytesScanned int64
TotalSize int64 // Total size of all files examined
ChunksCreated int
BlobsCreated int
BytesUploaded int64
@@ -900,6 +881,7 @@ type BackupStats struct {
type ExtendedBackupStats struct {
BackupStats
BlobSize int64 // Total compressed size of all referenced blobs
BlobUncompressedSize int64 // Total uncompressed size of all referenced blobs
CompressionLevel int // Compression level used for this snapshot
UploadDurationMs int64 // Total milliseconds spent uploading to S3
+52 -59
View File
@@ -189,6 +189,11 @@ type snapshotStats struct {
totalBytesUploaded int64
totalBlobsUploaded int
uploadDuration time.Duration
// The sizes of all blobs the snapshot references, set by
// finalizeSnapshotMetadata once snapshot_blobs is populated.
blobSize int64
blobUncompressedSize int64
}
// createNamedSnapshot creates a single named snapshot
@@ -228,8 +233,6 @@ func (v *Vaultik) createNamedSnapshot(
return err
}
v.collectUploadStats(scanner, stats)
err = v.finalizeSnapshotMetadata(snapshotID, stats)
if err != nil {
return err
@@ -314,6 +317,9 @@ func (v *Vaultik) scanAllDirectories(
stats.totalBytesSkipped += result.BytesSkipped
stats.totalFilesDeleted += result.FilesDeleted
stats.totalBytesDeleted += result.BytesDeleted
stats.totalBlobsUploaded += result.BlobsUploaded
stats.totalBytesUploaded += result.BytesUploaded
stats.uploadDuration += result.UploadDuration
log.Info("Directory scan complete",
"path", dir,
@@ -329,18 +335,6 @@ func (v *Vaultik) scanAllDirectories(
return stats, nil
}
// collectUploadStats gathers upload statistics from the scanner's
// progress reporter.
func (v *Vaultik) collectUploadStats(scanner *snapshot.Scanner, stats *snapshotStats) {
if s := scanner.GetProgress(); s != nil {
progressStats := s.GetStats()
stats.totalBytesUploaded = progressStats.BytesUploaded.Load()
stats.totalBlobsUploaded = int(progressStats.BlobsUploaded.Load())
stats.uploadDuration = time.Duration(
progressStats.UploadDurationMs.Load()) * time.Millisecond
}
}
// finalizeSnapshotMetadata updates stats, exports metadata, and only then
// marks the snapshot complete. Recording completion last is deliberate: an
// export interrupted by a crash leaves the snapshot incomplete rather than
@@ -350,31 +344,39 @@ func (v *Vaultik) collectUploadStats(scanner *snapshot.Scanner, stats *snapshotS
func (v *Vaultik) finalizeSnapshotMetadata(
snapshotID string, stats *snapshotStats,
) error {
// snapshot_blobs must be populated before the blob sizes below, which
// total the snapshot's blobs, and before the export, which builds the
// manifest and the trimmed metadata database from it.
err := v.SnapshotManager.PopulateSnapshotBlobs(v.ctx, snapshotID)
if err != nil {
return fmt.Errorf("populating snapshot blobs: %w", err)
}
stats.blobSize, stats.blobUncompressedSize, err =
v.Repositories.Snapshots.GetSnapshotBlobSizes(v.ctx, snapshotID)
if err != nil {
return fmt.Errorf("getting snapshot blob sizes: %w", err)
}
extStats := snapshot.ExtendedBackupStats{
BackupStats: snapshot.BackupStats{
FilesScanned: stats.totalFiles,
BytesScanned: stats.totalBytes,
TotalSize: stats.totalBytes + stats.totalBytesSkipped,
ChunksCreated: stats.totalChunks,
BlobsCreated: stats.totalBlobs,
BytesUploaded: stats.totalBytesUploaded,
},
BlobUncompressedSize: 0,
BlobSize: stats.blobSize,
BlobUncompressedSize: stats.blobUncompressedSize,
CompressionLevel: v.Config.CompressionLevel,
UploadDurationMs: stats.uploadDuration.Milliseconds(),
}
err := v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
err = v.SnapshotManager.UpdateSnapshotStatsExtended(v.ctx, snapshotID, extStats)
if err != nil {
return fmt.Errorf("updating snapshot stats: %w", err)
}
// snapshot_blobs must be populated before the export, which builds the
// manifest and the trimmed metadata database from it.
err = v.SnapshotManager.PopulateSnapshotBlobs(v.ctx, snapshotID)
if err != nil {
return fmt.Errorf("populating snapshot blobs: %w", err)
}
err = v.SnapshotManager.ExportSnapshotMetadata(
v.ctx, v.Config.IndexPath, snapshotID)
if err != nil {
@@ -409,12 +411,10 @@ func (v *Vaultik) printSnapshotSummary(
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
// Get total blob sizes from database
compressedSize, uncompressedSize := v.getSnapshotBlobSizes(snapshotID)
var compressionRatio float64
if uncompressedSize > 0 {
compressionRatio = float64(compressedSize) / float64(uncompressedSize)
if stats.blobUncompressedSize > 0 {
compressionRatio = float64(stats.blobSize) /
float64(stats.blobUncompressedSize)
} else {
compressionRatio = 1.0
}
@@ -442,8 +442,8 @@ func (v *Vaultik) printSnapshotSummary(
if stats.totalBlobsUploaded > 0 {
v.UI.Detailf("Storage: %s compressed from %s (%.2fx ratio).",
v.UI.Size(compressedSize),
v.UI.Size(uncompressedSize),
v.UI.Size(stats.blobSize),
v.UI.Size(stats.blobUncompressedSize),
compressionRatio)
v.UI.Detailf("Upload: %d blobs, %s in %s (%s).",
stats.totalBlobsUploaded,
@@ -455,27 +455,6 @@ func (v *Vaultik) printSnapshotSummary(
v.UI.Detailf("Snapshot create duration: %s.", v.UI.Duration(snapshotDuration))
}
// getSnapshotBlobSizes returns total compressed and uncompressed blob
// sizes for a snapshot.
func (v *Vaultik) getSnapshotBlobSizes(snapshotID string) (int64, int64) {
var compressed, uncompressed int64
blobHashes, err := v.Repositories.Snapshots.GetBlobHashes(v.ctx, snapshotID)
if err != nil {
return 0, 0
}
for _, hash := range blobHashes {
blob, err := v.Repositories.Blobs.GetByHash(v.ctx, hash)
if err == nil && blob != nil {
compressed += blob.CompressedSize
uncompressed += blob.UncompressedSize
}
}
return compressed, uncompressed
}
// SnapshotPurgeOptions contains options for the snapshot purge command.
type SnapshotPurgeOptions struct {
KeepLatest bool // Keep only the most recent snapshot per name
@@ -732,7 +711,7 @@ func (v *Vaultik) VerifySnapshotWithOptions(
result.BlobCount = manifest.BlobCount
result.TotalSize = manifest.TotalCompressedSize
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Snapshot information:\n")
v.stdoutf(" Blob count: %d\n", manifest.BlobCount)
v.stdoutf(" Total size: %s\n", ubytes(manifest.TotalCompressedSize))
@@ -787,7 +766,7 @@ func (v *Vaultik) printVerifyHeader(snapshotID string, opts *VerifyOptions) {
snapshotTime = t
}
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Verifying snapshot %s\n", snapshotID)
if !snapshotTime.IsZero() {
@@ -827,7 +806,7 @@ func (v *Vaultik) verifyManifestBlobs(
stat, err := v.Storage.Stat(v.ctx, blobPath)
switch {
case err != nil:
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf(" Missing: %s (%s)\n",
blob.Hash, ubytes(blob.CompressedSize))
}
@@ -835,7 +814,7 @@ func (v *Vaultik) verifyManifestBlobs(
missing++
missingSize += blob.CompressedSize
case stat.Size != blob.CompressedSize:
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf(" Wrong size: %s (store has %s, manifest lists %s)\n",
blob.Hash, ubytes(stat.Size), ubytes(blob.CompressedSize))
}
@@ -867,6 +846,22 @@ func (v *Vaultik) formatVerifyResult(
return v.outputVerifyJSON(result)
}
// Under --quiet a failure is still returned, and the cli layer
// prints it on stderr.
if !v.UI.Quiet() {
v.printVerifySummary(result, failure)
}
if failure != "" {
return fmt.Errorf("%w: %s", errSnapshotVerifyFailed, failure)
}
return nil
}
// printVerifySummary prints the counts and the status line that end the
// human-readable shallow verify report. failure is empty when it passed.
func (v *Vaultik) printVerifySummary(result *VerifyResult, failure string) {
v.stdoutf("\nVerification complete:\n")
v.stdoutf(" Present with listed size: %d blobs\n", result.Verified)
@@ -888,14 +883,12 @@ func (v *Vaultik) formatVerifyResult(
if failure != "" {
v.stdoutf("FAILED - %s\n", failure)
return fmt.Errorf("%w: %s", errSnapshotVerifyFailed, failure)
return
}
// Report only what was actually checked: presence and size, not contents.
v.stdoutf("OK - all %d blobs listed in the manifest are present with the "+
"listed size; contents not checked (use --deep)\n", result.Verified)
return nil
}
// shallowVerifyFailure returns a human-readable description of everything
+285
View File
@@ -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)
}
+11 -8
View File
@@ -103,7 +103,7 @@ func (v *Vaultik) RunDeepVerify(snapshotID string, opts *VerifyOptions) error {
log.Info("Starting snapshot verification", "snapshot_id", snapshotID, "mode", "deep")
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Deep verification of snapshot: %s\n\n", snapshotID)
}
@@ -143,10 +143,13 @@ func (v *Vaultik) RunDeepVerify(snapshotID string, opts *VerifyOptions) error {
log.Info("✓ Verification completed successfully",
"snapshot_id", snapshotID, "mode", "deep", "blobs_verified", len(dbBlobs))
if !v.UI.Quiet() {
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
}
@@ -170,7 +173,7 @@ func (v *Vaultik) loadVerificationData(
// remote manifests; see its doc comment.
log.Info("Downloading manifest", "remote_key", remoteKey)
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Downloading manifest...\n")
}
@@ -185,7 +188,7 @@ func (v *Vaultik) loadVerificationData(
"manifest_blob_count", manifest.BlobCount,
"manifest_total_size", ubytes(manifest.TotalCompressedSize))
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Manifest loaded: %d blobs (%s)\n",
manifest.BlobCount, ubytes(manifest.TotalCompressedSize))
v.stdoutf("Downloading and decrypting database...\n")
@@ -215,7 +218,7 @@ func (v *Vaultik) loadVerificationData(
"db_blob_count", len(dbBlobs),
"db_total_size", ubytes(dbTotalSize))
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Database loaded: %d blobs (%s)\n",
len(dbBlobs), ubytes(dbTotalSize))
}
@@ -273,7 +276,7 @@ func (v *Vaultik) runVerificationSteps(
totalSize int64,
identities []age.Identity,
) error {
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Verifying manifest against database...\n")
}
@@ -282,7 +285,7 @@ func (v *Vaultik) runVerificationSteps(
return v.deepVerifyFailure(result, opts, err.Error(), err)
}
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("Manifest verified.\n")
v.stdoutf("Checking blob existence in remote storage...\n")
}
@@ -292,7 +295,7 @@ func (v *Vaultik) runVerificationSteps(
return v.deepVerifyFailure(result, opts, err.Error(), err)
}
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf("All blobs exist.\n")
v.stdoutf("Downloading and verifying blob contents (%d blobs, %s)...\n",
len(dbBlobs), ubytes(totalSize))
@@ -748,7 +751,7 @@ func (v *Vaultik) performDeepVerificationFromDB(
"eta", eta.Round(time.Second),
)
if !opts.JSON {
if !opts.JSON && !v.UI.Quiet() {
v.stdoutf(" Verified %d/%d blobs (%d remaining) - %s/%s - elapsed %s, eta %s\n",
i+1, len(blobs), remaining,
ubytes(bytesProcessed),
+102
View File
@@ -0,0 +1,102 @@
package vaultik_test
import (
"bytes"
"context"
"io"
"os"
"path/filepath"
"testing"
"github.com/spf13/afero"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/snapshot"
"sneak.berlin/go/vaultik/internal/ui"
"sneak.berlin/go/vaultik/internal/vaultik"
)
// TestVerify_QuietSuppressesReport is the --quiet contract for
// `snapshot verify`: neither shallow nor deep verify writes its report,
// a failed verify still returns its error (which the cli layer prints
// on stderr), and the --json document still emits.
func TestVerify_QuietSuppressesReport(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs()
tempDir := t.TempDir()
dataDir := filepath.Join(tempDir, "source")
storeDir := filepath.Join(tempDir, "remote")
dbPath := filepath.Join(tempDir, "index.sqlite")
chunkSize := int64(32 * 1024)
maxBlobSize := int64(128 * 1024)
require.NoError(t, fs.MkdirAll(dataDir, 0o755))
require.NoError(t, afero.WriteFile(fs,
filepath.Join(dataDir, "data.bin"),
bytesPattern("quiet-", int(maxBlobSize*2)), 0o644))
ctx := context.Background()
cfg, storer, snapshotID := runFileStorageBackup(
ctx, t, fs, dataDir, storeDir, dbPath, chunkSize, maxBlobSize)
// The UI writes to the same buffer as Stdout, as both write to the
// process's stdout in production.
var stdout bytes.Buffer
newQuietVerifier := func() *vaultik.Vaultik {
v := &vaultik.Vaultik{
Config: cfg,
Storage: storer,
Fs: fs,
Stdout: &stdout,
Stderr: io.Discard,
UI: ui.NewWithColor(&stdout, false),
}
v.SetContext(ctx)
v.UI.SetQuiet(true)
return v
}
require.NoError(t, newQuietVerifier().VerifySnapshotWithOptions(
snapshotID, &vaultik.VerifyOptions{}))
require.Empty(t, stdout.String(),
"shallow verify must write no report under --quiet")
require.NoError(t, newQuietVerifier().VerifySnapshotWithOptions(
snapshotID, &vaultik.VerifyOptions{Deep: true}))
require.Empty(t, stdout.String(),
"deep verify must write no report under --quiet")
require.NoError(t, newQuietVerifier().VerifySnapshotWithOptions(
snapshotID, &vaultik.VerifyOptions{JSON: true}))
require.Equal(t, "ok", decodeVerifyResult(t, stdout.Bytes()).Status,
"the --json document must still emit under --quiet")
// A snapshot without its encrypted database fails shallow verify. A
// failed report also lists each missing blob and each blob of the
// wrong size, so remove one blob and grow another.
require.NoError(t, os.Remove(filepath.Join(storeDir, "metadata",
snapshot.RemoteSnapshotKey(snapshotID), "db.zst.age")))
blobFiles, err := filepath.Glob(
filepath.Join(storeDir, "blobs", "*", "*", "*"))
require.NoError(t, err)
require.GreaterOrEqual(t, len(blobFiles), 2,
"the snapshot must span two blobs, one to remove and one to grow")
require.NoError(t, os.Remove(blobFiles[0]))
growOneBlob(t, fs, filepath.Join(storeDir, "blobs"))
stdout.Reset()
require.Error(t, newQuietVerifier().VerifySnapshotWithOptions(
snapshotID, &vaultik.VerifyOptions{}),
"--quiet must not change the outcome of a failed verify")
require.Empty(t, stdout.String(),
"a failed verify must write no report under --quiet")
}