5 Commits
Author SHA1 Message Date
clawbot a50e3fa038 Add tests for internal/storage: URL parsing, backends, shared conformance suite (closes #66)
check / check (pull_request) In progress
check / check (push) Failing after 1s
internal/storage, the package that parses store URLs and selects the backend, had no tests.

Adds table-driven tests for URL parsing (each scheme, query parameters, malformed input, unknown scheme, backend type chosen); one shared conformance suite for the Storer interface, run against the file backend in a temp directory and the s3 backend on the in-process harness internal/s3 already uses, so a new backend inherits it; and rclone construction and argument tests using its in-process local backend. A comment records that rclone data operations need a configured remote and are not unit-tested. No production code changed and no defect surfaced.

model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)
2026-09-22 00:58:27 +02:00
clawbot 6fcd8e1668 Stamp Docker image version from the host; flush profiles on error exit (closes #75)
check / check (push) Failing after 1s
check / check (pull_request) Failing after 1s
Docker images reported commit unknown because the build ran git inside the container while .dockerignore excludes .git, and VERSION was never overridden. script/docker and script/cibuild now compute version, commit and date on the host and pass them as build args; the Dockerfile runs no git and falls back to dev and unknown, never empty, on a bare docker build.

Profiling a failing command gave a truncated or missing profile: Entry and each command goroutine called os.Exit(1), skipping the deferred profile writers in main. Entry now returns a status that main exits with after its defers run, and command goroutines report failure through one RunOperation helper, which also restores PID-lock release and graceful shutdown on failure.

model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)

Co-authored-by: clawbot <clawbot@noreply.example.org>
2026-09-21 22:01:05 +02:00
clawbot aab6a87f8c Reconcile the schema/migration docs with the code (closes #68)
check / check (pull_request) Failing after 1s
check / check (push) Successful in 2m46s
Four documents told different stories about the database schema. docs/DATAMODEL.md now owns the explanation and separates two things: the policy, which is unchanged (no supported upgrade path between versions; delete the local index with vaultik database delete and back up again), and the schema bootstrap that does exist (numbered files in internal/database/schema applied to a fresh database and recorded in schema_migrations).

README.md and AGENTS.md are reworded to match and link there. AGENTS.md names the real file to edit, internal/database/schema/001.sql, and notes that the pre-1.0 disposability clause expires on tagging. No code changed.

Judgement call: REPO_POLICIES.md still names a different schema file; it is cross-project policy and was left alone.

model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)
2026-09-21 21:58:37 +02:00
clawbot c355ef4d25 Report a prune count that could not be read as unknown, not 0 (closes #96)
check / check (pull_request) Failing after 1s
check / check (push) Successful in 2m46s
Prune read table row counts before and after to report how many orphaned files, chunks and blobs it removed, and discarded the error from every read. A failed query therefore reported as a count of 0, and the summary showed plausible wrong numbers.

A count that cannot be read is now logged as a warning (on stderr, also under --json) and shown as "unknown"; a difference computed from an unknown count is itself unknown. 0 still means the table was empty. No --json document carries these counts, so none can show a false 0.

model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)
2026-09-21 21:41:57 +02:00
clawbot 5927e1aa3d Write file:// blobs atomically via temp file and rename (closes #130)
check / check (push) Failing after 0s
check / check (pull_request) Failing after 0s
The file:// backend streamed each object straight to its final key, so an upload cut off mid-stream left a truncated object there. The next backup saw that Stat succeeded, recorded the blob as complete, and produced a snapshot that reported success but could not be restored.

Writes now go to a temporary file with a .partial suffix in the destination directory, are synced, then renamed onto the key. List and ListStream skip .partial files, so a leftover is never trusted as a blob and is overwritten when the key is written again. S3 PutObject is already atomic.

Disclosure: the containing directory is not synced after the rename, so a host crash right after it could still lose the object on some filesystems.

model: claude-opus-4-8 (implementation, review); claude-fable-5-1 (merge)
2026-09-21 21:24:42 +02:00
28 changed files with 1102 additions and 627 deletions
+8 -3
View File
@@ -104,7 +104,12 @@ Version: 2025-06-08
13. Pre-1.0: NEVER write database migrations. There are no live databases 13. Pre-1.0: NEVER write database migrations. There are no live databases
anywhere — every user's local index can be rebuilt from a fresh full anywhere — every user's local index can be rebuilt from a fresh full
backup. When the schema changes, just change `schema.sql` (and any code backup. To change the schema, edit `internal/database/schema/001.sql`
that touches the affected tables). The local index is disposable until (and any code that touches the affected tables) directly; do not add new
1.0 ships and is tagged. numbered schema files. Those numbered files and the `schema_migrations`
table they populate only bootstrap a fresh database — they are not an
upgrade path. The local index is disposable until 1.0 ships and is
tagged; once 1.0 is tagged that clause expires and the question of
upgrading existing indexes returns. See [`docs/DATAMODEL.md`](docs/DATAMODEL.md)
for the full explanation.
+24 -3
View File
@@ -20,8 +20,6 @@
# golang:1.26.1-alpine, 2026-03-17 # golang:1.26.1-alpine, 2026-03-17
FROM golang:1.26.1-alpine@sha256:2389ebfa5b7f43eeafbd6be0c3700cc46690ef842ad962f6c5bd6be49ed82039 AS builder FROM golang:1.26.1-alpine@sha256:2389ebfa5b7f43eeafbd6be0c3700cc46690ef842ad962f6c5bd6be49ed82039 AS builder
ARG VERSION=dev
# Build tooling: make, plus a C toolchain because `go test -race` needs cgo. # Build tooling: make, plus a C toolchain because `go test -race` needs cgo.
# The sqlite driver is pure Go (modernc.org/sqlite), so no sqlite library or # The sqlite driver is pure Go (modernc.org/sqlite), so no sqlite library or
# CLI is required. # CLI is required.
@@ -66,8 +64,31 @@ RUN [ -n "$CHECK_EPOCH" ] || exit 1
RUN echo "check epoch: ${CHECK_EPOCH}" && make fmt-check RUN echo "check epoch: ${CHECK_EPOCH}" && make fmt-check
RUN echo "check epoch: ${CHECK_EPOCH}" && make test RUN echo "check epoch: ${CHECK_EPOCH}" && make test
# Version, commit and build date are computed on the host by
# script/docker and script/cibuild (where .git exists) and passed in as
# build args. The build context excludes .git (see .dockerignore), so
# the build cannot derive them itself: it used to try, with `git
# rev-parse` inside this stage, and always got "unknown". VERSION comes
# from script/version, the source of truth shared with the Makefile, so
# it carries the same tag / dev-<sha> / -dirty rules and a Docker image
# reports the same string a local build of the same tree would.
#
# The defaults are the fallback for a bare `docker build .` that passes
# none of them: an unset arg would otherwise stamp an empty string and
# produce an image that cannot report its own version, commit or date.
# They match what an out-of-git build reports elsewhere.
#
# These ARGs sit here, after the checks, rather than at the top of the
# stage: every commit changes their values, and a value change
# invalidates all layers below the ARG. Declared up top they would bust
# `go mod download`; here they only rekey this build layer, which the
# COPY of the sources above already rebuilds on any change anyway.
ARG VERSION=dev
ARG COMMIT=unknown
ARG COMMIT_DATE=unknown
# Build (pure Go, no CGO required since we use modernc.org/sqlite) # Build (pure Go, no CGO required since we use modernc.org/sqlite)
RUN CGO_ENABLED=0 go build -ldflags "-X 'sneak.berlin/go/vaultik/internal/globals.Version=${VERSION}' -X 'sneak.berlin/go/vaultik/internal/globals.Commit=$(git rev-parse HEAD 2>/dev/null || echo unknown)' -X 'sneak.berlin/go/vaultik/internal/globals.CommitDate=$(git show -s --format=%cs HEAD 2>/dev/null || echo unknown)'" -o /vaultik ./cmd/vaultik RUN CGO_ENABLED=0 go build -ldflags "-X 'sneak.berlin/go/vaultik/internal/globals.Version=${VERSION}' -X 'sneak.berlin/go/vaultik/internal/globals.Commit=${COMMIT}' -X 'sneak.berlin/go/vaultik/internal/globals.CommitDate=${COMMIT_DATE}'" -o /vaultik ./cmd/vaultik
# Runtime stage # Runtime stage
# alpine:3.21, 2026-02-25 # alpine:3.21, 2026-02-25
+13 -7
View File
@@ -457,9 +457,13 @@ Key fields:
sequentially. Restore speed is bound by single-stream throughput. sequentially. Restore speed is bound by single-stream throughput.
* **Device nodes, named pipes, and sockets are silently skipped.** Only * **Device nodes, named pipes, and sockets are silently skipped.** Only
regular files, directories, and symlinks are backed up. regular files, directories, and symlinks are backed up.
* **No database migrations.** If the local SQLite schema changes between * **No upgrade path between versions.** There is no supported way to carry
versions, delete the local database (`vaultik database delete`) and run an existing local index across a schema change; if the local SQLite
a full backup. Remote storage is unaffected. schema changes between versions, delete the local database (`vaultik
database delete`) and run a full backup. Remote storage is unaffected.
(The binary does embed numbered schema files and a `schema_migrations`
table to bootstrap a fresh database — see [`docs/DATAMODEL.md`](docs/DATAMODEL.md)
— but that is not an upgrade path.)
* **Files that change during backup may be inconsistent.** There is no * **Files that change during backup may be inconsistent.** There is no
filesystem snapshot or freeze. If a file is modified between the scan filesystem snapshot or freeze. If a file is modified between the scan
and chunk phases, the backed-up copy may reflect a partial write. and chunk phases, the backed-up copy may reflect a partial write.
@@ -529,10 +533,12 @@ priority.
another host" workflow works but isn't documented as a another host" workflow works but isn't documented as a
first-class operation in this README. Worth a dedicated section first-class operation in this README. Worth a dedicated section
once it's settled. once it's settled.
* **Schema migrations.** Currently nonexistent — pre-1.0 schema * **Cross-version schema upgrades.** There is no upgrade path between
changes are handled by `vaultik database delete` plus a full released versions — pre-1.0 schema changes are handled by `vaultik
re-scan. Post-1.0 we'll need a migration story to keep existing database delete` plus a full re-scan (see
index databases usable across upgrades. [`docs/DATAMODEL.md`](docs/DATAMODEL.md)). Post-1.0 we'll need a
migration story to keep existing index databases usable across
upgrades.
* **Storage backend coverage tests.** S3, file://, and rclone:// * **Storage backend coverage tests.** S3, file://, and rclone://
all share the Storer interface but the rclone path is the least all share the Storer interface but the rclone path is the least
exercised in CI. exercised in CI.
+10 -1
View File
@@ -25,6 +25,16 @@ release" is exactly the contradiction
# Completed Steps # Completed Steps
- 2026-09-21: Stopped `prune` from reporting a failed row count as 0
([issue #96](https://git.eeqj.de/sneak/vaultik/issues/96)). The seven
`getTableCount` reads in `PruneDatabase` discarded their error, so a
query that could not run became a plausible `0` and the before/after
delta computed from it looked like real work. Each read now logs at
warn on failure and renders as `unknown`, never `0`, so an empty table
is distinguishable from one that could not be queried. The counts have
no `--json` representation — under `--json` the summary is suppressed
entirely — so nothing there can show a false `0`.
- 2026-09-21: Made the s3 storage backend report a missing object as - 2026-09-21: Made the s3 storage backend report a missing object as
`storage.ErrNotFound`, like the `file` and `rclone` backends and as the `storage.ErrNotFound`, like the `file` and `rclone` backends and as the
`Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw `Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw
@@ -33,7 +43,6 @@ release" is exactly the contradiction
helper (reused by `HeadObject`) and a test that a missing key maps to helper (reused by `HeadObject`) and a test that a missing key maps to
`ErrNotFound` `ErrNotFound`
([issue #129](https://git.eeqj.de/sneak/vaultik/issues/129)). ([issue #129](https://git.eeqj.de/sneak/vaultik/issues/129)).
- 2026-09-21: Fixed `verify --deep` reporting healthy snapshots as - 2026-09-21: Fixed `verify --deep` reporting healthy snapshots as
corrupt. Its final blob-integrity check hashed the encrypted corrupt. Its final blob-integrity check hashed the encrypted
downloaded bytes with a single SHA256 and compared that to the blob downloaded bytes with a single SHA256 and compared that to the blob
+102
View File
@@ -0,0 +1,102 @@
package main_test
import (
"strings"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
// This file guards the version stamping of the product image (issue
// #75). The failure it protects against is silent: the image still
// builds and runs, but `vaultik version` inside it reports "commit:
// unknown", so an operator cannot tell which source produced a given
// backup. .dockerignore excludes .git, so the build cannot derive the
// commit itself; the values must be computed on the host and passed in.
//
// These are parses of the committed files, for the same reason the lint
// guards next door are: shelling out to docker would nest a build
// inside `make test`. That `vaultik version` in the built image really
// prints the host's version is verified by hand and recorded on the
// pull request.
// dockerScript is script/docker, relative to the repository root.
const dockerScript = "script/docker"
// versionArgs are the ldflag targets the build stamps and, matching
// them, the build args the host must supply. The names line up so the
// same list checks both files.
func versionArgs() []string {
return []string{"VERSION", "COMMIT", "COMMIT_DATE"}
}
// TestProductDockerfileTakesVersionAsBuildArgs fails unless the build
// declares each version arg and stamps it into the binary by ldflag
// reference, rather than computing it in the container.
func TestProductDockerfileTakesVersionAsBuildArgs(t *testing.T) {
t.Parallel()
found := instructions(t, productDockerfile)
for _, arg := range versionArgs() {
require.GreaterOrEqual(t, indexOf(found, "ARG "+arg), 0,
"%s must declare `ARG %s` so the host can pass it in",
productDockerfile, arg)
assertLdflagReferences(t, found, arg)
}
}
// TestProductDockerfileDoesNotDeriveVersionItself is the anti-regression
// for the original defect: the container ran `git rev-parse`, but .git
// is not in the build context, so it always resolved to "unknown". No
// git command may reach into a build that cannot see the history.
func TestProductDockerfileDoesNotDeriveVersionItself(t *testing.T) {
t.Parallel()
text := instructionText(readRepoFile(t, productDockerfile))
assert.NotContains(t, text, "git ",
"%s must not run git: .git is excluded from the build context, so"+
" any value it derives is wrong. Pass version, commit and date"+
" in as build args instead.", productDockerfile)
}
// TestDockerScriptComputesVersionOnTheHost fails unless script/docker
// derives each value where .git exists and passes it as a build arg,
// with VERSION coming from script/version so a Docker build reports the
// same string a local build of the same tree would.
func TestDockerScriptComputesVersionOnTheHost(t *testing.T) {
t.Parallel()
script := readRepoFile(t, dockerScript)
for _, arg := range versionArgs() {
assert.Contains(t, script, "--build-arg "+arg+"=",
"%s must pass --build-arg %s to the build", dockerScript, arg)
}
assert.Contains(t, script, "/version",
"%s must take VERSION from script/version, the source of truth"+
" shared with the Makefile", dockerScript)
}
// assertLdflagReferences fails unless some build instruction stamps the
// named variable from the ARG (a ${arg} reference), not from a value
// computed inside the container.
func assertLdflagReferences(t *testing.T, found []string, arg string) {
t.Helper()
for _, instruction := range found {
if strings.HasPrefix(instruction, "RUN ") &&
strings.Contains(instruction, "go build") &&
strings.Contains(instruction, "${"+arg+"}") {
return
}
}
assert.Fail(t, "version arg is declared but never stamped",
"the go build in %s must reference ${%s} in its ldflags, or the"+
" arg is passed and discarded", productDockerfile, arg)
}
+6 -2
View File
@@ -304,10 +304,14 @@ func instructionText(contents string) string {
} }
// indexOf returns the position of the first instruction equal to, or // indexOf returns the position of the first instruction equal to, or
// beginning with, want; -1 if there is none. // beginning with, want; -1 if there is none. An `ARG NAME=default`
// counts as beginning with `ARG NAME`, so a declared arg is found
// whether or not it carries a default.
func indexOf(found []string, want string) int { func indexOf(found []string, want string) int {
for i, instruction := range found { for i, instruction := range found {
if instruction == want || strings.HasPrefix(instruction, want+" ") { if instruction == want ||
strings.HasPrefix(instruction, want+" ") ||
strings.HasPrefix(instruction, want+"=") {
return i return i
} }
} }
+11 -1
View File
@@ -10,6 +10,16 @@ import (
) )
func main() { func main() {
os.Exit(run())
}
// run sets up optional profiling, runs the CLI, and returns the process
// exit code. os.Exit lives in main so it fires only after run's deferred
// profile writers have flushed. cli.Entry returns a status code rather
// than calling os.Exit itself: an os.Exit from inside it would skip
// these defers and truncate the profile of a failing command -- exactly
// the command one most often wants to profile.
func run() int {
// CPU profiling: set VAULTIK_CPUPROFILE=/path/to/cpu.prof // CPU profiling: set VAULTIK_CPUPROFILE=/path/to/cpu.prof
if cpuProfile := os.Getenv("VAULTIK_CPUPROFILE"); cpuProfile != "" { if cpuProfile := os.Getenv("VAULTIK_CPUPROFILE"); cpuProfile != "" {
f, err := os.Create(cpuProfile) //nolint:gosec // G304: operator-set path f, err := os.Create(cpuProfile) //nolint:gosec // G304: operator-set path
@@ -46,5 +56,5 @@ func main() {
}() }()
} }
cli.Entry() return cli.Entry()
} }
+24 -5
View File
@@ -5,11 +5,30 @@
Vaultik uses a local SQLite database to track file metadata, chunk mappings, and blob associations during the backup process. This database serves as an index for incremental backups and enables efficient deduplication. Vaultik uses a local SQLite database to track file metadata, chunk mappings, and blob associations during the backup process. This database serves as an index for incremental backups and enables efficient deduplication.
**Important Notes:** **Important Notes:**
- **No Migration Support (pre-1.0)**: Vaultik does not support database schema
migrations. The local index is treated as disposable — if the schema changes, This section is the authoritative explanation of the schema/migration story;
delete the local SQLite database (`vaultik database delete`) and run a full other documents (the README and `AGENTS.md`) link here.
backup. The remote storage is unaffected; the new index will re-deduplicate
against existing remote blobs. - **No upgrade path between versions (pre-1.0)**: Vaultik has no supported way to
carry an existing local index across a schema change. The index is disposable
— if the on-disk schema changes between versions, delete the local SQLite
database (`vaultik database delete`) and run a full backup. Remote storage is
unaffected; the new index re-deduplicates against existing remote blobs. This
is the standing project policy, and it is separate from the schema bootstrap
described next.
- **Schema bootstrap**: a fresh database is populated from numbered SQL files
embedded in the binary under `internal/database/schema/`. `000.sql` creates the
`schema_migrations` table; `001.sql` creates the application tables. On opening
a database the code applies each numbered file that has not yet run and records
its version in `schema_migrations`. This bootstraps a new database; it does not
upgrade an existing one between released versions.
- **Changing the schema (pre-1.0)**: edit `internal/database/schema/001.sql` (and
the code that touches the affected tables) directly. Do not add new numbered
files — there is no installed base to migrate.
- **Disposability expires at 1.0**: the index is treated as disposable only until
1.0 ships and is tagged. Once 1.0 is tagged that clause expires and the
question of upgrading existing indexes returns. It is deliberately left open
here.
- **Version Compatibility**: In rare cases, you may need to use the same version - **Version Compatibility**: In rare cases, you may need to use the same version
of Vaultik to restore a backup as was used to create it. This ensures of Vaultik to restore a backup as was used to create it. This ensures
compatibility with the metadata format stored in S3. compatibility with the metadata format stored in S3.
+90 -39
View File
@@ -11,6 +11,7 @@ import (
"os/signal" "os/signal"
"path/filepath" "path/filepath"
"strings" "strings"
"sync"
"syscall" "syscall"
"time" "time"
@@ -196,13 +197,90 @@ 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.
var errReported = errors.New("operation failed")
// RunOperation runs op against the Vaultik instance inside the fx app
// and turns a failure into a returned error rather than an os.Exit from
// within the goroutine. An os.Exit there skipped main's deferred
// profile writers -- so profiling a failing command yielded a truncated
// profile (issue #75) -- and RunWithApp's PID-lock release, and denied
// the app any graceful shutdown; returning the error to the top runs
// all three.
//
// op runs in a goroutine so OnStart returns promptly and an interrupt
// can still cancel through OnStop; when it finishes, success or failure,
// it triggers shutdown, which is what lets RunWithApp return. 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.
func RunOperation(
ctx context.Context, opts AppOptions,
op func(v *vaultik.Vaultik) error, report func(err error),
) error {
var (
mu sync.Mutex
failed bool
)
opts.Invokes = append(opts.Invokes,
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
lc.Append(fx.Hook{
OnStart: func(_ context.Context) error {
go func() {
err := op(v)
if err != nil && !errors.Is(err, context.Canceled) {
report(err)
mu.Lock()
failed = true
mu.Unlock()
}
stopErr := v.Shutdowner.Shutdown()
if stopErr != nil {
log.Error("Failed to shutdown", "error", stopErr)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
v.Cancel()
return nil
},
})
}))
err := RunWithApp(ctx, opts)
if err != nil {
return err
}
// The goroutine sets failed before triggering the shutdown that lets
// RunWithApp return, so the write is in place by the time we read it.
mu.Lock()
defer mu.Unlock()
if failed {
return errReported
}
return nil
}
// runVaultikApp runs the standard single-operation command lifecycle // runVaultikApp runs the standard single-operation command lifecycle
// shared by the list/purge/verify/remove/remote-info subcommands: // shared by the list/purge/verify/remove/remote-info subcommands:
// resolve the config, start the fx app, run op against the Vaultik // resolve the config, then run op against the Vaultik instance through
// instance in a goroutine, report a failure prefixed with failMsg // RunOperation, reporting a failure prefixed with failMsg (suppressed
// (suppressed while suppressErrors is true, e.g. under --json), then // while suppressErrors is true, e.g. under --json). extraQuiet is OR-ed
// trigger shutdown. The operation is cancelled when the app stops. // into LogOptions.Quiet (e.g. --json output modes).
// extraQuiet is OR-ed into LogOptions.Quiet (e.g. --json output modes).
func runVaultikApp( func runVaultikApp(
cmd *cobra.Command, extraQuiet, suppressErrors bool, cmd *cobra.Command, extraQuiet, suppressErrors bool,
failMsg string, op func(v *vaultik.Vaultik) error, failMsg string, op func(v *vaultik.Vaultik) error,
@@ -214,47 +292,20 @@ func runVaultikApp(
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet || extraQuiet, Quiet: rootFlags.Quiet || extraQuiet,
}, },
Modules: []fx.Option{}, }, op, func(err error) {
Invokes: []fx.Option{ if suppressErrors {
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { return
lc.Append(fx.Hook{ }
OnStart: func(_ context.Context) error {
go func() {
err := op(v)
if err != nil {
if !errors.Is(err, context.Canceled) {
if !suppressErrors {
log.Error(failMsg, "error", err)
ReportErrorf("%s: %v", failMsg, err)
}
os.Exit(1) log.Error(failMsg, "error", err)
} ReportErrorf("%s: %v", failMsg, err)
}
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
v.Cancel()
return nil
},
})
}),
},
}) })
} }
+18 -3
View File
@@ -1,6 +1,7 @@
package cli package cli
import ( import (
"errors"
"io" "io"
"os" "os"
"strings" "strings"
@@ -19,7 +20,11 @@ const shortCommitLen = 12
// flag is present in os.Args — see bannerSuppressedInArgs), executes the // flag is present in os.Args — see bannerSuppressedInArgs), executes the
// root cobra command, and routes any returned error through the // root cobra command, and routes any returned error through the
// ui.Writer so the user sees a properly formatted "🛑 ERROR:" line. // ui.Writer so the user sees a properly formatted "🛑 ERROR:" line.
func Entry() { //
// 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.Stdout)
rootCmd := NewRootCommand() rootCmd := NewRootCommand()
@@ -27,9 +32,19 @@ func Entry() {
err := rootCmd.Execute() err := rootCmd.Execute()
if err != nil { if err != nil {
ReportErrorf("%s", err.Error()) // An operation that ran inside the fx app has already reported
os.Exit(1) // its own failure (and suppressed it under --json); 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) {
ReportErrorf("%s", err.Error())
}
return 1
} }
return 0
} }
// emitStartupBanner writes the startup banner to w unless args (the // emitStartupBanner writes the startup banner to w unless args (the
+1 -1
View File
@@ -230,7 +230,7 @@ func TestEntryJSONStdoutIsExactlyOneDocument(t *testing.T) {
programName, flagConfig, configPath, cmdSnapshot, cmdList, flagJSON, programName, flagConfig, configPath, cmdSnapshot, cmdList, flagJSON,
} }
stdout := captureProcessStdout(t, Entry) stdout := captureProcessStdout(t, func() { _ = Entry() })
requireExactlyOneJSONDocument(t, stdout) requireExactlyOneJSONDocument(t, stdout)
+1 -1
View File
@@ -81,7 +81,7 @@ func TestEntryPruneJSONStdoutIsExactlyOneDocument(t *testing.T) {
programName, flagConfig, configPath, cmdPrune, flagJSON, programName, flagConfig, configPath, cmdPrune, flagJSON,
} }
stdout := captureProcessStdout(t, Entry) stdout := captureProcessStdout(t, func() { _ = Entry() })
requireExactlyOneJSONDocument(t, stdout) requireExactlyOneJSONDocument(t, stdout)
+58
View File
@@ -0,0 +1,58 @@
package cli //nolint:testpackage // shares programName and the capture helpers
import (
"os"
"testing"
"github.com/stretchr/testify/assert"
)
// TestEntryReturnsStatusCode pins the contract main() relies on for
// issue #75: Entry reports success or failure through its return value
// and never calls os.Exit. An os.Exit from inside Entry would skip
// main's deferred profile writers and truncate the profile of a failing
// command. main turns this code into os.Exit only after those defers
// 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.
//
//nolint:paralleltest // replaces os.Args and rootFlags
func TestEntryReturnsStatusCode(t *testing.T) {
for _, testCase := range []struct {
name string
args []string
want int
}{
{
// version is self-contained: it needs no config and no
// destination store, so it exercises the success path.
name: "successful command returns zero",
args: []string{programName, "version"},
want: 0,
},
{
name: "unknown command returns one",
args: []string{programName, "no-such-command"},
want: 1,
},
} {
t.Run(testCase.name, func(t *testing.T) {
previousArgs := os.Args
t.Cleanup(func() {
os.Args = previousArgs
rootFlags = RootFlags{}
})
os.Args = testCase.args
var code int
_ = captureProcessStdout(t, func() { code = Entry() })
assert.Equal(t, testCase.want, code)
})
}
}
+6 -37
View File
@@ -1,12 +1,7 @@
package cli package cli
import ( import (
"context"
"errors"
"os"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/fx"
"sneak.berlin/go/vaultik/internal/log" "sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/vaultik" "sneak.berlin/go/vaultik/internal/vaultik"
) )
@@ -33,44 +28,18 @@ func NewInfoCommand() *cobra.Command {
// Use the app framework // Use the app framework
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet, Quiet: rootFlags.Quiet,
}, },
Modules: []fx.Option{}, }, func(v *vaultik.Vaultik) error {
Invokes: []fx.Option{ return v.ShowInfo()
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { }, func(err error) {
lc.Append(fx.Hook{ log.Error("Failed to show info", "error", err)
OnStart: func(_ context.Context) error { ReportErrorf("Failed to show info: %v", err)
go func() {
err := v.ShowInfo()
if err != nil {
if !errors.Is(err, context.Canceled) {
log.Error("Failed to show info", "error", err)
ReportErrorf("Failed to show info: %v", err)
os.Exit(1)
}
}
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
v.Cancel()
return nil
},
})
}),
},
}) })
}, },
} }
+9 -43
View File
@@ -1,12 +1,7 @@
package cli package cli
import ( import (
"context"
"errors"
"os"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/fx"
"sneak.berlin/go/vaultik/internal/log" "sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/vaultik" "sneak.berlin/go/vaultik/internal/vaultik"
) )
@@ -41,51 +36,22 @@ work (e.g. after a crashed backup or to reclaim storage).`,
// Use the app framework like other commands // Use the app framework like other commands
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet || opts.JSON, Quiet: rootFlags.Quiet || opts.JSON,
}, },
Modules: []fx.Option{}, }, func(v *vaultik.Vaultik) error {
Invokes: []fx.Option{ return v.Prune(opts)
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { }, func(err error) {
lc.Append(fx.Hook{ if opts.JSON {
OnStart: func(_ context.Context) error { return
// Start the prune operation in a goroutine }
go func() {
// Run the prune operation
err := v.Prune(opts)
if err != nil {
if !errors.Is(err, context.Canceled) {
if !opts.JSON {
log.Error("Prune operation failed", "error", err)
ReportErrorf("Prune failed: %v", err)
}
os.Exit(1) log.Error("Prune operation failed", "error", err)
} ReportErrorf("Prune failed: %v", err)
}
// Shutdown the app when prune completes
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
log.Debug("Stopping prune operation")
v.Cancel()
return nil
},
})
}),
},
}) })
}, },
} }
+9 -37
View File
@@ -1,12 +1,9 @@
package cli package cli
import ( import (
"context"
"errors" "errors"
"os"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/fx"
"sneak.berlin/go/vaultik/internal/log" "sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/vaultik" "sneak.berlin/go/vaultik/internal/vaultik"
) )
@@ -83,47 +80,22 @@ func newRemoteInfoCommand() *cobra.Command {
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet || jsonOutput, Quiet: rootFlags.Quiet || jsonOutput,
}, },
Modules: []fx.Option{}, }, func(v *vaultik.Vaultik) error {
Invokes: []fx.Option{ return v.RemoteInfo(jsonOutput)
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { }, func(err error) {
lc.Append(fx.Hook{ if jsonOutput {
OnStart: func(_ context.Context) error { return
go func() { }
err := v.RemoteInfo(jsonOutput)
if err != nil {
if !errors.Is(err, context.Canceled) {
if !jsonOutput {
log.Error("Failed to get remote info", "error", err)
ReportErrorf("Failed to get remote info: %v", err)
}
os.Exit(1) log.Error("Failed to get remote info", "error", err)
} ReportErrorf("Failed to get remote info: %v", err)
}
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
v.Cancel()
return nil
},
})
}),
},
}) })
}, },
} }
+16 -74
View File
@@ -1,13 +1,10 @@
package cli package cli
import ( import (
"context"
"errors" "errors"
"fmt" "fmt"
"os"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/fx"
"sneak.berlin/go/vaultik/internal/log" "sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/vaultik" "sneak.berlin/go/vaultik/internal/vaultik"
) )
@@ -86,7 +83,8 @@ specifying a path using --config or by setting VAULTIK_CONFIG to a path.`,
// Use the backup functionality from cli package // Use the backup functionality from cli package
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ // --cron suppression is wired through v.UI by setupGlobals.
return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
@@ -94,42 +92,11 @@ specifying a path using --config or by setting VAULTIK_CONFIG to a path.`,
Cron: opts.Cron, Cron: opts.Cron,
Quiet: rootFlags.Quiet, Quiet: rootFlags.Quiet,
}, },
Modules: []fx.Option{}, }, func(v *vaultik.Vaultik) error {
Invokes: []fx.Option{ return v.CreateSnapshot(opts)
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { }, func(err error) {
lc.Append(fx.Hook{ log.Error("Snapshot creation failed", "error", err)
OnStart: func(_ context.Context) error { ReportErrorf("Snapshot creation failed: %v", err)
// Start the snapshot creation in a goroutine
go func() {
// --cron suppression is wired through v.UI by setupGlobals.
err := v.CreateSnapshot(opts)
if err != nil {
if !errors.Is(err, context.Canceled) {
log.Error("Snapshot creation failed", "error", err)
ReportErrorf("Snapshot creation failed: %v", err)
os.Exit(1)
}
}
// Shutdown the app when snapshot completes
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
log.Debug("Stopping snapshot creation")
// Cancel the Vaultik context
v.Cancel()
return nil
},
})
}),
},
}) })
}, },
} }
@@ -234,47 +201,22 @@ func newSnapshotVerifyCommand() *cobra.Command {
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet || opts.JSON, Quiet: rootFlags.Quiet || opts.JSON,
}, },
Modules: []fx.Option{}, }, func(v *vaultik.Vaultik) error {
Invokes: []fx.Option{ return v.VerifySnapshotWithOptions(snapshotID, opts)
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) { }, func(err error) {
lc.Append(fx.Hook{ if opts.JSON {
OnStart: func(_ context.Context) error { return
go func() { }
err := v.VerifySnapshotWithOptions(snapshotID, opts)
if err != nil {
if !errors.Is(err, context.Canceled) {
if !opts.JSON {
log.Error("Verification failed", "error", err)
ReportErrorf("Verification failed: %v", err)
}
os.Exit(1) log.Error("Verification failed", "error", err)
} ReportErrorf("Verification failed: %v", err)
}
err = v.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
v.Cancel()
return nil
},
})
}),
},
}) })
}, },
} }
+14 -87
View File
@@ -1,16 +1,8 @@
package cli package cli
import ( import (
"context"
"errors"
"os"
"github.com/spf13/cobra" "github.com/spf13/cobra"
"go.uber.org/fx"
"sneak.berlin/go/vaultik/internal/config"
"sneak.berlin/go/vaultik/internal/globals"
"sneak.berlin/go/vaultik/internal/log" "sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/storage"
"sneak.berlin/go/vaultik/internal/vaultik" "sneak.berlin/go/vaultik/internal/vaultik"
) )
@@ -25,15 +17,6 @@ type RestoreOptions struct {
Verify bool // Verify restored files after restore Verify bool // Verify restored files after restore
} }
// RestoreApp contains all dependencies needed for restore
type RestoreApp struct {
Globals *globals.Globals
Config *config.Config
Storage storage.Storer
Vaultik *vaultik.Vaultik
Shutdowner fx.Shutdowner
}
// newSnapshotRestoreCommand creates the 'snapshot restore' subcommand // newSnapshotRestoreCommand creates the 'snapshot restore' subcommand
func newSnapshotRestoreCommand() *cobra.Command { func newSnapshotRestoreCommand() *cobra.Command {
opts := &RestoreOptions{} opts := &RestoreOptions{}
@@ -77,7 +60,8 @@ Examples:
return cmd return cmd
} }
// runRestore parses arguments and runs the restore operation through the app framework // runRestore parses arguments and runs the restore operation through the
// app framework.
func runRestore(cmd *cobra.Command, args []string, opts *RestoreOptions) error { func runRestore(cmd *cobra.Command, args []string, opts *RestoreOptions) error {
snapshotID := args[0] snapshotID := args[0]
@@ -86,87 +70,30 @@ func runRestore(cmd *cobra.Command, args []string, opts *RestoreOptions) error {
opts.Paths = args[restoreMinArgs:] opts.Paths = args[restoreMinArgs:]
} }
// Use unified config resolution
configPath, err := ResolveConfigPath() configPath, err := ResolveConfigPath()
if err != nil { if err != nil {
return err return err
} }
// Use the app framework like other commands
rootFlags := GetRootFlags() rootFlags := GetRootFlags()
return RunWithApp(cmd.Context(), AppOptions{ return RunOperation(cmd.Context(), AppOptions{
ConfigPath: configPath, ConfigPath: configPath,
LogOptions: log.Options{ LogOptions: log.Options{
Verbose: rootFlags.Verbose, Verbose: rootFlags.Verbose,
Debug: rootFlags.Debug, Debug: rootFlags.Debug,
Quiet: rootFlags.Quiet, Quiet: rootFlags.Quiet,
}, },
Modules: buildRestoreModules(), }, func(v *vaultik.Vaultik) error {
Invokes: buildRestoreInvokes(snapshotID, opts), return v.Restore(&vaultik.RestoreOptions{
SnapshotID: snapshotID,
TargetDir: opts.TargetDir,
Paths: opts.Paths,
Verify: opts.Verify,
SkipErrors: rootFlags.SkipErrors,
})
}, func(err error) {
log.Error("Restore operation failed", "error", err)
ReportErrorf("Restore failed: %v", err)
}) })
} }
// buildRestoreModules returns the fx.Options for dependency injection in restore
func buildRestoreModules() []fx.Option {
return []fx.Option{
fx.Provide(fx.Annotate(
func(g *globals.Globals, cfg *config.Config,
storer storage.Storer, v *vaultik.Vaultik, shutdowner fx.Shutdowner) *RestoreApp {
return &RestoreApp{
Globals: g,
Config: cfg,
Storage: storer,
Vaultik: v,
Shutdowner: shutdowner,
}
},
)),
}
}
// buildRestoreInvokes returns the fx.Options that wire up the restore lifecycle
func buildRestoreInvokes(snapshotID string, opts *RestoreOptions) []fx.Option {
return []fx.Option{
fx.Invoke(func(app *RestoreApp, lc fx.Lifecycle) {
lc.Append(fx.Hook{
OnStart: func(_ context.Context) error {
// Start the restore operation in a goroutine
go func() {
// Run the restore operation
restoreOpts := &vaultik.RestoreOptions{
SnapshotID: snapshotID,
TargetDir: opts.TargetDir,
Paths: opts.Paths,
Verify: opts.Verify,
SkipErrors: GetRootFlags().SkipErrors,
}
err := app.Vaultik.Restore(restoreOpts)
if err != nil {
if !errors.Is(err, context.Canceled) {
log.Error("Restore operation failed", "error", err)
ReportErrorf("Restore failed: %v", err)
os.Exit(1)
}
}
// Shutdown the app when restore completes
err = app.Shutdowner.Shutdown()
if err != nil {
log.Error("Failed to shutdown", "error", err)
}
}()
return nil
},
OnStop: func(_ context.Context) error {
log.Debug("Stopping restore operation")
app.Vaultik.Cancel()
return nil
},
})
}),
}
}
+198
View File
@@ -0,0 +1,198 @@
package storage_test
import (
"bytes"
"context"
"errors"
"io"
"reflect"
"sort"
"testing"
"sneak.berlin/go/vaultik/internal/storage"
)
// runStorerConformance is the shared Storer contract. Every backend that
// can run in-process is expected to pass it: TestFileStorer runs it against
// file://, TestS3Storer against s3://. A new backend inherits this coverage
// by passing its own constructor, so the contract is defined once.
//
// It exercises the public Storer interface: round-trip, stat, list with
// prefix filtering, overwrite, delete, delete-of-missing, and not-found on
// Get and Stat. Each section takes its own fresh backend instance, so the
// order of sections never matters and no section sees another's objects.
func runStorerConformance(t *testing.T, newStorer func(*testing.T) storage.Storer) {
t.Helper()
conformanceRoundTrip(t, newStorer(t))
conformanceOverwrite(t, newStorer(t))
conformanceList(t, newStorer(t))
conformanceDelete(t, newStorer(t))
conformanceNotFound(t, newStorer(t))
}
// conformanceRoundTrip stores a nested key, then reads it back and stats it.
func conformanceRoundTrip(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "blobs/aa/bb/object.bin"
want := []byte("round-trip payload")
err := s.Put(ctx, key, bytes.NewReader(want))
if err != nil {
t.Fatalf("Put: %v", err)
}
got := getBytes(t, s, key)
if !bytes.Equal(got, want) {
t.Errorf("Get returned %q, want %q", got, want)
}
info, err := s.Stat(ctx, key)
if err != nil {
t.Fatalf("Stat: %v", err)
}
if info.Key != key {
t.Errorf("Stat key = %q, want %q", info.Key, key)
}
if info.Size != int64(len(want)) {
t.Errorf("Stat size = %d, want %d", info.Size, len(want))
}
}
// conformanceOverwrite checks that a second Put replaces the first.
func conformanceOverwrite(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "meta/snapshot.json"
err := s.Put(ctx, key, bytes.NewReader([]byte("first")))
if err != nil {
t.Fatalf("first Put: %v", err)
}
want := []byte("second and longer payload")
err = s.Put(ctx, key, bytes.NewReader(want))
if err != nil {
t.Fatalf("second Put: %v", err)
}
got := getBytes(t, s, key)
if !bytes.Equal(got, want) {
t.Errorf("after overwrite Get returned %q, want %q", got, want)
}
}
// conformanceList checks prefix filtering and the empty result for a
// prefix that matches nothing.
func conformanceList(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
keys := []string{"blobs/aa/one", "blobs/bb/two", "meta/three"}
for _, k := range keys {
err := s.Put(ctx, k, bytes.NewReader([]byte("data")))
if err != nil {
t.Fatalf("Put %q: %v", k, err)
}
}
if got := listSorted(t, s, ""); !reflect.DeepEqual(got, keys) {
t.Errorf("List(\"\") = %v, want %v", got, keys)
}
wantBlobs := []string{"blobs/aa/one", "blobs/bb/two"}
if got := listSorted(t, s, "blobs/"); !reflect.DeepEqual(got, wantBlobs) {
t.Errorf("List(\"blobs/\") = %v, want %v", got, wantBlobs)
}
if got := listSorted(t, s, "absent/"); len(got) != 0 {
t.Errorf("List(\"absent/\") = %v, want empty", got)
}
}
// conformanceDelete checks that Delete removes an object and that deleting
// a missing key is not an error.
func conformanceDelete(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "blobs/cc/gone.bin"
err := s.Put(ctx, key, bytes.NewReader([]byte("temporary")))
if err != nil {
t.Fatalf("Put: %v", err)
}
err = s.Delete(ctx, key)
if err != nil {
t.Fatalf("Delete: %v", err)
}
_, err = s.Get(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Get after Delete error = %v, want ErrNotFound", err)
}
err = s.Delete(ctx, key)
if err != nil {
t.Errorf("Delete of missing key = %v, want nil", err)
}
}
// conformanceNotFound checks Get and Stat on an absent key.
func conformanceNotFound(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "never/written"
_, err := s.Get(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Get error = %v, want ErrNotFound", err)
}
_, err = s.Stat(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Stat error = %v, want ErrNotFound", err)
}
}
// getBytes reads a key fully and closes the reader.
func getBytes(t *testing.T, s storage.Storer, key string) []byte {
t.Helper()
rc, err := s.Get(context.Background(), key)
if err != nil {
t.Fatalf("Get %q: %v", key, err)
}
defer func() { _ = rc.Close() }()
data, err := io.ReadAll(rc)
if err != nil {
t.Fatalf("read %q: %v", key, err)
}
return data
}
// listSorted returns the keys under a prefix in a stable order.
func listSorted(t *testing.T, s storage.Storer, prefix string) []string {
t.Helper()
keys, err := s.List(context.Background(), prefix)
if err != nil {
t.Fatalf("List %q: %v", prefix, err)
}
sort.Strings(keys)
return keys
}
+79 -54
View File
@@ -46,31 +46,18 @@ func (f *FileStorer) SetFilesystem(fs afero.Fs) {
// storage base path. // storage base path.
const storageDirPerm = 0o755 const storageDirPerm = 0o755
// tempSuffix marks a partially written object. writeAtomic streams into a
// temp file carrying this suffix and only renames it onto the real key once
// the whole object is on disk, so an interrupted write can never leave a
// truncated object at the key a later run would Stat and trust as a complete
// blob. List and ListStream skip these files, so a leftover from an
// interrupted write is never listed or trusted as a blob; it is otherwise
// harmless and is overwritten when the same key is written again.
const tempSuffix = ".partial"
// Put stores data at the specified key. // Put stores data at the specified key.
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error { func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
path := f.fullPath(key) return f.writeAtomic(key, data, nil)
// Create parent directories
dir := filepath.Dir(path)
err := f.fs.MkdirAll(dir, storageDirPerm)
if err != nil {
return fmt.Errorf("creating directories: %w", err)
}
file, err := f.fs.Create(path)
if err != nil {
return fmt.Errorf("creating file: %w", err)
}
defer func() { _ = file.Close() }()
_, err = io.Copy(file, data)
if err != nil {
return fmt.Errorf("writing file: %w", err)
}
return nil
} }
// PutWithProgress stores data with progress reporting. // PutWithProgress stores data with progress reporting.
@@ -78,35 +65,7 @@ func (f *FileStorer) PutWithProgress(
_ context.Context, key string, data io.Reader, _ context.Context, key string, data io.Reader,
_ int64, progress ProgressCallback, _ int64, progress ProgressCallback,
) error { ) error {
path := f.fullPath(key) return f.writeAtomic(key, data, progress)
// Create parent directories
dir := filepath.Dir(path)
err := f.fs.MkdirAll(dir, storageDirPerm)
if err != nil {
return fmt.Errorf("creating directories: %w", err)
}
file, err := f.fs.Create(path)
if err != nil {
return fmt.Errorf("creating file: %w", err)
}
defer func() { _ = file.Close() }()
// Wrap with progress tracking
pw := &progressWriter{
writer: file,
callback: progress,
}
_, err = io.Copy(pw, data)
if err != nil {
return fmt.Errorf("writing file: %w", err)
}
return nil
} }
// Get retrieves data from the specified key. // Get retrieves data from the specified key.
@@ -188,7 +147,7 @@ func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error)
default: default:
} }
if !info.IsDir() { if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
// Convert back to key (relative path from basePath) // Convert back to key (relative path from basePath)
relPath, err := filepath.Rel(f.basePath, path) relPath, err := filepath.Rel(f.basePath, path)
if err != nil { if err != nil {
@@ -245,7 +204,7 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec
return nil //nolint:nilerr // continue walking despite errors return nil //nolint:nilerr // continue walking despite errors
} }
if !info.IsDir() { if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
relPath, err := filepath.Rel(f.basePath, path) relPath, err := filepath.Rel(f.basePath, path)
if err != nil { if err != nil {
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)} ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
@@ -275,6 +234,72 @@ func (f *FileStorer) Info() Info {
} }
} }
// writeAtomic streams data into a temp file in the destination directory,
// fsyncs it, and renames it onto the final key. The key therefore appears
// only once the whole object has been durably written; a failure part-way
// leaves a temp file (removed here on the failing path) rather than a
// truncated object at the key.
func (f *FileStorer) writeAtomic(
key string, data io.Reader, progress ProgressCallback,
) error {
path := f.fullPath(key)
dir := filepath.Dir(path)
err := f.fs.MkdirAll(dir, storageDirPerm)
if err != nil {
return fmt.Errorf("creating directories: %w", err)
}
tmp, err := afero.TempFile(f.fs, dir, filepath.Base(path)+"-*"+tempSuffix)
if err != nil {
return fmt.Errorf("creating temp file: %w", err)
}
tmpPath := tmp.Name()
// Remove the temp file unless the rename below claims it. On the success
// path renamed is true, so the deferred Close and Remove are harmless
// no-ops on a name that no longer exists.
renamed := false
defer func() {
_ = tmp.Close()
if !renamed {
_ = f.fs.Remove(tmpPath)
}
}()
var w io.Writer = tmp
if progress != nil {
w = &progressWriter{writer: tmp, callback: progress}
}
_, err = io.Copy(w, data)
if err != nil {
return fmt.Errorf("writing file: %w", err)
}
err = tmp.Sync()
if err != nil {
return fmt.Errorf("syncing temp file: %w", err)
}
err = tmp.Close()
if err != nil {
return fmt.Errorf("closing temp file: %w", err)
}
err = f.fs.Rename(tmpPath, path)
if err != nil {
return fmt.Errorf("renaming temp file: %w", err)
}
renamed = true
return nil
}
// fullPath returns the full filesystem path for a key. // fullPath returns the full filesystem path for a key.
func (f *FileStorer) fullPath(key string) string { func (f *FileStorer) fullPath(key string) string {
return filepath.Join(f.basePath, key) return filepath.Join(f.basePath, key)
+119
View File
@@ -0,0 +1,119 @@
package storage_test
import (
"context"
"errors"
"os"
"path/filepath"
"strings"
"testing"
"sneak.berlin/go/vaultik/internal/storage"
)
// errStreamInterrupted stands in for an upload cut off mid-stream.
var errStreamInterrupted = errors.New("connection reset mid-upload")
// failingReader yields its data once, then fails.
type failingReader struct {
data []byte
done bool
}
func (r *failingReader) Read(p []byte) (int, error) {
if r.done {
return 0, errStreamInterrupted
}
n := copy(p, r.data)
r.done = true
return n, nil
}
// TestFileStorer_InterruptedWriteLeavesNoTrustedObject checks that a write
// cut off mid-stream leaves nothing at the destination key, so a later run
// cannot Stat a truncated object and trust it as a complete blob.
func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
t.Parallel()
f, err := storage.NewFileStorer(t.TempDir())
if err != nil {
t.Fatalf("NewFileStorer: %v", err)
}
ctx := context.Background()
key := "blobs/aa/bb/aabbccddeeff"
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
if err == nil {
t.Fatal("expected the interrupted write to fail, got nil")
}
_, err = f.Stat(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Fatalf("expected key absent after interrupted write, got Stat err %v", err)
}
keys, err := f.List(ctx, "blobs/")
if err != nil {
t.Fatalf("List: %v", err)
}
if len(keys) != 0 {
t.Fatalf("expected no keys listed after interrupted write, got %v", keys)
}
}
// TestFileStorer_ListSkipsPartialFiles checks that a leftover temp file (the
// storage layer names them with a ".partial" suffix) is never surfaced as a
// key by List or ListStream.
func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
t.Parallel()
base := t.TempDir()
f, err := storage.NewFileStorer(base)
if err != nil {
t.Fatalf("NewFileStorer: %v", err)
}
ctx := context.Background()
realKey := "blobs/aa/bb/aabbccddeeff"
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
if err != nil {
t.Fatalf("Put: %v", err)
}
// A stray temp file, as an interrupted write would leave behind.
leftover := filepath.Join(base, "blobs/aa/bb/aabbccddeeff-123456.partial")
err = os.WriteFile(leftover, []byte("half"), 0o600)
if err != nil {
t.Fatalf("writing leftover temp file: %v", err)
}
keys, err := f.List(ctx, "blobs/")
if err != nil {
t.Fatalf("List: %v", err)
}
if len(keys) != 1 || keys[0] != realKey {
t.Fatalf("List should return only the real key, got %v", keys)
}
var streamed []string
for obj := range f.ListStream(ctx, "blobs/") {
if obj.Err != nil {
t.Fatalf("ListStream: %v", obj.Err)
}
streamed = append(streamed, obj.Key)
}
if len(streamed) != 1 || streamed[0] != realKey {
t.Fatalf("ListStream should return only the real key, got %v", streamed)
}
}
+1 -188
View File
@@ -1,12 +1,6 @@
package storage_test package storage_test
import ( import (
"bytes"
"context"
"errors"
"io"
"reflect"
"sort"
"testing" "testing"
"sneak.berlin/go/vaultik/internal/storage" "sneak.berlin/go/vaultik/internal/storage"
@@ -26,189 +20,8 @@ func newFileStorer(t *testing.T) storage.Storer {
return s return s
} }
// TestFileStorer runs the Storer contract against the file:// backend. // TestFileStorer runs the shared Storer contract against the file:// backend.
// The conformance helper is backend-agnostic, so a new backend inherits
// this coverage by passing its own constructor.
func TestFileStorer(t *testing.T) { func TestFileStorer(t *testing.T) {
t.Parallel() t.Parallel()
runStorerConformance(t, newFileStorer) runStorerConformance(t, newFileStorer)
} }
// runStorerConformance exercises the public Storer contract: round-trip,
// stat, list, overwrite, delete, and not-found behaviour. Each section
// uses its own backend instance so ordering never matters.
func runStorerConformance(t *testing.T, newStorer func(*testing.T) storage.Storer) {
t.Helper()
conformanceRoundTrip(t, newStorer(t))
conformanceOverwrite(t, newStorer(t))
conformanceList(t, newStorer(t))
conformanceDelete(t, newStorer(t))
conformanceNotFound(t, newStorer(t))
}
// conformanceRoundTrip stores a nested key, then reads it back and stats it.
func conformanceRoundTrip(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "blobs/aa/bb/object.bin"
want := []byte("round-trip payload")
err := s.Put(ctx, key, bytes.NewReader(want))
if err != nil {
t.Fatalf("Put: %v", err)
}
got := getBytes(t, s, key)
if !bytes.Equal(got, want) {
t.Errorf("Get returned %q, want %q", got, want)
}
info, err := s.Stat(ctx, key)
if err != nil {
t.Fatalf("Stat: %v", err)
}
if info.Key != key {
t.Errorf("Stat key = %q, want %q", info.Key, key)
}
if info.Size != int64(len(want)) {
t.Errorf("Stat size = %d, want %d", info.Size, len(want))
}
}
// conformanceOverwrite checks that a second Put replaces the first.
func conformanceOverwrite(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "meta/snapshot.json"
err := s.Put(ctx, key, bytes.NewReader([]byte("first")))
if err != nil {
t.Fatalf("first Put: %v", err)
}
want := []byte("second and longer payload")
err = s.Put(ctx, key, bytes.NewReader(want))
if err != nil {
t.Fatalf("second Put: %v", err)
}
got := getBytes(t, s, key)
if !bytes.Equal(got, want) {
t.Errorf("after overwrite Get returned %q, want %q", got, want)
}
}
// conformanceList checks prefix filtering and the empty result for a
// prefix that matches nothing.
func conformanceList(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
keys := []string{"blobs/aa/one", "blobs/bb/two", "meta/three"}
for _, k := range keys {
err := s.Put(ctx, k, bytes.NewReader([]byte("data")))
if err != nil {
t.Fatalf("Put %q: %v", k, err)
}
}
if got := listSorted(t, s, ""); !reflect.DeepEqual(got, keys) {
t.Errorf("List(\"\") = %v, want %v", got, keys)
}
wantBlobs := []string{"blobs/aa/one", "blobs/bb/two"}
if got := listSorted(t, s, "blobs/"); !reflect.DeepEqual(got, wantBlobs) {
t.Errorf("List(\"blobs/\") = %v, want %v", got, wantBlobs)
}
if got := listSorted(t, s, "absent/"); len(got) != 0 {
t.Errorf("List(\"absent/\") = %v, want empty", got)
}
}
// conformanceDelete checks that Delete removes an object and that deleting
// a missing key is not an error.
func conformanceDelete(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "blobs/cc/gone.bin"
err := s.Put(ctx, key, bytes.NewReader([]byte("temporary")))
if err != nil {
t.Fatalf("Put: %v", err)
}
err = s.Delete(ctx, key)
if err != nil {
t.Fatalf("Delete: %v", err)
}
_, err = s.Get(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Get after Delete error = %v, want ErrNotFound", err)
}
err = s.Delete(ctx, key)
if err != nil {
t.Errorf("Delete of missing key = %v, want nil", err)
}
}
// conformanceNotFound checks Get and Stat on an absent key.
func conformanceNotFound(t *testing.T, s storage.Storer) {
t.Helper()
ctx := context.Background()
key := "never/written"
_, err := s.Get(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Get error = %v, want ErrNotFound", err)
}
_, err = s.Stat(ctx, key)
if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Stat error = %v, want ErrNotFound", err)
}
}
// getBytes reads a key fully and closes the reader.
func getBytes(t *testing.T, s storage.Storer, key string) []byte {
t.Helper()
rc, err := s.Get(context.Background(), key)
if err != nil {
t.Fatalf("Get %q: %v", key, err)
}
defer func() { _ = rc.Close() }()
data, err := io.ReadAll(rc)
if err != nil {
t.Fatalf("read %q: %v", key, err)
}
return data
}
// listSorted returns the keys under a prefix in a stable order.
func listSorted(t *testing.T, s storage.Storer, prefix string) []string {
t.Helper()
keys, err := s.List(context.Background(), prefix)
if err != nil {
t.Fatalf("List %q: %v", prefix, err)
}
sort.Strings(keys)
return keys
}
+58
View File
@@ -0,0 +1,58 @@
package storage_test
import (
"context"
"errors"
"testing"
"sneak.berlin/go/vaultik/internal/storage"
)
// The rclone backend is a thin adapter over the rclone library: it turns a
// (remote, path) pair into rclone's "remote:path" string, hands it to
// rclone, and maps rclone's own results back to the Storer interface. What
// can be tested in-process, without a configured remote or network, is that
// adapter layer — how the arguments are shaped and how construction errors
// are reported. The data-plane operations (Put/Get/List/Delete) are rclone's
// own, exercised against a real provider (drive, s3-via-rclone, ...), which
// needs a configured remote with credentials and network access and so is
// out of reach of a unit test. The shared Storer conformance suite therefore
// runs against the in-process file and s3 backends; the rclone backend
// inherits that contract once a remote is configured.
//
// These tests use rclone's ":local:" on-the-fly backend, which addresses the
// local filesystem directly without any configured remote, so construction
// runs entirely in-process.
// TestNewRcloneStorerConstruction checks that a valid remote constructs a
// backend and that Info() reports the shaped "remote:path" location.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestNewRcloneStorerConstruction(t *testing.T) {
dir := t.TempDir()
s, err := storage.NewRcloneStorer(context.Background(), ":local", dir)
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
// Info().Location is the "remote:path" string the adapter builds from
// its two arguments, so asserting it confirms the argument shaping.
want := ":local:" + dir
if got := s.Info().Location; got != want {
t.Errorf("Info().Location = %q, want %q", got, want)
}
}
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
// rclone config fails construction with the ErrRemoteNotFound sentinel,
// rather than silently returning a backend pointed nowhere.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestNewRcloneStorerUnknownRemote(t *testing.T) {
_, err := storage.NewRcloneStorer(
context.Background(), "vaultik-no-such-remote", "path")
if !errors.Is(err, storage.ErrRemoteNotFound) {
t.Errorf("NewRcloneStorer error = %v, want ErrRemoteNotFound", err)
}
}
+36 -14
View File
@@ -13,18 +13,23 @@ import (
"sneak.berlin/go/vaultik/internal/storage" "sneak.berlin/go/vaultik/internal/storage"
) )
// TestS3StorerMissingKeyMapsToErrNotFound verifies that the s3 backend reports // s3TestBucket is the bucket created for each in-process S3 server.
// a missing object as storage.ErrNotFound, matching the file and rclone const s3TestBucket = "test-bucket"
// backends and the Storer contract. Without the mapping, Get and Stat leak the
// raw SDK error and errors.Is(err, storage.ErrNotFound) is false. // newS3Storer builds an s3:// backend backed by a fresh in-process
// S3 server. It reuses the same in-memory S3 harness (gofakes3 + s3mem
// over httptest) that internal/s3 and the not-found regression test use,
// so no new mock or dependency is introduced. Each call gets its own
// server, bucket, and client, so the conformance suite's per-section
// instances stay isolated.
// //
//nolint:paralleltest // shares an in-process S3 server via t.Cleanup //nolint:ireturn // conformance runs against the Storer interface by design
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) { func newS3Storer(t *testing.T) storage.Storer {
const bucket = "test-bucket" t.Helper()
backend := s3mem.New() backend := s3mem.New()
err := backend.CreateBucket(bucket) err := backend.CreateBucket(s3TestBucket)
if err != nil { if err != nil {
t.Fatalf("create bucket: %v", err) t.Fatalf("create bucket: %v", err)
} }
@@ -32,11 +37,9 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
srv := httptest.NewServer(gofakes3.New(backend).Server()) srv := httptest.NewServer(gofakes3.New(backend).Server())
t.Cleanup(srv.Close) t.Cleanup(srv.Close)
ctx := context.Background() client, err := s3.NewClient(context.Background(), s3.Config{
client, err := s3.NewClient(ctx, s3.Config{
Endpoint: srv.URL, Endpoint: srv.URL,
Bucket: bucket, Bucket: s3TestBucket,
AccessKeyID: "test", AccessKeyID: "test",
SecretAccessKey: "test", SecretAccessKey: "test",
Region: "us-east-1", Region: "us-east-1",
@@ -45,9 +48,28 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
t.Fatalf("new client: %v", err) t.Fatalf("new client: %v", err)
} }
storer := storage.NewS3Storer(client) return storage.NewS3Storer(client)
}
_, err = storer.Get(ctx, "does-not-exist") // TestS3Storer runs the shared Storer contract against the s3:// backend,
// so it is held to the same round-trip, list, delete, and not-found
// behaviour as the file:// backend.
func TestS3Storer(t *testing.T) {
t.Parallel()
runStorerConformance(t, newS3Storer)
}
// TestS3StorerMissingKeyMapsToErrNotFound pins the specific contract that a
// missing object surfaces as storage.ErrNotFound rather than the raw AWS SDK
// error. Without the mapping, errors.Is(err, storage.ErrNotFound) is false on
// s3 and callers would branch differently per backend.
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
t.Parallel()
storer := newS3Storer(t)
ctx := context.Background()
_, err := storer.Get(ctx, "does-not-exist")
if !errors.Is(err, storage.ErrNotFound) { if !errors.Is(err, storage.ErrNotFound) {
t.Errorf("Get on missing key: got %v, want ErrNotFound", err) t.Errorf("Get on missing key: got %v, want ErrNotFound", err)
} }
+79
View File
@@ -0,0 +1,79 @@
package vaultik //nolint:testpackage // exercises unexported count helpers
import (
"context"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/log"
)
// TestTableCountForReportSurfacesReadFailure is the regression guard for
// the discarded-error bug: getTableCount for a table its query cannot
// resolve must not silently become 0. A count that could not be read is
// reported as unknown, which a reader can tell apart from an empty table.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestTableCountForReportSurfacesReadFailure(t *testing.T) {
log.Initialize(log.Config{})
ctx := context.Background()
db, err := database.New(ctx, ":memory:")
require.NoError(t, err)
t.Cleanup(func() { _ = db.Close() })
v := &Vaultik{DB: db}
v.SetContext(ctx)
// A table present in the schema reads as a real count.
blobs := v.tableCountForReport("blobs")
require.NotNil(t, blobs, "an existing table must read as a real count")
assert.Equal(t, int64(0), *blobs)
// A syntactically valid name the sanitizer accepts but whose table
// the query cannot resolve is the exact shape #96 describes: a
// would-be loud failure that used to be discarded into a 0.
_, err = v.getTableCount("snapshots_missing")
require.Error(t, err, "a query against a nonexistent table must fail")
missing := v.tableCountForReport("snapshots_missing")
assert.Nil(t, missing, "a failed read is unknown, not a count")
// The rendered count for a failed read must say unknown, never 0.
assert.Equal(t, countUnknown, countText(missing))
assert.NotEqual(t, "0", countText(missing))
}
// TestCountTextDistinguishesEmptyFromUnknown pins the distinction the
// output has to preserve: 0 means the table was empty, "unknown" means
// the count could not be read.
func TestCountTextDistinguishesEmptyFromUnknown(t *testing.T) {
t.Parallel()
zero := int64(0)
seven := int64(7)
assert.Equal(t, "0", countText(&zero))
assert.Equal(t, "7", countText(&seven))
assert.Equal(t, countUnknown, countText(nil))
}
// TestCountDiffUnknownWhenEitherSideUnknown checks that a delta computed
// from an unreadable count is itself unknown rather than a plausible
// number.
func TestCountDiffUnknownWhenEitherSideUnknown(t *testing.T) {
t.Parallel()
before := int64(10)
after := int64(3)
require.NotNil(t, countDiff(&before, &after))
assert.Equal(t, int64(7), *countDiff(&before, &after))
assert.Nil(t, countDiff(nil, &after), "unknown before yields unknown delta")
assert.Nil(t, countDiff(&before, nil), "unknown after yields unknown delta")
assert.Nil(t, countDiff(nil, nil))
}
+79 -26
View File
@@ -8,6 +8,7 @@ import (
"path/filepath" "path/filepath"
"regexp" "regexp"
"sort" "sort"
"strconv"
"strings" "strings"
"time" "time"
@@ -1540,12 +1541,17 @@ func (v *Vaultik) outputRemoveJSON(result *RemoveResult) error {
return encoder.Encode(result) return encoder.Encode(result)
} }
// PruneResult contains statistics about the prune operation // PruneResult contains statistics about the prune operation.
// SnapshotsDeleted counts snapshots actually deleted. FilesDeleted,
// ChunksDeleted, and BlobsDeleted are derived from before/after row
// counts of the local index; each is nil when a count could not be read,
// so an unreadable count is reported as unknown rather than silently
// as 0.
type PruneResult struct { type PruneResult struct {
SnapshotsDeleted int64 SnapshotsDeleted int64
FilesDeleted int64 FilesDeleted *int64
ChunksDeleted int64 ChunksDeleted *int64
BlobsDeleted int64 BlobsDeleted *int64
} }
// PruneDatabase removes incomplete snapshots and orphaned files, chunks, // PruneDatabase removes incomplete snapshots and orphaned files, chunks,
@@ -1560,7 +1566,7 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
result := &PruneResult{} result := &PruneResult{}
// Snapshot counts before deletion of incompletes. // Snapshot counts before deletion of incompletes.
snapshotCountBefore, _ := v.getTableCount("snapshots") snapshotCountBefore := v.tableCountForReport("snapshots")
// First, delete any incomplete snapshots // First, delete any incomplete snapshots
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx) incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
@@ -1575,9 +1581,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
} }
// Get counts before cleanup for reporting // Get counts before cleanup for reporting
fileCountBefore, _ := v.getTableCount("files") fileCountBefore := v.tableCountForReport("files")
chunkCountBefore, _ := v.getTableCount("chunks") chunkCountBefore := v.tableCountForReport("chunks")
blobCountBefore, _ := v.getTableCount("blobs") blobCountBefore := v.tableCountForReport("blobs")
// Run the cleanup // Run the cleanup
err = v.SnapshotManager.CleanupOrphanedData(v.ctx) err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
@@ -1586,36 +1592,83 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
} }
// Get counts after cleanup // Get counts after cleanup
fileCountAfter, _ := v.getTableCount("files") fileCountAfter := v.tableCountForReport("files")
chunkCountAfter, _ := v.getTableCount("chunks") chunkCountAfter := v.tableCountForReport("chunks")
blobCountAfter, _ := v.getTableCount("blobs") blobCountAfter := v.tableCountForReport("blobs")
result.FilesDeleted = fileCountBefore - fileCountAfter result.FilesDeleted = countDiff(fileCountBefore, fileCountAfter)
result.ChunksDeleted = chunkCountBefore - chunkCountAfter result.ChunksDeleted = countDiff(chunkCountBefore, chunkCountAfter)
result.BlobsDeleted = blobCountBefore - blobCountAfter result.BlobsDeleted = countDiff(blobCountBefore, blobCountAfter)
log.Info("Local database prune complete", log.Info("Local database prune complete",
"incomplete_snapshots", result.SnapshotsDeleted, "incomplete_snapshots", result.SnapshotsDeleted,
"orphaned_files", result.FilesDeleted, "orphaned_files", countText(result.FilesDeleted),
"orphaned_chunks", result.ChunksDeleted, "orphaned_chunks", countText(result.ChunksDeleted),
"orphaned_blobs", result.BlobsDeleted, "orphaned_blobs", countText(result.BlobsDeleted),
) )
snapshotCountAfter := snapshotCountBefore - result.SnapshotsDeleted // Snapshots remaining after removing the incomplete ones; unknown if
// the pre-prune snapshot count could not be read.
snapshotsRemain := countDiff(snapshotCountBefore, &result.SnapshotsDeleted)
v.UI.Completef("Pruned local index database.") v.UI.Completef("Pruned local index database.")
v.UI.Detailf("Incomplete snapshots: %d removed (%d remain).", v.UI.Detailf("Incomplete snapshots: %s removed (%s remain).",
result.SnapshotsDeleted, snapshotCountAfter) countText(&result.SnapshotsDeleted), countText(snapshotsRemain))
v.UI.Detailf("Orphaned files: %d removed (%d remain).", v.UI.Detailf("Orphaned files: %s removed (%s remain).",
result.FilesDeleted, fileCountAfter) countText(result.FilesDeleted), countText(fileCountAfter))
v.UI.Detailf("Orphaned chunks: %d removed (%d remain).", v.UI.Detailf("Orphaned chunks: %s removed (%s remain).",
result.ChunksDeleted, chunkCountAfter) countText(result.ChunksDeleted), countText(chunkCountAfter))
v.UI.Detailf("Orphaned blobs: %d removed (%d remain).", v.UI.Detailf("Orphaned blobs: %s removed (%s remain).",
result.BlobsDeleted, blobCountAfter) countText(result.BlobsDeleted), countText(blobCountAfter))
return result, nil return result, nil
} }
// countUnknown is what a count reads as when its query could not be run,
// distinct from "0", which means the table really was empty.
const countUnknown = "unknown"
// tableCountForReport returns the row count of a table for the prune
// summary, or nil if the count could not be read. A read failure is
// logged at warn — visible even under --json, which routes warnings to
// stderr — and then rendered as unknown rather than silently becoming 0,
// so a broken query is a visible failure instead of a plausible wrong
// number.
func (v *Vaultik) tableCountForReport(tableName string) *int64 {
count, err := v.getTableCount(tableName)
if err != nil {
log.Warn("could not read table row count for prune summary",
"table", tableName, "error", err)
return nil
}
return &count
}
// countDiff returns before-after, or nil if either count is unknown so
// that an unreadable count does not collapse into a plausible delta.
func countDiff(before, after *int64) *int64 {
if before == nil || after == nil {
return nil
}
diff := *before - *after
return &diff
}
// countText renders a count that may be unknown: nil (the read failed)
// becomes "unknown", never "0", so a reader can tell an empty table from
// one that could not be queried.
func countText(count *int64) string {
if count == nil {
return countUnknown
}
return strconv.FormatInt(*count, 10)
}
// validTableNameRe matches table names containing only lowercase // validTableNameRe matches table names containing only lowercase
// alphanumeric characters and underscores. // alphanumeric characters and underscores.
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`) var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
+16 -1
View File
@@ -56,8 +56,23 @@ main() {
docker build --output=type=cacheonly \ docker build --output=type=cacheonly \
--build-arg CHECK_EPOCH="$epoch" -f Dockerfile.lint . --build-arg CHECK_EPOCH="$epoch" -f Dockerfile.lint .
# Version, commit and build date are computed here on the host, the
# same way script/docker does, and passed into the product build so
# the CI-built image reports its real source. The build context
# excludes .git (see .dockerignore), so the build cannot derive them
# itself; without these it would stamp the Dockerfile's dev/unknown
# fallbacks. VERSION comes from script/version, the source of truth
# shared with the Makefile.
version="$("$ROOT/script/version")"
commit="$(git rev-parse HEAD 2>/dev/null || echo unknown)"
commit_date="$(git show -s --format=%cs HEAD 2>/dev/null || echo unknown)"
epoch="$(date +%s%N)$$" epoch="$(date +%s%N)$$"
docker build --build-arg CHECK_EPOCH="$epoch" . docker build --build-arg CHECK_EPOCH="$epoch" \
--build-arg VERSION="$version" \
--build-arg COMMIT="$commit" \
--build-arg COMMIT_DATE="$commit_date" \
.
} }
main "$@" main "$@"
+17
View File
@@ -24,7 +24,24 @@ main() {
# whether the tree is clean. The Dockerfile now refuses to build # whether the tree is clean. The Dockerfile now refuses to build
# without a non-empty value, so this is required, not optional. # without a non-empty value, so this is required, not optional.
epoch="$(date +%s%N)$$" epoch="$(date +%s%N)$$"
# Version, commit and build date are computed here on the host,
# where .git exists, and passed into the build. The build context
# excludes .git (see .dockerignore), so the container cannot derive
# them itself -- it used to try and always got "unknown", giving
# every image a "commit: unknown" it could not be traced from.
# VERSION comes from script/version, the source of truth shared with
# the Makefile, so a Docker build reports the same string (tag,
# dev-<sha>, or a -dirty variant) that a local build of the same
# tree would.
version="$("$SCRIPT_DIR/version")"
commit="$(git rev-parse HEAD 2>/dev/null || echo unknown)"
commit_date="$(git show -s --format=%cs HEAD 2>/dev/null || echo unknown)"
docker build --build-arg CHECK_EPOCH="$epoch" \ docker build --build-arg CHECK_EPOCH="$epoch" \
--build-arg VERSION="$version" \
--build-arg COMMIT="$commit" \
--build-arg COMMIT_DATE="$commit_date" \
-t "$("$SCRIPT_DIR/projectname")" . -t "$("$SCRIPT_DIR/projectname")" .
} }