Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1073420e8b | ||
|
|
343129f891 | ||
|
|
a50e3fa038 | ||
|
|
6fcd8e1668 | ||
|
|
aab6a87f8c | ||
|
|
c355ef4d25 | ||
|
|
5927e1aa3d |
@@ -104,7 +104,12 @@ Version: 2025-06-08
|
||||
|
||||
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
|
||||
backup. When the schema changes, just change `schema.sql` (and any code
|
||||
that touches the affected tables). The local index is disposable until
|
||||
1.0 ships and is tagged.
|
||||
backup. To change the schema, edit `internal/database/schema/001.sql`
|
||||
(and any code that touches the affected tables) directly; do not add new
|
||||
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
@@ -20,8 +20,6 @@
|
||||
# golang:1.26.1-alpine, 2026-03-17
|
||||
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.
|
||||
# The sqlite driver is pure Go (modernc.org/sqlite), so no sqlite library or
|
||||
# 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 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)
|
||||
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
|
||||
# alpine:3.21, 2026-02-25
|
||||
|
||||
@@ -84,6 +84,57 @@ VAULTIK_AGE_SECRET_KEY='AGE-SECRET-KEY-...' vaultik snapshot restore <snapshot-i
|
||||
# 0 3 * * * vaultik snapshot create --cron --prune --keep-newer-than 4w
|
||||
```
|
||||
|
||||
## restoring on another machine
|
||||
|
||||
Restoring on a host that never ran the backup — a replacement machine
|
||||
after the original is gone — is the case vaultik is built for. That host
|
||||
needs only three things: the `vaultik` binary, the age **private** key,
|
||||
and the storage credentials for the destination. It does **not** need the
|
||||
local index, the original config file, or the original hostname.
|
||||
|
||||
```sh
|
||||
# install
|
||||
go install sneak.berlin/go/vaultik/cmd/vaultik@latest
|
||||
|
||||
# create a config and point it at the ORIGINAL backup destination
|
||||
vaultik config init
|
||||
vaultik config set storage_url "s3://bucket/prefix?endpoint=https://s3.example.com"
|
||||
vaultik config set s3.access_key_id "..."
|
||||
vaultik config set s3.secret_access_key "..."
|
||||
|
||||
# see what is on the destination store
|
||||
vaultik snapshot list
|
||||
```
|
||||
|
||||
`snapshot list` reads the destination store without the private key. A
|
||||
snapshot that is not in this host's (empty) local index is shown as
|
||||
remote-only: its row is identified by `<remote only:...>` rather than by
|
||||
a `hostname_name_timestamp` name, because the name lives only in the
|
||||
local index and the encrypted database and cannot be recovered from the
|
||||
store. Its timestamp and compressed size are real. (See the `snapshot
|
||||
list` description under [command details](#command-details) for the full
|
||||
explanation.)
|
||||
|
||||
Use that remote key — the hex printed inside `<remote only:...>`, or the
|
||||
full `remote_key` from `snapshot list --json` — to restore and verify:
|
||||
|
||||
```sh
|
||||
# restore everything to /tmp/restored, then check every restored file's
|
||||
# chunk hashes
|
||||
VAULTIK_AGE_SECRET_KEY='AGE-SECRET-KEY-...' \
|
||||
vaultik snapshot restore --verify <remote-key> /tmp/restored
|
||||
|
||||
# optionally, deep-verify the snapshot against the store (downloads and
|
||||
# cryptographically checks every blob)
|
||||
VAULTIK_AGE_SECRET_KEY='AGE-SECRET-KEY-...' \
|
||||
vaultik snapshot verify --deep <remote-key>
|
||||
```
|
||||
|
||||
`age_recipients` (the public key) is not needed to restore — only the
|
||||
private key in `VAULTIK_AGE_SECRET_KEY`. Both the abbreviated key printed
|
||||
in the table and the full 64-character key from `--json` are accepted; a
|
||||
leading part of the key is enough as long as it is unambiguous.
|
||||
|
||||
---
|
||||
|
||||
## cli
|
||||
@@ -245,6 +296,8 @@ local index alone, and still exits zero.
|
||||
* Default (shallow): checks that all blobs referenced in the manifest exist in storage
|
||||
* `--deep`: Downloads and decrypts each blob, verifies chunk hashes against the
|
||||
encrypted metadata database
|
||||
* Accepts the same identifiers as `snapshot restore`: a snapshot ID, or a
|
||||
remote-only snapshot's remote key (or an unambiguous leading part of it)
|
||||
* `--json`: Output results as JSON
|
||||
|
||||
**`snapshot purge`**: Remove old snapshots based on criteria. Retention is
|
||||
@@ -275,6 +328,10 @@ on the destination in one go, use `vaultik remote nuke --force`.
|
||||
|
||||
**`snapshot restore`**: Restore files from a backup snapshot.
|
||||
* Requires `VAULTIK_AGE_SECRET_KEY` environment variable
|
||||
* Accepts a snapshot ID, or — for a snapshot only on the destination
|
||||
store — its remote key (or an unambiguous leading part of it) as shown
|
||||
by `snapshot list`. See
|
||||
[restoring on another machine](#restoring-on-another-machine).
|
||||
* Optional path arguments to restore specific files/directories (default: all)
|
||||
* Preserves file permissions, timestamps, ownership (ownership requires root),
|
||||
symlinks, and empty directories
|
||||
@@ -457,9 +514,13 @@ Key fields:
|
||||
sequentially. Restore speed is bound by single-stream throughput.
|
||||
* **Device nodes, named pipes, and sockets are silently skipped.** Only
|
||||
regular files, directories, and symlinks are backed up.
|
||||
* **No database migrations.** If the local SQLite schema changes between
|
||||
versions, delete the local database (`vaultik database delete`) and run
|
||||
a full backup. Remote storage is unaffected.
|
||||
* **No upgrade path between versions.** There is no supported way to carry
|
||||
an existing local index across a schema change; if the local SQLite
|
||||
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
|
||||
filesystem snapshot or freeze. If a file is modified between the scan
|
||||
and chunk phases, the backed-up copy may reflect a partial write.
|
||||
@@ -525,14 +586,12 @@ priority.
|
||||
|
||||
### infrastructure
|
||||
|
||||
* **Cross-machine restore documentation.** The "restore from
|
||||
another host" workflow works but isn't documented as a
|
||||
first-class operation in this README. Worth a dedicated section
|
||||
once it's settled.
|
||||
* **Schema migrations.** Currently nonexistent — pre-1.0 schema
|
||||
changes are handled by `vaultik database delete` plus a full
|
||||
re-scan. Post-1.0 we'll need a migration story to keep existing
|
||||
index databases usable across upgrades.
|
||||
* **Cross-version schema upgrades.** There is no upgrade path between
|
||||
released versions — pre-1.0 schema changes are handled by `vaultik
|
||||
database delete` plus a full re-scan (see
|
||||
[`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://
|
||||
all share the Storer interface but the rclone path is the least
|
||||
exercised in CI.
|
||||
|
||||
@@ -25,6 +25,29 @@ release" is exactly the contradiction
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-09-21: Stopped an interrupted blob upload from making a later
|
||||
backup deduplicate against data that was never stored
|
||||
([issue #148](https://git.eeqj.de/sneak/vaultik/issues/148)). The
|
||||
packer commits a blob's `chunks`, `blob_chunks`, and `blobs` rows
|
||||
before the upload is attempted, so a failed upload left chunk rows
|
||||
behind and the next run skipped re-uploading them, producing a
|
||||
snapshot that reported success but could not be restored. A run now
|
||||
deduplicates only against chunks held by a blob whose `uploaded_ts` is
|
||||
set, and at startup drops any un-uploaded blob rows (and the chunks
|
||||
they orphan) so the affected data is re-chunked and re-uploaded. Blobs
|
||||
recorded with no remote backend are marked uploaded so this invariant
|
||||
holds uniformly.
|
||||
|
||||
- 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
|
||||
`storage.ErrNotFound`, like the `file` and `rclone` backends and as the
|
||||
`Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw
|
||||
@@ -33,7 +56,6 @@ release" is exactly the contradiction
|
||||
helper (reused by `HeadObject`) and a test that a missing key maps to
|
||||
`ErrNotFound`
|
||||
([issue #129](https://git.eeqj.de/sneak/vaultik/issues/129)).
|
||||
|
||||
- 2026-09-21: Fixed `verify --deep` reporting healthy snapshots as
|
||||
corrupt. Its final blob-integrity check hashed the encrypted
|
||||
downloaded bytes with a single SHA256 and compared that to the blob
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -304,10 +304,14 @@ func instructionText(contents string) string {
|
||||
}
|
||||
|
||||
// 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 {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
+11
-1
@@ -10,6 +10,16 @@ import (
|
||||
)
|
||||
|
||||
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
|
||||
if cpuProfile := os.Getenv("VAULTIK_CPUPROFILE"); cpuProfile != "" {
|
||||
f, err := os.Create(cpuProfile) //nolint:gosec // G304: operator-set path
|
||||
@@ -46,5 +56,5 @@ func main() {
|
||||
}()
|
||||
}
|
||||
|
||||
cli.Entry()
|
||||
return cli.Entry()
|
||||
}
|
||||
|
||||
+24
-5
@@ -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.
|
||||
|
||||
**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,
|
||||
delete the local SQLite database (`vaultik database delete`) and run a full
|
||||
backup. The remote storage is unaffected; the new index will re-deduplicate
|
||||
against existing remote blobs.
|
||||
|
||||
This section is the authoritative explanation of the schema/migration story;
|
||||
other documents (the README and `AGENTS.md`) link here.
|
||||
|
||||
- **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
|
||||
of Vaultik to restore a backup as was used to create it. This ensures
|
||||
compatibility with the metadata format stored in S3.
|
||||
|
||||
+90
-39
@@ -11,6 +11,7 @@ import (
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"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
|
||||
// shared by the list/purge/verify/remove/remote-info subcommands:
|
||||
// resolve the config, start the fx app, run op against the Vaultik
|
||||
// instance in a goroutine, report a failure prefixed with failMsg
|
||||
// (suppressed while suppressErrors is true, e.g. under --json), then
|
||||
// trigger shutdown. The operation is cancelled when the app stops.
|
||||
// extraQuiet is OR-ed into LogOptions.Quiet (e.g. --json output modes).
|
||||
// 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). extraQuiet is OR-ed
|
||||
// into LogOptions.Quiet (e.g. --json output modes).
|
||||
func runVaultikApp(
|
||||
cmd *cobra.Command, extraQuiet, suppressErrors bool,
|
||||
failMsg string, op func(v *vaultik.Vaultik) error,
|
||||
@@ -214,47 +292,20 @@ func runVaultikApp(
|
||||
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet || extraQuiet,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
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 {
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
if !suppressErrors {
|
||||
log.Error(failMsg, "error", err)
|
||||
ReportErrorf("%s: %v", failMsg, err)
|
||||
}
|
||||
}, op, func(err error) {
|
||||
if suppressErrors {
|
||||
return
|
||||
}
|
||||
|
||||
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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
log.Error(failMsg, "error", err)
|
||||
ReportErrorf("%s: %v", failMsg, err)
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
+18
-3
@@ -1,6 +1,7 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
@@ -19,7 +20,11 @@ const shortCommitLen = 12
|
||||
// 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.
|
||||
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)
|
||||
|
||||
rootCmd := NewRootCommand()
|
||||
@@ -27,9 +32,19 @@ func Entry() {
|
||||
|
||||
err := rootCmd.Execute()
|
||||
if err != nil {
|
||||
ReportErrorf("%s", err.Error())
|
||||
os.Exit(1)
|
||||
// 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.
|
||||
// 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
|
||||
|
||||
@@ -230,7 +230,7 @@ func TestEntryJSONStdoutIsExactlyOneDocument(t *testing.T) {
|
||||
programName, flagConfig, configPath, cmdSnapshot, cmdList, flagJSON,
|
||||
}
|
||||
|
||||
stdout := captureProcessStdout(t, Entry)
|
||||
stdout := captureProcessStdout(t, func() { _ = Entry() })
|
||||
|
||||
requireExactlyOneJSONDocument(t, stdout)
|
||||
|
||||
|
||||
@@ -81,7 +81,7 @@ func TestEntryPruneJSONStdoutIsExactlyOneDocument(t *testing.T) {
|
||||
programName, flagConfig, configPath, cmdPrune, flagJSON,
|
||||
}
|
||||
|
||||
stdout := captureProcessStdout(t, Entry)
|
||||
stdout := captureProcessStdout(t, func() { _ = Entry() })
|
||||
|
||||
requireExactlyOneJSONDocument(t, stdout)
|
||||
|
||||
|
||||
@@ -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
@@ -1,12 +1,7 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
@@ -33,44 +28,18 @@ func NewInfoCommand() *cobra.Command {
|
||||
// Use the app framework
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(_ context.Context) error {
|
||||
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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
return v.ShowInfo()
|
||||
}, func(err error) {
|
||||
log.Error("Failed to show info", "error", err)
|
||||
ReportErrorf("Failed to show info: %v", err)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
+9
-43
@@ -1,12 +1,7 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"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
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet || opts.JSON,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(_ context.Context) error {
|
||||
// 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)
|
||||
}
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
return v.Prune(opts)
|
||||
}, func(err error) {
|
||||
if opts.JSON {
|
||||
return
|
||||
}
|
||||
|
||||
os.Exit(1)
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
log.Error("Prune operation failed", "error", err)
|
||||
ReportErrorf("Prune failed: %v", err)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
+9
-37
@@ -1,12 +1,9 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
@@ -83,47 +80,22 @@ func newRemoteInfoCommand() *cobra.Command {
|
||||
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet || jsonOutput,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(_ context.Context) error {
|
||||
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)
|
||||
}
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
return v.RemoteInfo(jsonOutput)
|
||||
}, func(err error) {
|
||||
if jsonOutput {
|
||||
return
|
||||
}
|
||||
|
||||
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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
log.Error("Failed to get remote info", "error", err)
|
||||
ReportErrorf("Failed to get remote info: %v", err)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
+21
-76
@@ -1,13 +1,10 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
"go.uber.org/fx"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"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
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
// --cron suppression is wired through v.UI by setupGlobals.
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
@@ -94,42 +92,11 @@ specifying a path using --config or by setting VAULTIK_CONFIG to a path.`,
|
||||
Cron: opts.Cron,
|
||||
Quiet: rootFlags.Quiet,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(_ context.Context) error {
|
||||
// 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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
return v.CreateSnapshot(opts)
|
||||
}, func(err error) {
|
||||
log.Error("Snapshot creation failed", "error", err)
|
||||
ReportErrorf("Snapshot creation failed: %v", err)
|
||||
})
|
||||
},
|
||||
}
|
||||
@@ -221,8 +188,11 @@ func newSnapshotVerifyCommand() *cobra.Command {
|
||||
cmd := &cobra.Command{
|
||||
Use: "verify <snapshot-id>",
|
||||
Short: "Verify snapshot integrity",
|
||||
Long: "Verifies that all blobs referenced in a snapshot exist",
|
||||
Args: requireSnapshotIDArg,
|
||||
Long: "Verifies that all blobs referenced in a snapshot exist.\n\n" +
|
||||
"The snapshot may be named by its ID or, on a host with no local\n" +
|
||||
"index, by the remote key that 'snapshot list' prints for a\n" +
|
||||
"remote-only snapshot (an unambiguous leading part is enough).",
|
||||
Args: requireSnapshotIDArg,
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
snapshotID := args[0]
|
||||
|
||||
@@ -234,47 +204,22 @@ func newSnapshotVerifyCommand() *cobra.Command {
|
||||
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet || opts.JSON,
|
||||
},
|
||||
Modules: []fx.Option{},
|
||||
Invokes: []fx.Option{
|
||||
fx.Invoke(func(v *vaultik.Vaultik, lc fx.Lifecycle) {
|
||||
lc.Append(fx.Hook{
|
||||
OnStart: func(_ context.Context) error {
|
||||
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)
|
||||
}
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
return v.VerifySnapshotWithOptions(snapshotID, opts)
|
||||
}, func(err error) {
|
||||
if opts.JSON {
|
||||
return
|
||||
}
|
||||
|
||||
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
|
||||
},
|
||||
})
|
||||
}),
|
||||
},
|
||||
log.Error("Verification failed", "error", err)
|
||||
ReportErrorf("Verification failed: %v", err)
|
||||
})
|
||||
},
|
||||
}
|
||||
|
||||
@@ -1,16 +1,8 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
|
||||
"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/storage"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
|
||||
@@ -25,15 +17,6 @@ type RestoreOptions struct {
|
||||
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
|
||||
func newSnapshotRestoreCommand() *cobra.Command {
|
||||
opts := &RestoreOptions{}
|
||||
@@ -48,6 +31,10 @@ target directory.
|
||||
If no paths are specified, all files are restored.
|
||||
If paths are specified, only matching files/directories are restored.
|
||||
|
||||
The snapshot may be named by its ID or, when restoring on a host with no
|
||||
local index, by the remote key that 'snapshot list' prints for a
|
||||
remote-only snapshot (an unambiguous leading part is enough).
|
||||
|
||||
Requires the VAULTIK_AGE_SECRET_KEY environment variable to be set with
|
||||
the age private key.
|
||||
|
||||
@@ -77,7 +64,8 @@ Examples:
|
||||
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 {
|
||||
snapshotID := args[0]
|
||||
|
||||
@@ -86,87 +74,30 @@ func runRestore(cmd *cobra.Command, args []string, opts *RestoreOptions) error {
|
||||
opts.Paths = args[restoreMinArgs:]
|
||||
}
|
||||
|
||||
// Use unified config resolution
|
||||
configPath, err := ResolveConfigPath()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Use the app framework like other commands
|
||||
rootFlags := GetRootFlags()
|
||||
|
||||
return RunWithApp(cmd.Context(), AppOptions{
|
||||
return RunOperation(cmd.Context(), AppOptions{
|
||||
ConfigPath: configPath,
|
||||
LogOptions: log.Options{
|
||||
Verbose: rootFlags.Verbose,
|
||||
Debug: rootFlags.Debug,
|
||||
Quiet: rootFlags.Quiet,
|
||||
},
|
||||
Modules: buildRestoreModules(),
|
||||
Invokes: buildRestoreInvokes(snapshotID, opts),
|
||||
}, func(v *vaultik.Vaultik) error {
|
||||
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
|
||||
},
|
||||
})
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -208,6 +208,30 @@ func (r *BlobRepository) DeleteOrphaned(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteUnuploaded deletes blob rows whose upload never completed
|
||||
// (uploaded_ts IS NULL) and returns how many were removed. Their
|
||||
// blob_chunks rows are removed by the ON DELETE CASCADE foreign key.
|
||||
// A blob is only ever attached to a snapshot once its upload has been
|
||||
// recorded, so an un-uploaded blob is never referenced by a completed
|
||||
// snapshot: dropping it discards chunk rows that point at data which
|
||||
// was never stored remotely, so the affected content is re-chunked and
|
||||
// re-uploaded on the next run.
|
||||
func (r *BlobRepository) DeleteUnuploaded(ctx context.Context) (int64, error) {
|
||||
query := `DELETE FROM blobs WHERE uploaded_ts IS NULL`
|
||||
|
||||
result, err := r.db.ExecWithLog(ctx, query)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("deleting un-uploaded blobs: %w", err)
|
||||
}
|
||||
|
||||
rowsAffected, _ := result.RowsAffected()
|
||||
if rowsAffected > 0 {
|
||||
log.Debug("Deleted un-uploaded blobs", "count", rowsAffected)
|
||||
}
|
||||
|
||||
return rowsAffected, nil
|
||||
}
|
||||
|
||||
// getOne fetches a single blob row matched on the given column, or
|
||||
// (nil, nil) when no row matches.
|
||||
func (r *BlobRepository) getOne(
|
||||
|
||||
@@ -7,12 +7,32 @@ import (
|
||||
|
||||
// List returns every chunk in the index, ordered by chunk hash.
|
||||
func (r *ChunkRepository) List(ctx context.Context) ([]*Chunk, error) {
|
||||
query := `
|
||||
return r.list(ctx, `
|
||||
SELECT chunk_hash, size
|
||||
FROM chunks
|
||||
ORDER BY chunk_hash
|
||||
`
|
||||
`)
|
||||
}
|
||||
|
||||
// ListInUploadedBlobs returns the chunks that are stored in a blob whose
|
||||
// upload has completed (uploaded_ts set), ordered by chunk hash. These
|
||||
// are the only chunks a backup may safely deduplicate against: a chunk
|
||||
// recorded solely in a blob that was never uploaded refers to data that
|
||||
// is not in remote storage, so trusting it would silently drop that data
|
||||
// from later snapshots.
|
||||
func (r *ChunkRepository) ListInUploadedBlobs(ctx context.Context) ([]*Chunk, error) {
|
||||
return r.list(ctx, `
|
||||
SELECT DISTINCT c.chunk_hash, c.size
|
||||
FROM chunks c
|
||||
JOIN blob_chunks bc ON c.chunk_hash = bc.chunk_hash
|
||||
JOIN blobs b ON bc.blob_id = b.id
|
||||
WHERE b.uploaded_ts IS NOT NULL
|
||||
ORDER BY c.chunk_hash
|
||||
`)
|
||||
}
|
||||
|
||||
// list runs a chunk-selecting query and scans the (chunk_hash, size) rows.
|
||||
func (r *ChunkRepository) list(ctx context.Context, query string) ([]*Chunk, error) {
|
||||
rows, err := r.db.conn.QueryContext(ctx, query)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("querying chunks: %w", err)
|
||||
|
||||
@@ -220,7 +220,14 @@ func (s *Scanner) Scan(
|
||||
defer s.progress.Stop()
|
||||
}
|
||||
|
||||
// Phase 0: Load known files and chunks from database into memory for fast lookup
|
||||
// Phase 0: Repair any state left by an interrupted previous run, then
|
||||
// load known files and chunks from the database into memory for fast
|
||||
// lookup.
|
||||
err := s.repairInterruptedBlobs(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
knownFiles, err := s.loadDatabaseState(ctx, path)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -317,6 +324,38 @@ func (s *Scanner) loadDatabaseState(
|
||||
return knownFiles, nil
|
||||
}
|
||||
|
||||
// repairInterruptedBlobs discards blob rows left by a previous run whose
|
||||
// upload never completed. Such a blob has its chunks, blob_chunks, and
|
||||
// blobs rows committed to the local index before the upload is attempted,
|
||||
// so a crash or dropped connection mid-upload leaves them behind while the
|
||||
// data never reaches remote storage. Deduplicating against those chunks on
|
||||
// a later run would produce a snapshot that reports success but cannot be
|
||||
// restored. Dropping the un-uploaded blobs (their blob_chunks cascade) and
|
||||
// then any chunks left unreferenced forces the affected data to be
|
||||
// re-chunked and re-uploaded this run. A blob is attached to a snapshot
|
||||
// only once its upload is recorded, so this never touches a completed
|
||||
// snapshot's data.
|
||||
func (s *Scanner) repairInterruptedBlobs(ctx context.Context) error {
|
||||
removed, err := s.repos.Blobs.DeleteUnuploaded(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("removing un-uploaded blob records: %w", err)
|
||||
}
|
||||
|
||||
if removed == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
log.Warn("Discarded blob records from an interrupted previous run; "+
|
||||
"their data will be re-uploaded", "blobs", removed)
|
||||
|
||||
err = s.repos.Chunks.DeleteOrphaned(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("removing orphaned chunks: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// summarizeScanPhase calculates total size to process, updates progress tracking,
|
||||
// and prints the scan phase summary with file counts and sizes
|
||||
func (s *Scanner) summarizeScanPhase(
|
||||
@@ -392,11 +431,14 @@ func (s *Scanner) loadKnownFiles(
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// loadKnownChunks loads all known chunk hashes from the database into a
|
||||
// map for fast lookup. This avoids per-chunk database queries during file
|
||||
// processing.
|
||||
// loadKnownChunks loads the chunk hashes safe to deduplicate against into
|
||||
// an in-memory map for fast lookup, avoiding per-chunk database queries
|
||||
// during file processing. Only chunks held by a blob whose upload
|
||||
// completed are loaded: a chunk left behind by an interrupted upload
|
||||
// refers to data that never reached remote storage, and deduplicating
|
||||
// against it would silently produce an unrestorable snapshot.
|
||||
func (s *Scanner) loadKnownChunks(ctx context.Context) error {
|
||||
chunks, err := s.repos.Chunks.List(ctx)
|
||||
chunks, err := s.repos.Chunks.ListInUploadedBlobs(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("listing chunks: %w", err)
|
||||
}
|
||||
@@ -1401,7 +1443,17 @@ func (s *Scanner) finalizeProcessPhase(ctx context.Context, result *ScanResult)
|
||||
return fmt.Errorf("parsing blob ID: %w", err)
|
||||
}
|
||||
|
||||
// With no remote backend the blob's lifecycle ends here, so
|
||||
// mark it uploaded in the same transaction that attaches it to
|
||||
// the snapshot. This keeps the invariant that any blob a
|
||||
// snapshot references has uploaded_ts set, so deduplication and
|
||||
// interrupted-run repair treat these blobs as trustworthy.
|
||||
err = s.repos.WithTx(ctx, func(ctx context.Context, tx *sql.Tx) error {
|
||||
err := s.repos.Blobs.UpdateUploaded(ctx, tx, b.ID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("marking blob uploaded: %w", err)
|
||||
}
|
||||
|
||||
return s.repos.Snapshots.AddBlob(ctx, tx, s.snapshotID, blobID,
|
||||
types.BlobHash(b.Hash))
|
||||
})
|
||||
|
||||
@@ -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
@@ -46,31 +46,18 @@ func (f *FileStorer) SetFilesystem(fs afero.Fs) {
|
||||
// storage base path.
|
||||
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.
|
||||
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
|
||||
path := f.fullPath(key)
|
||||
|
||||
// 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
|
||||
return f.writeAtomic(key, data, nil)
|
||||
}
|
||||
|
||||
// PutWithProgress stores data with progress reporting.
|
||||
@@ -78,35 +65,7 @@ func (f *FileStorer) PutWithProgress(
|
||||
_ context.Context, key string, data io.Reader,
|
||||
_ int64, progress ProgressCallback,
|
||||
) error {
|
||||
path := f.fullPath(key)
|
||||
|
||||
// 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
|
||||
return f.writeAtomic(key, data, progress)
|
||||
}
|
||||
|
||||
// Get retrieves data from the specified key.
|
||||
@@ -188,7 +147,7 @@ func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error)
|
||||
default:
|
||||
}
|
||||
|
||||
if !info.IsDir() {
|
||||
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
||||
// Convert back to key (relative path from basePath)
|
||||
relPath, err := filepath.Rel(f.basePath, path)
|
||||
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
|
||||
}
|
||||
|
||||
if !info.IsDir() {
|
||||
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
||||
relPath, err := filepath.Rel(f.basePath, path)
|
||||
if err != nil {
|
||||
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.
|
||||
func (f *FileStorer) fullPath(key string) string {
|
||||
return filepath.Join(f.basePath, key)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
package storage_test
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
)
|
||||
|
||||
// newFileStorer builds a file:// backend rooted at a fresh temp directory.
|
||||
//
|
||||
//nolint:ireturn // conformance runs against the Storer interface by design
|
||||
func newFileStorer(t *testing.T) storage.Storer {
|
||||
t.Helper()
|
||||
|
||||
s, err := storage.NewFileStorer(t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("NewFileStorer: %v", err)
|
||||
}
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
// TestFileStorer runs the shared Storer contract against the file:// backend.
|
||||
func TestFileStorer(t *testing.T) {
|
||||
t.Parallel()
|
||||
runStorerConformance(t, newFileStorer)
|
||||
}
|
||||
@@ -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
@@ -13,18 +13,23 @@ import (
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
)
|
||||
|
||||
// TestS3StorerMissingKeyMapsToErrNotFound verifies that the s3 backend reports
|
||||
// a missing object as storage.ErrNotFound, matching the file and rclone
|
||||
// backends and the Storer contract. Without the mapping, Get and Stat leak the
|
||||
// raw SDK error and errors.Is(err, storage.ErrNotFound) is false.
|
||||
// s3TestBucket is the bucket created for each in-process S3 server.
|
||||
const s3TestBucket = "test-bucket"
|
||||
|
||||
// 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
|
||||
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
||||
const bucket = "test-bucket"
|
||||
//nolint:ireturn // conformance runs against the Storer interface by design
|
||||
func newS3Storer(t *testing.T) storage.Storer {
|
||||
t.Helper()
|
||||
|
||||
backend := s3mem.New()
|
||||
|
||||
err := backend.CreateBucket(bucket)
|
||||
err := backend.CreateBucket(s3TestBucket)
|
||||
if err != nil {
|
||||
t.Fatalf("create bucket: %v", err)
|
||||
}
|
||||
@@ -32,11 +37,9 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
||||
srv := httptest.NewServer(gofakes3.New(backend).Server())
|
||||
t.Cleanup(srv.Close)
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
client, err := s3.NewClient(ctx, s3.Config{
|
||||
client, err := s3.NewClient(context.Background(), s3.Config{
|
||||
Endpoint: srv.URL,
|
||||
Bucket: bucket,
|
||||
Bucket: s3TestBucket,
|
||||
AccessKeyID: "test",
|
||||
SecretAccessKey: "test",
|
||||
Region: "us-east-1",
|
||||
@@ -45,9 +48,28 @@ func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
||||
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) {
|
||||
t.Errorf("Get on missing key: got %v, want ErrNotFound", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,110 @@
|
||||
package storage_test
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
)
|
||||
|
||||
// TestParseStorageURLValid checks that each supported scheme parses into
|
||||
// the expected fields, since those fields decide which backend is built.
|
||||
func TestParseStorageURLValid(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const bucket = "mybucket"
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
raw string
|
||||
want *storage.URL
|
||||
}{
|
||||
{
|
||||
name: "file absolute path",
|
||||
raw: "file:///var/backups/vaultik",
|
||||
want: &storage.URL{Scheme: "file", Prefix: "/var/backups/vaultik"},
|
||||
},
|
||||
{
|
||||
name: "s3 bucket and prefix, ssl defaults on",
|
||||
raw: "s3://mybucket/backups/host",
|
||||
want: &storage.URL{
|
||||
Scheme: "s3", Bucket: bucket,
|
||||
Prefix: "backups/host", UseSSL: true,
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "s3 bucket only",
|
||||
raw: "s3://mybucket",
|
||||
want: &storage.URL{Scheme: "s3", Bucket: bucket, UseSSL: true},
|
||||
},
|
||||
{
|
||||
name: "s3 with endpoint, region, ssl off",
|
||||
raw: "s3://mybucket?endpoint=minio.example.com®ion=us-west-2&ssl=false",
|
||||
want: &storage.URL{
|
||||
Scheme: "s3", Bucket: bucket,
|
||||
Endpoint: "minio.example.com", Region: "us-west-2", UseSSL: false,
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "rclone remote and path",
|
||||
raw: "rclone://gdrive/backups/host",
|
||||
want: &storage.URL{
|
||||
Scheme: "rclone", RcloneRemote: "gdrive", Prefix: "backups/host",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "rclone remote only",
|
||||
raw: "rclone://gdrive",
|
||||
want: &storage.URL{Scheme: "rclone", RcloneRemote: "gdrive"},
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
got, err := storage.ParseStorageURL(tc.raw)
|
||||
if err != nil {
|
||||
t.Fatalf("ParseStorageURL(%q) returned error: %v", tc.raw, err)
|
||||
}
|
||||
|
||||
if !reflect.DeepEqual(got, tc.want) {
|
||||
t.Errorf("ParseStorageURL(%q) = %+v, want %+v", tc.raw, got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestParseStorageURLErrors checks that empty, missing, and unknown-scheme
|
||||
// inputs fail with the documented sentinel errors instead of parsing to a
|
||||
// wrong destination.
|
||||
func TestParseStorageURLErrors(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
raw string
|
||||
wantErr error
|
||||
}{
|
||||
{"empty url", "", storage.ErrEmptyStorageURL},
|
||||
{"file empty path", "file://", storage.ErrEmptyFilePath},
|
||||
{"s3 missing bucket", "s3://", storage.ErrMissingBucket},
|
||||
{"s3 missing bucket with path", "s3:///justprefix", storage.ErrMissingBucket},
|
||||
{"rclone missing remote", "rclone://", storage.ErrMissingRemote},
|
||||
{"unknown scheme", "gs://bucket/x", storage.ErrUnsupportedScheme},
|
||||
{"no scheme", "/local/path", storage.ErrUnsupportedScheme},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
_, err := storage.ParseStorageURL(tc.raw)
|
||||
if !errors.Is(err, tc.wantErr) {
|
||||
t.Errorf("ParseStorageURL(%q) error = %v, want %v",
|
||||
tc.raw, err, tc.wantErr)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,198 @@
|
||||
package vaultik //nolint:testpackage // constructs Vaultik with unexported fields
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/spf13/afero"
|
||||
"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/snapshot"
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
"sneak.berlin/go/vaultik/internal/ui"
|
||||
)
|
||||
|
||||
// errUploadInterrupted stands in for a dropped connection or kill -9
|
||||
// partway through a blob upload.
|
||||
var errUploadInterrupted = errors.New("simulated interrupted blob upload")
|
||||
|
||||
// interruptedBlobStorer wraps a real Storer but fails every blob upload,
|
||||
// modelling a run that dies mid-blob after the packer has already
|
||||
// committed the blob's chunk rows to the local index.
|
||||
type interruptedBlobStorer struct {
|
||||
storage.Storer
|
||||
}
|
||||
|
||||
func (s *interruptedBlobStorer) PutWithProgress(
|
||||
ctx context.Context, key string, reader io.Reader,
|
||||
size int64, cb storage.ProgressCallback,
|
||||
) error {
|
||||
if strings.HasPrefix(key, "blobs/") {
|
||||
return errUploadInterrupted
|
||||
}
|
||||
|
||||
return s.Storer.PutWithProgress(ctx, key, reader, size, cb)
|
||||
}
|
||||
|
||||
func (s *interruptedBlobStorer) Put(
|
||||
ctx context.Context, key string, reader io.Reader,
|
||||
) error {
|
||||
if strings.HasPrefix(key, "blobs/") {
|
||||
return errUploadInterrupted
|
||||
}
|
||||
|
||||
return s.Storer.Put(ctx, key, reader)
|
||||
}
|
||||
|
||||
// TestBackupRetryAfterInterruptedUploadIsRestorable reproduces the
|
||||
// silent-data-loss defect in
|
||||
// https://git.eeqj.de/sneak/vaultik/issues/148: a blob upload is
|
||||
// interrupted, leaving chunk rows in the local index for data that never
|
||||
// reached storage. The retry run reuses the same index. Before the fix it
|
||||
// deduplicated against those orphaned chunks, uploaded nothing for them,
|
||||
// and produced a snapshot that reported success but could not be
|
||||
// restored. The retried snapshot must instead restore byte-for-byte.
|
||||
func TestBackupRetryAfterInterruptedUploadIsRestorable(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")
|
||||
restoreDir := filepath.Join(tempDir, "restored")
|
||||
dbPath := filepath.Join(tempDir, "index.sqlite")
|
||||
|
||||
require.NoError(t, fs.MkdirAll(dataDir, 0o755))
|
||||
|
||||
// Random content forces real chunks; small blobs guarantee at least
|
||||
// one blob is finalized (and its upload attempted) during the run.
|
||||
sources := map[string][]byte{
|
||||
"a.bin": randomBytes(t, 128*1024),
|
||||
"b.bin": randomBytes(t, 128*1024),
|
||||
"c.bin": randomBytes(t, 128*1024),
|
||||
}
|
||||
for name, data := range sources {
|
||||
require.NoError(t, afero.WriteFile(
|
||||
fs, filepath.Join(dataDir, name), data, 0o644))
|
||||
}
|
||||
|
||||
cfg := &config.Config{
|
||||
AgeRecipients: []string{"age1ezrjmfpwsc95svdg0y54mums3zevgzu0x0ecq2" +
|
||||
"f7tp8a05gl0sjq9q9wjg"},
|
||||
AgeSecretKey: "AGE-SECRET-KEY-19CR5YSFW59HM4TLD6GXVEDMZFTVVF7PPHKU" +
|
||||
"T68TXSFPK7APHXA2QS2NJA5",
|
||||
CompressionLevel: 3,
|
||||
Hostname: "test-host",
|
||||
ChunkSize: config.Size(16 * 1024),
|
||||
BlobSizeLimit: config.Size(64 * 1024),
|
||||
}
|
||||
|
||||
working, err := storage.NewFileStorer(storeDir)
|
||||
require.NoError(t, err)
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
// Run 1: the upload is interrupted, so the scan fails but leaves the
|
||||
// interrupted blob's chunk rows committed in the index.
|
||||
_, err = runBackup(ctx, fs, cfg, &interruptedBlobStorer{Storer: working},
|
||||
dataDir, dbPath, "interrupted")
|
||||
require.Error(t, err, "interrupted upload must fail the run")
|
||||
|
||||
// Run 2: retry on the same index with a working backend. This must
|
||||
// succeed and produce a fully restorable snapshot.
|
||||
snapshotID, err := runBackup(ctx, fs, cfg, working, dataDir, dbPath, "retry")
|
||||
require.NoError(t, err, "retry backup must succeed")
|
||||
|
||||
v := &Vaultik{
|
||||
Config: cfg,
|
||||
Storage: working,
|
||||
Fs: fs,
|
||||
Stdout: io.Discard,
|
||||
Stderr: io.Discard,
|
||||
UI: ui.NewWithColor(io.Discard, false),
|
||||
}
|
||||
v.SetContext(ctx)
|
||||
|
||||
require.NoError(t, v.Restore(&RestoreOptions{
|
||||
SnapshotID: snapshotID,
|
||||
TargetDir: restoreDir,
|
||||
}), "the retried snapshot must be restorable")
|
||||
|
||||
for name, data := range sources {
|
||||
restored := filepath.Join(restoreDir, dataDir, name)
|
||||
got, err := afero.ReadFile(fs, restored)
|
||||
require.NoErrorf(t, err, "restored file missing: %s", name)
|
||||
require.Truef(t, bytes.Equal(got, data),
|
||||
"restored bytes differ from original for %s", name)
|
||||
}
|
||||
}
|
||||
|
||||
// runBackup performs one backup of dataDir into a fresh snapshot on the
|
||||
// index at dbPath, returning the snapshot ID. When the scan succeeds it
|
||||
// also completes and exports the snapshot metadata so the result can be
|
||||
// restored. The index database is always closed before returning.
|
||||
func runBackup(
|
||||
ctx context.Context,
|
||||
fs afero.Fs,
|
||||
cfg *config.Config,
|
||||
storer storage.Storer,
|
||||
dataDir, dbPath, name string,
|
||||
) (string, error) {
|
||||
db, err := database.New(ctx, dbPath)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
defer func() { _ = db.Close() }()
|
||||
|
||||
repos := database.NewRepositories(db)
|
||||
|
||||
sm := snapshot.NewSnapshotManager(snapshot.SnapshotManagerParams{
|
||||
Repos: repos,
|
||||
Storage: storer,
|
||||
Config: cfg,
|
||||
})
|
||||
sm.SetFilesystem(fs)
|
||||
|
||||
scanner := snapshot.NewScanner(snapshot.ScannerConfig{
|
||||
FS: fs,
|
||||
Storage: storer,
|
||||
ChunkSize: cfg.ChunkSize.Int64(),
|
||||
MaxBlobSize: cfg.BlobSizeLimit.Int64(),
|
||||
CompressionLevel: cfg.CompressionLevel,
|
||||
AgeRecipients: cfg.AgeRecipients,
|
||||
Repositories: repos,
|
||||
})
|
||||
|
||||
snapshotID, err := sm.CreateSnapshotWithName(
|
||||
ctx, cfg.Hostname, name, "test-version", "test-git")
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
_, err = scanner.Scan(ctx, dataDir, snapshotID)
|
||||
if err != nil {
|
||||
return snapshotID, err
|
||||
}
|
||||
|
||||
err = sm.CompleteSnapshot(ctx, snapshotID)
|
||||
if err != nil {
|
||||
return snapshotID, err
|
||||
}
|
||||
|
||||
err = sm.ExportSnapshotMetadata(ctx, dbPath, snapshotID)
|
||||
if err != nil {
|
||||
return snapshotID, err
|
||||
}
|
||||
|
||||
return snapshotID, nil
|
||||
}
|
||||
@@ -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))
|
||||
}
|
||||
@@ -18,7 +18,6 @@ import (
|
||||
"sneak.berlin/go/vaultik/internal/blobgen"
|
||||
"sneak.berlin/go/vaultik/internal/database"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||
"sneak.berlin/go/vaultik/internal/types"
|
||||
)
|
||||
|
||||
@@ -577,14 +576,20 @@ func (v *Vaultik) handleRestoreVerification(
|
||||
}
|
||||
|
||||
// downloadSnapshotDB downloads and decrypts the snapshot metadata
|
||||
// database. The snapshotID is the human ID; we hash it to the remote
|
||||
// key for the storage path.
|
||||
// database. The identifier is resolved to the snapshot's remote key: a
|
||||
// human ID is hashed, and a remote key (or its abbreviation, as printed
|
||||
// for a remote-only snapshot) is used as-is, so a host with no local
|
||||
// index can restore the snapshots it can only see on the store.
|
||||
func (v *Vaultik) downloadSnapshotDB(
|
||||
snapshotID string, identity age.Identity,
|
||||
) (*database.DB, error) {
|
||||
remoteKey, err := v.resolveSnapshotRemoteKey(snapshotID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Download encrypted database from storage
|
||||
dbKey := fmt.Sprintf("metadata/%s/db.zst.age",
|
||||
snapshot.RemoteSnapshotKey(snapshotID))
|
||||
dbKey := fmt.Sprintf("metadata/%s/db.zst.age", remoteKey)
|
||||
|
||||
reader, err := v.Storage.Get(v.ctx, dbKey)
|
||||
if err != nil {
|
||||
|
||||
@@ -0,0 +1,167 @@
|
||||
package vaultik_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"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/snapshot"
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
"sneak.berlin/go/vaultik/internal/ui"
|
||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
||||
)
|
||||
|
||||
// TestRestoreOnAnotherMachine proves the disaster-recovery path: a host
|
||||
// that has only the vaultik binary, the age secret key, and the storage
|
||||
// credentials — no local index, a different hostname, and no
|
||||
// age_recipients configured — can list, restore, and verify a snapshot
|
||||
// straight from the destination store.
|
||||
//
|
||||
// The backup half writes a snapshot with one index and hostname. The
|
||||
// restore half throws that index away entirely: a fresh, empty index and
|
||||
// a config that shares nothing with the original but the storage location
|
||||
// and the secret key. If restore or verify needed the original local
|
||||
// index — or the human snapshot ID that only that index holds — this test
|
||||
// could not run, because the recovery host can know neither.
|
||||
func TestRestoreOnAnotherMachine(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")
|
||||
restoreDir := filepath.Join(tempDir, "restored")
|
||||
dbPath := filepath.Join(tempDir, "index.sqlite")
|
||||
|
||||
chunkSize := int64(64 * 1024)
|
||||
maxBlobSize := int64(512 * 1024)
|
||||
|
||||
sourceFiles := writeRecoverySourceTree(t, fs, dataDir, chunkSize)
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
// Backup host: one index, hostname test-host, age_recipients set.
|
||||
// runFileStorageBackup closes the index before returning, so nothing
|
||||
// below can lean on it.
|
||||
_, storer, originalID := runFileStorageBackup(
|
||||
ctx, t, fs, dataDir, storeDir, dbPath, chunkSize, maxBlobSize)
|
||||
|
||||
// Recovery host: a fresh empty index, a different hostname, and no
|
||||
// age_recipients — only the secret key and the same storage location.
|
||||
recovery, stdout := newRecoveryHost(ctx, t, fs, storer)
|
||||
|
||||
// The recovery index really is empty. This is the assertion that makes
|
||||
// the test a guard against restore quietly depending on the original
|
||||
// index: if it did, an empty index would make restore fail.
|
||||
localSnaps, err := recovery.Repositories.Snapshots.ListRecent(ctx, 100)
|
||||
require.NoError(t, err)
|
||||
require.Empty(t, localSnaps, "recovery host must start with no local index")
|
||||
|
||||
// List: the snapshot shows up as remote-only, identified by its remote
|
||||
// key, with no recoverable human ID.
|
||||
require.NoError(t, recovery.ListSnapshots(true))
|
||||
|
||||
rows := decodeListJSON(t, stdout.String())
|
||||
require.Len(t, rows, 1)
|
||||
|
||||
remote := rows[0]
|
||||
assert.False(t, remote.LocallyTracked, "snapshot must be remote-only here")
|
||||
assert.Empty(t, remote.ID, "the human ID is unknown to the recovery host")
|
||||
require.Len(t, remote.RemoteKey, 64)
|
||||
assert.Equal(t, snapshot.RemoteSnapshotKey(originalID), remote.RemoteKey,
|
||||
"the listed key is the hashed snapshot ID")
|
||||
|
||||
// Restore driven by the abbreviated identifier the table prints (the
|
||||
// first 12 hex of the remote key), then deep-verify from the store
|
||||
// keyed by the full remote key. Both are what a recovery host can know.
|
||||
require.NoError(t, recovery.Restore(&vaultik.RestoreOptions{
|
||||
SnapshotID: remote.RemoteKey[:12],
|
||||
TargetDir: restoreDir,
|
||||
Verify: true,
|
||||
}))
|
||||
require.NoError(t, recovery.RunDeepVerify(
|
||||
remote.RemoteKey, &vaultik.VerifyOptions{Deep: true}))
|
||||
|
||||
assertRestoredTreeMatches(t, fs, restoreDir, sourceFiles)
|
||||
}
|
||||
|
||||
// writeRecoverySourceTree writes a small source tree spanning several
|
||||
// chunks (so restore reassembles real multi-chunk files) and returns the
|
||||
// content keyed by absolute path.
|
||||
func writeRecoverySourceTree(
|
||||
t *testing.T, fs afero.Fs, dataDir string, chunkSize int64,
|
||||
) map[string][]byte {
|
||||
t.Helper()
|
||||
|
||||
sourceFiles := map[string][]byte{
|
||||
filepath.Join(dataDir, "notes.txt"): []byte("recover me"),
|
||||
filepath.Join(dataDir, "sub", "big.bin"): bytesPattern("big-", int(chunkSize*3)),
|
||||
filepath.Join(dataDir, "sub", "small.bin"): bytesPattern("small-", 128),
|
||||
}
|
||||
|
||||
for path, content := range sourceFiles {
|
||||
require.NoError(t, fs.MkdirAll(filepath.Dir(path), 0o755))
|
||||
require.NoError(t, afero.WriteFile(fs, path, content, 0o644))
|
||||
}
|
||||
|
||||
return sourceFiles
|
||||
}
|
||||
|
||||
// newRecoveryHost builds the Vaultik a replacement machine would run: an
|
||||
// empty in-memory index, a hostname different from the backup host, no
|
||||
// age_recipients, and only the secret key plus the shared storer. It
|
||||
// returns the instance and the buffer its stdout is wired to.
|
||||
func newRecoveryHost(
|
||||
ctx context.Context, t *testing.T, fs afero.Fs, storer storage.Storer,
|
||||
) (*vaultik.Vaultik, *bytes.Buffer) {
|
||||
t.Helper()
|
||||
|
||||
recoveryDB, err := database.New(ctx, ":memory:")
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() { _ = recoveryDB.Close() })
|
||||
|
||||
stdout := &bytes.Buffer{}
|
||||
|
||||
recovery := &vaultik.Vaultik{
|
||||
Config: &config.Config{
|
||||
AgeSecretKey: testAgeSecretKey,
|
||||
Hostname: "recovery-host",
|
||||
},
|
||||
Storage: storer,
|
||||
Fs: fs,
|
||||
Repositories: database.NewRepositories(recoveryDB),
|
||||
DB: recoveryDB,
|
||||
Stdout: stdout,
|
||||
Stderr: io.Discard,
|
||||
UI: ui.NewWithColor(io.Discard, false),
|
||||
}
|
||||
recovery.SetContext(ctx)
|
||||
|
||||
return recovery, stdout
|
||||
}
|
||||
|
||||
// assertRestoredTreeMatches byte-compares every restored file against its
|
||||
// source content.
|
||||
func assertRestoredTreeMatches(
|
||||
t *testing.T, fs afero.Fs, restoreDir string, sourceFiles map[string][]byte,
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
for origPath, expected := range sourceFiles {
|
||||
restored := filepath.Join(restoreDir, origPath)
|
||||
got, err := afero.ReadFile(fs, restored)
|
||||
require.NoErrorf(t, err, "restored file missing: %s", restored)
|
||||
require.Truef(t, bytes.Equal(got, expected),
|
||||
"byte mismatch for %s", origPath)
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -669,9 +670,11 @@ func (v *Vaultik) VerifySnapshotWithOptions(
|
||||
|
||||
v.printVerifyHeader(snapshotID, opts)
|
||||
|
||||
// Download and parse manifest. The caller supplies a human
|
||||
// snapshot ID; we hash it to address remote storage.
|
||||
manifest, err := v.downloadManifestByKey(snapshot.RemoteSnapshotKey(snapshotID))
|
||||
// Resolve the identifier to the snapshot's remote key and download the
|
||||
// manifest. A human ID is hashed; a remote key (or its abbreviation,
|
||||
// as printed for a remote-only snapshot) is used as-is, so a host with
|
||||
// no local index can verify a snapshot it can only see on the store.
|
||||
manifest, err := v.resolveAndDownloadManifest(snapshotID)
|
||||
if err != nil {
|
||||
if opts.JSON {
|
||||
result.Status = verifyStatusFailed
|
||||
@@ -1540,12 +1543,17 @@ func (v *Vaultik) outputRemoveJSON(result *RemoveResult) error {
|
||||
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 {
|
||||
SnapshotsDeleted int64
|
||||
FilesDeleted int64
|
||||
ChunksDeleted int64
|
||||
BlobsDeleted int64
|
||||
FilesDeleted *int64
|
||||
ChunksDeleted *int64
|
||||
BlobsDeleted *int64
|
||||
}
|
||||
|
||||
// PruneDatabase removes incomplete snapshots and orphaned files, chunks,
|
||||
@@ -1560,7 +1568,7 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
||||
result := &PruneResult{}
|
||||
|
||||
// Snapshot counts before deletion of incompletes.
|
||||
snapshotCountBefore, _ := v.getTableCount("snapshots")
|
||||
snapshotCountBefore := v.tableCountForReport("snapshots")
|
||||
|
||||
// First, delete any incomplete snapshots
|
||||
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
|
||||
@@ -1575,9 +1583,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
||||
}
|
||||
|
||||
// Get counts before cleanup for reporting
|
||||
fileCountBefore, _ := v.getTableCount("files")
|
||||
chunkCountBefore, _ := v.getTableCount("chunks")
|
||||
blobCountBefore, _ := v.getTableCount("blobs")
|
||||
fileCountBefore := v.tableCountForReport("files")
|
||||
chunkCountBefore := v.tableCountForReport("chunks")
|
||||
blobCountBefore := v.tableCountForReport("blobs")
|
||||
|
||||
// Run the cleanup
|
||||
err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
|
||||
@@ -1586,36 +1594,83 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
||||
}
|
||||
|
||||
// Get counts after cleanup
|
||||
fileCountAfter, _ := v.getTableCount("files")
|
||||
chunkCountAfter, _ := v.getTableCount("chunks")
|
||||
blobCountAfter, _ := v.getTableCount("blobs")
|
||||
fileCountAfter := v.tableCountForReport("files")
|
||||
chunkCountAfter := v.tableCountForReport("chunks")
|
||||
blobCountAfter := v.tableCountForReport("blobs")
|
||||
|
||||
result.FilesDeleted = fileCountBefore - fileCountAfter
|
||||
result.ChunksDeleted = chunkCountBefore - chunkCountAfter
|
||||
result.BlobsDeleted = blobCountBefore - blobCountAfter
|
||||
result.FilesDeleted = countDiff(fileCountBefore, fileCountAfter)
|
||||
result.ChunksDeleted = countDiff(chunkCountBefore, chunkCountAfter)
|
||||
result.BlobsDeleted = countDiff(blobCountBefore, blobCountAfter)
|
||||
|
||||
log.Info("Local database prune complete",
|
||||
"incomplete_snapshots", result.SnapshotsDeleted,
|
||||
"orphaned_files", result.FilesDeleted,
|
||||
"orphaned_chunks", result.ChunksDeleted,
|
||||
"orphaned_blobs", result.BlobsDeleted,
|
||||
"orphaned_files", countText(result.FilesDeleted),
|
||||
"orphaned_chunks", countText(result.ChunksDeleted),
|
||||
"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.Detailf("Incomplete snapshots: %d removed (%d remain).",
|
||||
result.SnapshotsDeleted, snapshotCountAfter)
|
||||
v.UI.Detailf("Orphaned files: %d removed (%d remain).",
|
||||
result.FilesDeleted, fileCountAfter)
|
||||
v.UI.Detailf("Orphaned chunks: %d removed (%d remain).",
|
||||
result.ChunksDeleted, chunkCountAfter)
|
||||
v.UI.Detailf("Orphaned blobs: %d removed (%d remain).",
|
||||
result.BlobsDeleted, blobCountAfter)
|
||||
v.UI.Detailf("Incomplete snapshots: %s removed (%s remain).",
|
||||
countText(&result.SnapshotsDeleted), countText(snapshotsRemain))
|
||||
v.UI.Detailf("Orphaned files: %s removed (%s remain).",
|
||||
countText(result.FilesDeleted), countText(fileCountAfter))
|
||||
v.UI.Detailf("Orphaned chunks: %s removed (%s remain).",
|
||||
countText(result.ChunksDeleted), countText(chunkCountAfter))
|
||||
v.UI.Detailf("Orphaned blobs: %s removed (%s remain).",
|
||||
countText(result.BlobsDeleted), countText(blobCountAfter))
|
||||
|
||||
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
|
||||
// alphanumeric characters and underscores.
|
||||
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
|
||||
|
||||
@@ -0,0 +1,101 @@
|
||||
package vaultik
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||
)
|
||||
|
||||
// remoteKeyHexLen is the length of a full remote snapshot key: a SHA256
|
||||
// digest rendered as lowercase hex.
|
||||
const remoteKeyHexLen = 64
|
||||
|
||||
// Sentinel errors for resolving a snapshot identifier against the store.
|
||||
var (
|
||||
errSnapshotKeyNotFound = errors.New(
|
||||
"no snapshot on the destination store matches this identifier")
|
||||
errSnapshotKeyAmbiguous = errors.New(
|
||||
"identifier matches more than one snapshot on the destination store")
|
||||
)
|
||||
|
||||
// resolveSnapshotRemoteKey turns a snapshot identifier supplied on the
|
||||
// command line into the remote key that names the snapshot's metadata
|
||||
// directory on the destination store. Every remote path a restore or
|
||||
// verify reads is built from that key.
|
||||
//
|
||||
// Two forms are accepted, matching the two things a host can know:
|
||||
//
|
||||
// - A human snapshot ID (hostname_name_timestamp), which a host holding
|
||||
// the local index has. It is hashed to its remote key; the store is
|
||||
// not consulted.
|
||||
// - A remote key, or the leading part of one, which is all a host with
|
||||
// no local index can know — it is exactly what `snapshot list` prints
|
||||
// for a remote-only snapshot (see formatRemoteOnlyID). It is resolved
|
||||
// against the destination store's metadata listing; an identifier that
|
||||
// matches no snapshot, or more than one, is an error.
|
||||
//
|
||||
// The two are told apart by shape: a remote key is lowercase hex, and a
|
||||
// human snapshot ID never is (it carries a hostname, underscores, and an
|
||||
// RFC3339 timestamp).
|
||||
func (v *Vaultik) resolveSnapshotRemoteKey(identifier string) (string, error) {
|
||||
if !isRemoteKeyOrPrefix(identifier) {
|
||||
return snapshot.RemoteSnapshotKey(identifier), nil
|
||||
}
|
||||
|
||||
keys, err := v.listAllRemoteSnapshotKeys()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf(
|
||||
"listing destination store to resolve %q: %w", identifier, err)
|
||||
}
|
||||
|
||||
var matches []string
|
||||
|
||||
for _, key := range keys {
|
||||
if strings.HasPrefix(key, identifier) {
|
||||
matches = append(matches, key)
|
||||
}
|
||||
}
|
||||
|
||||
switch len(matches) {
|
||||
case 1:
|
||||
return matches[0], nil
|
||||
case 0:
|
||||
return "", fmt.Errorf("%w: %s", errSnapshotKeyNotFound, identifier)
|
||||
default:
|
||||
return "", fmt.Errorf("%w: %s (%d matches)",
|
||||
errSnapshotKeyAmbiguous, identifier, len(matches))
|
||||
}
|
||||
}
|
||||
|
||||
// resolveAndDownloadManifest resolves a snapshot identifier to its remote
|
||||
// key (see resolveSnapshotRemoteKey) and downloads that snapshot's
|
||||
// manifest.
|
||||
func (v *Vaultik) resolveAndDownloadManifest(
|
||||
identifier string,
|
||||
) (*snapshot.Manifest, error) {
|
||||
remoteKey, err := v.resolveSnapshotRemoteKey(identifier)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return v.downloadManifestByKey(remoteKey)
|
||||
}
|
||||
|
||||
// isRemoteKeyOrPrefix reports whether s is a full remote key or the
|
||||
// leading part of one: 1 to 64 lowercase hex characters. A human snapshot
|
||||
// ID is never all hex, so this shape test is enough to tell the two apart.
|
||||
func isRemoteKeyOrPrefix(s string) bool {
|
||||
if s == "" || len(s) > remoteKeyHexLen {
|
||||
return false
|
||||
}
|
||||
|
||||
for _, r := range s {
|
||||
if (r < '0' || r > '9') && (r < 'a' || r > 'f') {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
@@ -138,8 +138,15 @@ func (v *Vaultik) RunDeepVerify(snapshotID string, opts *VerifyOptions) error {
|
||||
func (v *Vaultik) loadVerificationData(
|
||||
snapshotID string, opts *VerifyOptions, result *VerifyResult,
|
||||
) (*snapshot.Manifest, *tempDB, []snapshot.BlobInfo, error) {
|
||||
// All remote paths use the hashed key derived from the human ID.
|
||||
remoteKey := snapshot.RemoteSnapshotKey(snapshotID)
|
||||
// Resolve the identifier to the snapshot's remote key. A human ID is
|
||||
// hashed; a remote key (or its abbreviation, as printed for a
|
||||
// remote-only snapshot) is used as-is, so a host with no local index
|
||||
// can verify a snapshot it can only see on the store.
|
||||
remoteKey, err := v.resolveSnapshotRemoteKey(snapshotID)
|
||||
if err != nil {
|
||||
return nil, nil, nil, v.deepVerifyFailure(result, opts,
|
||||
fmt.Sprintf("resolving snapshot identifier: %v", err), err)
|
||||
}
|
||||
|
||||
// Download manifest. downloadManifestByKey is the single reader for
|
||||
// remote manifests; see its doc comment.
|
||||
@@ -186,7 +193,7 @@ func (v *Vaultik) loadVerificationData(
|
||||
fmt.Errorf("failed to decrypt database: %w", err))
|
||||
}
|
||||
|
||||
dbBlobs, err := v.getBlobsFromDatabase(snapshotID, tdb.DB)
|
||||
dbBlobs, err := v.getBlobsFromDatabase(tdb.DB)
|
||||
if err != nil {
|
||||
_ = tdb.Close()
|
||||
|
||||
@@ -501,19 +508,21 @@ func (v *Vaultik) verifyBlobFinalIntegrity(
|
||||
return nil
|
||||
}
|
||||
|
||||
// getBlobsFromDatabase gets all blobs for the snapshot from the database
|
||||
func (v *Vaultik) getBlobsFromDatabase(
|
||||
snapshotID string, db *sql.DB,
|
||||
) ([]snapshot.BlobInfo, error) {
|
||||
// getBlobsFromDatabase gets all blobs for the snapshot from the database.
|
||||
//
|
||||
// The exported per-snapshot database holds exactly one snapshot's data
|
||||
// (see cleanSnapshotDB), so every row in snapshot_blobs belongs to it.
|
||||
// We select them directly rather than filtering by the human snapshot ID,
|
||||
// which a host restoring from the store alone does not have.
|
||||
func (v *Vaultik) getBlobsFromDatabase(db *sql.DB) ([]snapshot.BlobInfo, error) {
|
||||
query := `
|
||||
SELECT b.blob_hash, b.compressed_size
|
||||
FROM snapshot_blobs sb
|
||||
JOIN blobs b ON sb.blob_hash = b.blob_hash
|
||||
WHERE sb.snapshot_id = ?
|
||||
ORDER BY b.blob_hash
|
||||
`
|
||||
|
||||
rows, err := db.QueryContext(v.ctx, query, snapshotID)
|
||||
rows, err := db.QueryContext(v.ctx, query)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to query snapshot blobs: %w", err)
|
||||
}
|
||||
|
||||
+16
-1
@@ -56,8 +56,23 @@ main() {
|
||||
docker build --output=type=cacheonly \
|
||||
--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)$$"
|
||||
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 "$@"
|
||||
|
||||
@@ -24,7 +24,24 @@ main() {
|
||||
# whether the tree is clean. The Dockerfile now refuses to build
|
||||
# without a non-empty value, so this is required, not optional.
|
||||
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" \
|
||||
--build-arg VERSION="$version" \
|
||||
--build-arg COMMIT="$commit" \
|
||||
--build-arg COMMIT_DATE="$commit_date" \
|
||||
-t "$("$SCRIPT_DIR/projectname")" .
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user