Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d4df9701f6 |
@@ -104,12 +104,7 @@ Version: 2025-06-08
|
|||||||
|
|
||||||
13. Pre-1.0: NEVER write database migrations. There are no live databases
|
13. Pre-1.0: NEVER write database migrations. There are no live databases
|
||||||
anywhere — every user's local index can be rebuilt from a fresh full
|
anywhere — every user's local index can be rebuilt from a fresh full
|
||||||
backup. To change the schema, edit `internal/database/schema/001.sql`
|
backup. When the schema changes, just change `schema.sql` (and any code
|
||||||
(and any code that touches the affected tables) directly; do not add new
|
that touches the affected tables). The local index is disposable until
|
||||||
numbered schema files. Those numbered files and the `schema_migrations`
|
1.0 ships and is tagged.
|
||||||
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.
|
|
||||||
|
|
||||||
|
|||||||
@@ -84,57 +84,6 @@ VAULTIK_AGE_SECRET_KEY='AGE-SECRET-KEY-...' vaultik snapshot restore <snapshot-i
|
|||||||
# 0 3 * * * vaultik snapshot create --cron --prune --keep-newer-than 4w
|
# 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
|
## cli
|
||||||
@@ -296,16 +245,13 @@ local index alone, and still exits zero.
|
|||||||
* Default (shallow): checks that all blobs referenced in the manifest exist in storage
|
* Default (shallow): checks that all blobs referenced in the manifest exist in storage
|
||||||
* `--deep`: Downloads and decrypts each blob, verifies chunk hashes against the
|
* `--deep`: Downloads and decrypts each blob, verifies chunk hashes against the
|
||||||
encrypted metadata database
|
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
|
* `--json`: Output results as JSON
|
||||||
|
|
||||||
**`snapshot purge`**: Remove old snapshots based on criteria. Retention is
|
**`snapshot purge`**: Remove old snapshots based on criteria. Retention is
|
||||||
per-snapshot-name (`--keep-latest` keeps the latest of each name, not the
|
per-snapshot-name (`--keep-latest` keeps the latest of each name, not the
|
||||||
latest globally).
|
latest globally).
|
||||||
* `--keep-latest`: Keep only the most recent snapshot of each name
|
* `--keep-latest`: Keep only the most recent snapshot of each name
|
||||||
* `--older-than <duration>`: Remove snapshots older than duration (e.g. `30d`,
|
* `--older-than <duration>`: Remove snapshots older than duration (e.g. `30d`, `6m`, `1y`)
|
||||||
`4w`, `6mo`, `1y`; `m` is minutes, `mo` is months)
|
|
||||||
* `--snapshot <name>`: Restrict to specific snapshot names (repeat for multiple)
|
* `--snapshot <name>`: Restrict to specific snapshot names (repeat for multiple)
|
||||||
* `--force`: Skip confirmation prompt
|
* `--force`: Skip confirmation prompt
|
||||||
|
|
||||||
@@ -328,10 +274,6 @@ on the destination in one go, use `vaultik remote nuke --force`.
|
|||||||
|
|
||||||
**`snapshot restore`**: Restore files from a backup snapshot.
|
**`snapshot restore`**: Restore files from a backup snapshot.
|
||||||
* Requires `VAULTIK_AGE_SECRET_KEY` environment variable
|
* 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)
|
* Optional path arguments to restore specific files/directories (default: all)
|
||||||
* Preserves file permissions, timestamps, ownership (ownership requires root),
|
* Preserves file permissions, timestamps, ownership (ownership requires root),
|
||||||
symlinks, and empty directories
|
symlinks, and empty directories
|
||||||
@@ -514,13 +456,9 @@ Key fields:
|
|||||||
sequentially. Restore speed is bound by single-stream throughput.
|
sequentially. Restore speed is bound by single-stream throughput.
|
||||||
* **Device nodes, named pipes, and sockets are silently skipped.** Only
|
* **Device nodes, named pipes, and sockets are silently skipped.** Only
|
||||||
regular files, directories, and symlinks are backed up.
|
regular files, directories, and symlinks are backed up.
|
||||||
* **No upgrade path between versions.** There is no supported way to carry
|
* **No database migrations.** If the local SQLite schema changes between
|
||||||
an existing local index across a schema change; if the local SQLite
|
versions, delete the local database (`vaultik database delete`) and run
|
||||||
schema changes between versions, delete the local database (`vaultik
|
a full backup. Remote storage is unaffected.
|
||||||
database delete`) and run a full backup. Remote storage is unaffected.
|
|
||||||
(The binary does embed numbered schema files and a `schema_migrations`
|
|
||||||
table to bootstrap a fresh database — see [`docs/DATAMODEL.md`](docs/DATAMODEL.md)
|
|
||||||
— but that is not an upgrade path.)
|
|
||||||
* **Files that change during backup may be inconsistent.** There is no
|
* **Files that change during backup may be inconsistent.** There is no
|
||||||
filesystem snapshot or freeze. If a file is modified between the scan
|
filesystem snapshot or freeze. If a file is modified between the scan
|
||||||
and chunk phases, the backed-up copy may reflect a partial write.
|
and chunk phases, the backed-up copy may reflect a partial write.
|
||||||
@@ -586,12 +524,14 @@ priority.
|
|||||||
|
|
||||||
### infrastructure
|
### infrastructure
|
||||||
|
|
||||||
* **Cross-version schema upgrades.** There is no upgrade path between
|
* **Cross-machine restore documentation.** The "restore from
|
||||||
released versions — pre-1.0 schema changes are handled by `vaultik
|
another host" workflow works but isn't documented as a
|
||||||
database delete` plus a full re-scan (see
|
first-class operation in this README. Worth a dedicated section
|
||||||
[`docs/DATAMODEL.md`](docs/DATAMODEL.md)). Post-1.0 we'll need a
|
once it's settled.
|
||||||
migration story to keep existing index databases usable across
|
* **Schema migrations.** Currently nonexistent — pre-1.0 schema
|
||||||
upgrades.
|
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.
|
||||||
* **Storage backend coverage tests.** S3, file://, and rclone://
|
* **Storage backend coverage tests.** S3, file://, and rclone://
|
||||||
all share the Storer interface but the rclone path is the least
|
all share the Storer interface but the rclone path is the least
|
||||||
exercised in CI.
|
exercised in CI.
|
||||||
|
|||||||
@@ -25,34 +25,6 @@ release" is exactly the contradiction
|
|||||||
|
|
||||||
# Completed Steps
|
# Completed Steps
|
||||||
|
|
||||||
- 2026-09-21: Stopped `prune` from reporting a failed row count as 0
|
|
||||||
([issue #96](https://git.eeqj.de/sneak/vaultik/issues/96)). The seven
|
|
||||||
`getTableCount` reads in `PruneDatabase` discarded their error, so a
|
|
||||||
query that could not run became a plausible `0` and the before/after
|
|
||||||
delta computed from it looked like real work. Each read now logs at
|
|
||||||
warn on failure and renders as `unknown`, never `0`, so an empty table
|
|
||||||
is distinguishable from one that could not be queried. The counts have
|
|
||||||
no `--json` representation — under `--json` the summary is suppressed
|
|
||||||
entirely — so nothing there can show a false `0`.
|
|
||||||
|
|
||||||
- 2026-09-21: Made the s3 storage backend report a missing object as
|
|
||||||
`storage.ErrNotFound`, like the `file` and `rclone` backends and as the
|
|
||||||
`Storer` interface documents. `S3Storer.Get` and `Stat` returned the raw
|
|
||||||
AWS SDK error, so `errors.Is(err, storage.ErrNotFound)` was false on s3
|
|
||||||
and callers branched differently per backend. Added a small `s3.IsNotFound`
|
|
||||||
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
|
|
||||||
ID, which is the double SHA256 of the plaintext, so the two could
|
|
||||||
never match. It now hashes the decompressed plaintext and compares the
|
|
||||||
double SHA256. Added a test that backs up a real snapshot, deep-verifies
|
|
||||||
it, then flips a byte in one stored blob and confirms deep verification
|
|
||||||
then fails
|
|
||||||
([issue #131](https://git.eeqj.de/sneak/vaultik/issues/131)).
|
|
||||||
|
|
||||||
- 2026-09-21: Made `snapshot create` VACUUM the per-snapshot metadata
|
- 2026-09-21: Made `snapshot create` VACUUM the per-snapshot metadata
|
||||||
database through the `modernc.org/sqlite` driver instead of shelling
|
database through the `modernc.org/sqlite` driver instead of shelling
|
||||||
out to the external `sqlite` command-line binary (issue #120). A
|
out to the external `sqlite` command-line binary (issue #120). A
|
||||||
@@ -78,24 +50,6 @@ release" is exactly the contradiction
|
|||||||
keeps that exact compiler from auto-switching. Bumping Go now touches
|
keeps that exact compiler from auto-switching. Bumping Go now touches
|
||||||
`go.mod`, the checksum, and the `Dockerfile` `golang` digest together.
|
`go.mod`, the checksum, and the `Dockerfile` `golang` digest together.
|
||||||
|
|
||||||
- 2026-09-21: Collapsed the two duration parsers into one and fixed the
|
|
||||||
`--older-than` months example
|
|
||||||
([issue #123](https://git.eeqj.de/sneak/vaultik/issues/123)). Two
|
|
||||||
functions named `parseDuration` existed with different grammars;
|
|
||||||
`snapshot purge --older-than` and `--keep-newer-than` both already went
|
|
||||||
through the one in `internal/vaultik`, while the richer copy in
|
|
||||||
`internal/cli/duration.go` was reachable only from its own test. Kept
|
|
||||||
the live-path parser and deleted the unused one, so no flag's accepted
|
|
||||||
grammar changes. The trap the issue was filed over: `README.md`
|
|
||||||
documented `6m` as the months example for `--older-than`, but `m` is
|
|
||||||
minutes, so the documented command deleted every snapshot older than
|
|
||||||
six minutes on a destructive flag. Corrected the doc to `6mo` and put
|
|
||||||
both flags' help text on one example list that states `m` is minutes
|
|
||||||
and `mo` is months. The surviving parser now rejects negatives, which
|
|
||||||
it previously accepted (`-5h`) or silently made positive (`-5d`).
|
|
||||||
Table-driven tests cover every unit, `6m` as six minutes, `6mo` as 180
|
|
||||||
days, and rejection of a bare number, an unknown unit, and a negative.
|
|
||||||
|
|
||||||
- 2026-08-10: Moved every lint run into its own container, as a build
|
- 2026-08-10: Moved every lint run into its own container, as a build
|
||||||
step ([issue #113](https://git.eeqj.de/sneak/vaultik/issues/113)).
|
step ([issue #113](https://git.eeqj.de/sneak/vaultik/issues/113)).
|
||||||
New root `Dockerfile.lint`, built by `script/lint`, runs
|
New root `Dockerfile.lint`, built by `script/lint`, runs
|
||||||
|
|||||||
+5
-24
@@ -5,30 +5,11 @@
|
|||||||
Vaultik uses a local SQLite database to track file metadata, chunk mappings, and blob associations during the backup process. This database serves as an index for incremental backups and enables efficient deduplication.
|
Vaultik uses a local SQLite database to track file metadata, chunk mappings, and blob associations during the backup process. This database serves as an index for incremental backups and enables efficient deduplication.
|
||||||
|
|
||||||
**Important Notes:**
|
**Important Notes:**
|
||||||
|
- **No Migration Support (pre-1.0)**: Vaultik does not support database schema
|
||||||
This section is the authoritative explanation of the schema/migration story;
|
migrations. The local index is treated as disposable — if the schema changes,
|
||||||
other documents (the README and `AGENTS.md`) link here.
|
delete the local SQLite database (`vaultik database delete`) and run a full
|
||||||
|
backup. The remote storage is unaffected; the new index will re-deduplicate
|
||||||
- **No upgrade path between versions (pre-1.0)**: Vaultik has no supported way to
|
against existing remote blobs.
|
||||||
carry an existing local index across a schema change. The index is disposable
|
|
||||||
— if the on-disk schema changes between versions, delete the local SQLite
|
|
||||||
database (`vaultik database delete`) and run a full backup. Remote storage is
|
|
||||||
unaffected; the new index re-deduplicates against existing remote blobs. This
|
|
||||||
is the standing project policy, and it is separate from the schema bootstrap
|
|
||||||
described next.
|
|
||||||
- **Schema bootstrap**: a fresh database is populated from numbered SQL files
|
|
||||||
embedded in the binary under `internal/database/schema/`. `000.sql` creates the
|
|
||||||
`schema_migrations` table; `001.sql` creates the application tables. On opening
|
|
||||||
a database the code applies each numbered file that has not yet run and records
|
|
||||||
its version in `schema_migrations`. This bootstraps a new database; it does not
|
|
||||||
upgrade an existing one between released versions.
|
|
||||||
- **Changing the schema (pre-1.0)**: edit `internal/database/schema/001.sql` (and
|
|
||||||
the code that touches the affected tables) directly. Do not add new numbered
|
|
||||||
files — there is no installed base to migrate.
|
|
||||||
- **Disposability expires at 1.0**: the index is treated as disposable only until
|
|
||||||
1.0 ships and is tagged. Once 1.0 is tagged that clause expires and the
|
|
||||||
question of upgrading existing indexes returns. It is deliberately left open
|
|
||||||
here.
|
|
||||||
- **Version Compatibility**: In rare cases, you may need to use the same version
|
- **Version Compatibility**: In rare cases, you may need to use the same version
|
||||||
of Vaultik to restore a backup as was used to create it. This ensures
|
of Vaultik to restore a backup as was used to create it. This ensures
|
||||||
compatibility with the metadata format stored in S3.
|
compatibility with the metadata format stored in S3.
|
||||||
|
|||||||
@@ -0,0 +1,126 @@
|
|||||||
|
package cli
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"regexp"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Approximate lengths of the extended calendar units accepted by
|
||||||
|
// parseDuration.
|
||||||
|
const (
|
||||||
|
durationDay = 24 * time.Hour
|
||||||
|
durationWeek = 7 * durationDay
|
||||||
|
durationMonth = 30 * durationDay
|
||||||
|
durationYear = 365 * durationDay
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
errNegativeDuration = errors.New("negative durations are not supported")
|
||||||
|
errInvalidDuration = errors.New("invalid duration format")
|
||||||
|
errUnknownTimeUnit = errors.New("unknown time unit")
|
||||||
|
)
|
||||||
|
|
||||||
|
// parseDuration parses duration strings. Supports standard Go duration format
|
||||||
|
// (e.g., "3h30m", "1h45m30s") as well as extended units:
|
||||||
|
// - d: days (e.g., "30d", "7d")
|
||||||
|
// - w: weeks (e.g., "2w", "4w")
|
||||||
|
// - mo: months (30 days) (e.g., "6mo", "1mo")
|
||||||
|
// - y: years (365 days) (e.g., "1y", "2y")
|
||||||
|
//
|
||||||
|
// Can combine units: "1y6mo", "2w3d", "1d12h30m"
|
||||||
|
func parseDuration(s string) (time.Duration, error) {
|
||||||
|
// First try standard Go duration parsing
|
||||||
|
d, err := time.ParseDuration(s)
|
||||||
|
if err == nil {
|
||||||
|
return d, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Extended duration parsing
|
||||||
|
// Check for negative values
|
||||||
|
if strings.HasPrefix(strings.TrimSpace(s), "-") {
|
||||||
|
return 0, errNegativeDuration
|
||||||
|
}
|
||||||
|
|
||||||
|
// Pattern matches: number + unit, repeated
|
||||||
|
re := regexp.MustCompile(`(\d+(?:\.\d+)?)\s*([a-zA-Z]+)`)
|
||||||
|
matches := re.FindAllStringSubmatch(s, -1)
|
||||||
|
|
||||||
|
if len(matches) == 0 {
|
||||||
|
return 0, fmt.Errorf("%w: %q", errInvalidDuration, s)
|
||||||
|
}
|
||||||
|
|
||||||
|
var total time.Duration
|
||||||
|
|
||||||
|
for _, match := range matches {
|
||||||
|
valueStr := match[1]
|
||||||
|
unit := strings.ToLower(match[2])
|
||||||
|
|
||||||
|
value, err := strconv.ParseFloat(valueStr, 64)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("invalid number %q: %w", valueStr, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
d, err := durationForUnit(value, unit)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
|
||||||
|
total += d
|
||||||
|
}
|
||||||
|
|
||||||
|
return total, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// durationForUnit converts a value with a (case-normalized) unit suffix
|
||||||
|
// into a time.Duration, accepting Go's standard units plus the extended
|
||||||
|
// calendar units.
|
||||||
|
func durationForUnit(value float64, unit string) (time.Duration, error) {
|
||||||
|
switch unit {
|
||||||
|
// Standard time units
|
||||||
|
case "ns", "nanosecond", "nanoseconds":
|
||||||
|
return time.Duration(value), nil
|
||||||
|
case "us", "µs", "microsecond", "microseconds":
|
||||||
|
return time.Duration(value * float64(time.Microsecond)), nil
|
||||||
|
case "ms", "millisecond", "milliseconds":
|
||||||
|
return time.Duration(value * float64(time.Millisecond)), nil
|
||||||
|
case "s", "sec", "second", "seconds":
|
||||||
|
return time.Duration(value * float64(time.Second)), nil
|
||||||
|
case "m", "min", "minute", "minutes":
|
||||||
|
return time.Duration(value * float64(time.Minute)), nil
|
||||||
|
case "h", "hr", "hour", "hours":
|
||||||
|
return time.Duration(value * float64(time.Hour)), nil
|
||||||
|
// Extended units
|
||||||
|
case "d", "day", "days":
|
||||||
|
return time.Duration(value * float64(durationDay)), nil
|
||||||
|
case "w", "week", "weeks":
|
||||||
|
return time.Duration(value * float64(durationWeek)), nil
|
||||||
|
case "mo", "month", "months":
|
||||||
|
// Using 30 days as approximation
|
||||||
|
return time.Duration(value * float64(durationMonth)), nil
|
||||||
|
case "y", "year", "years":
|
||||||
|
// Using 365 days as approximation
|
||||||
|
return time.Duration(value * float64(durationYear)), nil
|
||||||
|
default:
|
||||||
|
// Try parsing as standard Go duration unit
|
||||||
|
testStr := "1" + unit
|
||||||
|
|
||||||
|
_, err := time.ParseDuration(testStr)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("%w: %q", errUnknownTimeUnit, unit)
|
||||||
|
}
|
||||||
|
|
||||||
|
// It's a valid Go duration unit, parse the full value
|
||||||
|
fullStr := fmt.Sprintf("%g%s", value, unit)
|
||||||
|
|
||||||
|
d, err := time.ParseDuration(fullStr)
|
||||||
|
if err != nil {
|
||||||
|
return 0, fmt.Errorf("invalid duration %q: %w", fullStr, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return d, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,299 @@
|
|||||||
|
package cli //nolint:testpackage // needs access to unexported parseDuration
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
type parseDurationCase struct {
|
||||||
|
name string
|
||||||
|
input string
|
||||||
|
expected time.Duration
|
||||||
|
wantErr bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// runParseDurationCases executes a table of parseDuration cases as
|
||||||
|
// parallel subtests.
|
||||||
|
func runParseDurationCases(t *testing.T, tests []parseDurationCase) {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
got, err := parseDuration(tt.input)
|
||||||
|
|
||||||
|
if tt.wantErr {
|
||||||
|
require.Error(t, err, "expected error for input %q", tt.input)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
require.NoError(t, err, "unexpected error for input %q", tt.input)
|
||||||
|
assert.Equal(t, tt.expected, got, "duration mismatch for input %q", tt.input)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDurationStandard(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
runParseDurationCases(t, []parseDurationCase{
|
||||||
|
{
|
||||||
|
name: "standard seconds",
|
||||||
|
input: "30s",
|
||||||
|
expected: 30 * time.Second,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "standard minutes",
|
||||||
|
input: "45m",
|
||||||
|
expected: 45 * time.Minute,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "standard hours",
|
||||||
|
input: "2h",
|
||||||
|
expected: 2 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "standard combined",
|
||||||
|
input: "3h30m",
|
||||||
|
expected: 3*time.Hour + 30*time.Minute,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "standard complex",
|
||||||
|
input: "1h45m30s",
|
||||||
|
expected: 1*time.Hour + 45*time.Minute + 30*time.Second,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "standard with milliseconds",
|
||||||
|
input: "1s500ms",
|
||||||
|
expected: 1*time.Second + 500*time.Millisecond,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDurationExtendedUnits(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
runParseDurationCases(t, []parseDurationCase{
|
||||||
|
// Extended units - days
|
||||||
|
{
|
||||||
|
name: "single day",
|
||||||
|
input: "1d",
|
||||||
|
expected: 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "multiple days",
|
||||||
|
input: "7d",
|
||||||
|
expected: 7 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "fractional days",
|
||||||
|
input: "1.5d",
|
||||||
|
expected: 36 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "days spelled out",
|
||||||
|
input: "3days",
|
||||||
|
expected: 3 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
// Extended units - weeks
|
||||||
|
{
|
||||||
|
name: "single week",
|
||||||
|
input: "1w",
|
||||||
|
expected: 7 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "multiple weeks",
|
||||||
|
input: "4w",
|
||||||
|
expected: 4 * 7 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "weeks spelled out",
|
||||||
|
input: "2weeks",
|
||||||
|
expected: 2 * 7 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
// Extended units - months
|
||||||
|
{
|
||||||
|
name: "single month",
|
||||||
|
input: "1mo",
|
||||||
|
expected: 30 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "multiple months",
|
||||||
|
input: "6mo",
|
||||||
|
expected: 6 * 30 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "months spelled out",
|
||||||
|
input: "3months",
|
||||||
|
expected: 3 * 30 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
// Extended units - years
|
||||||
|
{
|
||||||
|
name: "single year",
|
||||||
|
input: "1y",
|
||||||
|
expected: 365 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "multiple years",
|
||||||
|
input: "2y",
|
||||||
|
expected: 2 * 365 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "years spelled out",
|
||||||
|
input: "1year",
|
||||||
|
expected: 365 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDurationCombinedAndErrors(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
runParseDurationCases(t, []parseDurationCase{
|
||||||
|
// Combined extended units
|
||||||
|
{
|
||||||
|
name: "weeks and days",
|
||||||
|
input: "2w3d",
|
||||||
|
expected: 2*7*24*time.Hour + 3*24*time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "years and months",
|
||||||
|
input: "1y6mo",
|
||||||
|
expected: 365*24*time.Hour + 6*30*24*time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "days and hours",
|
||||||
|
input: "1d12h",
|
||||||
|
expected: 24*time.Hour + 12*time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "complex combination",
|
||||||
|
input: "1y2mo3w4d5h6m7s",
|
||||||
|
expected: 365*24*time.Hour + 2*30*24*time.Hour +
|
||||||
|
3*7*24*time.Hour + 4*24*time.Hour +
|
||||||
|
5*time.Hour + 6*time.Minute + 7*time.Second,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "with spaces",
|
||||||
|
input: "1d 12h 30m",
|
||||||
|
expected: 24*time.Hour + 12*time.Hour + 30*time.Minute,
|
||||||
|
},
|
||||||
|
// Edge cases
|
||||||
|
{
|
||||||
|
name: "zero duration",
|
||||||
|
input: "0s",
|
||||||
|
expected: 0,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "large duration",
|
||||||
|
input: "10y",
|
||||||
|
expected: 10 * 365 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
// Error cases
|
||||||
|
{
|
||||||
|
name: "empty string",
|
||||||
|
input: "",
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "invalid format",
|
||||||
|
input: "abc",
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "unknown unit",
|
||||||
|
input: "5x",
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "invalid number",
|
||||||
|
input: "xyzd",
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "negative not supported",
|
||||||
|
input: "-5d",
|
||||||
|
wantErr: true,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDurationSpecialCases(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
// Test that standard Go durations work exactly as expected
|
||||||
|
standardDurations := []string{
|
||||||
|
"300ms",
|
||||||
|
"1.5h",
|
||||||
|
"2h45m",
|
||||||
|
"72h",
|
||||||
|
"1us",
|
||||||
|
"1µs",
|
||||||
|
"1ns",
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, d := range standardDurations {
|
||||||
|
expected, err := time.ParseDuration(d)
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
got, err := parseDuration(d)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, expected, got, "standard duration %q should parse identically", d)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestParseDurationRealWorldExamples(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
// Test real-world snapshot purge scenarios
|
||||||
|
tests := []struct {
|
||||||
|
description string
|
||||||
|
input string
|
||||||
|
olderThan time.Duration
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
description: "keep snapshots from last 30 days",
|
||||||
|
input: "30d",
|
||||||
|
olderThan: 30 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
description: "keep snapshots from last 6 months",
|
||||||
|
input: "6mo",
|
||||||
|
olderThan: 6 * 30 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
description: "keep snapshots from last year",
|
||||||
|
input: "1y",
|
||||||
|
olderThan: 365 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
description: "keep snapshots from last week and a half",
|
||||||
|
input: "1w3d",
|
||||||
|
olderThan: 10 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
description: "keep snapshots from last 90 days",
|
||||||
|
input: "90d",
|
||||||
|
olderThan: 90 * 24 * time.Hour,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.description, func(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
got, err := parseDuration(tt.input)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, tt.olderThan, got)
|
||||||
|
|
||||||
|
// Verify the duration makes sense for snapshot purging
|
||||||
|
assert.Greater(t, got, time.Hour,
|
||||||
|
"snapshot purge duration should be at least an hour")
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -141,8 +141,7 @@ specifying a path using --config or by setting VAULTIK_CONFIG to a path.`,
|
|||||||
"orphaned blobs")
|
"orphaned blobs")
|
||||||
cmd.Flags().StringVar(&opts.KeepNewerThan, "keep-newer-than", "",
|
cmd.Flags().StringVar(&opts.KeepNewerThan, "keep-newer-than", "",
|
||||||
"With --prune: keep snapshots newer than this duration "+
|
"With --prune: keep snapshots newer than this duration "+
|
||||||
"(e.g. 30d, 4w, 6mo, 1y; m is minutes, mo is months) "+
|
"(e.g. 4w, 30d, 6mo) instead of only the latest")
|
||||||
"instead of only the latest")
|
|
||||||
|
|
||||||
return cmd
|
return cmd
|
||||||
}
|
}
|
||||||
@@ -205,8 +204,7 @@ restrict the operation to specific snapshot names.`,
|
|||||||
cmd.Flags().BoolVar(&opts.KeepLatest, "keep-latest", false,
|
cmd.Flags().BoolVar(&opts.KeepLatest, "keep-latest", false,
|
||||||
"Keep only the latest snapshot of each name")
|
"Keep only the latest snapshot of each name")
|
||||||
cmd.Flags().StringVar(&opts.OlderThan, "older-than", "",
|
cmd.Flags().StringVar(&opts.OlderThan, "older-than", "",
|
||||||
"Remove snapshots older than duration "+
|
"Remove snapshots older than duration (e.g., 30d, 6m, 1y)")
|
||||||
"(e.g. 30d, 4w, 6mo, 1y; m is minutes, mo is months)")
|
|
||||||
cmd.Flags().BoolVar(&opts.Force, "force", false, "Skip confirmation prompt")
|
cmd.Flags().BoolVar(&opts.Force, "force", false, "Skip confirmation prompt")
|
||||||
cmd.Flags().StringArrayVar(&opts.Names, "snapshot", nil,
|
cmd.Flags().StringArrayVar(&opts.Names, "snapshot", nil,
|
||||||
"Restrict to snapshots with these names (repeat for multiple)")
|
"Restrict to snapshots with these names (repeat for multiple)")
|
||||||
@@ -221,11 +219,8 @@ func newSnapshotVerifyCommand() *cobra.Command {
|
|||||||
cmd := &cobra.Command{
|
cmd := &cobra.Command{
|
||||||
Use: "verify <snapshot-id>",
|
Use: "verify <snapshot-id>",
|
||||||
Short: "Verify snapshot integrity",
|
Short: "Verify snapshot integrity",
|
||||||
Long: "Verifies that all blobs referenced in a snapshot exist.\n\n" +
|
Long: "Verifies that all blobs referenced in a snapshot exist",
|
||||||
"The snapshot may be named by its ID or, on a host with no local\n" +
|
Args: requireSnapshotIDArg,
|
||||||
"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 {
|
RunE: func(cmd *cobra.Command, args []string) error {
|
||||||
snapshotID := args[0]
|
snapshotID := args[0]
|
||||||
|
|
||||||
|
|||||||
@@ -48,10 +48,6 @@ target directory.
|
|||||||
If no paths are specified, all files are restored.
|
If no paths are specified, all files are restored.
|
||||||
If paths are specified, only matching files/directories 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
|
Requires the VAULTIK_AGE_SECRET_KEY environment variable to be set with
|
||||||
the age private key.
|
the age private key.
|
||||||
|
|
||||||
|
|||||||
+5
-13
@@ -219,7 +219,11 @@ func (c *Client) HeadObject(ctx context.Context, key string) (bool, error) {
|
|||||||
Key: aws.String(fullKey),
|
Key: aws.String(fullKey),
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if IsNotFound(err) {
|
var (
|
||||||
|
notFound *s3types.NotFound
|
||||||
|
noSuchKey *s3types.NoSuchKey
|
||||||
|
)
|
||||||
|
if errors.As(err, ¬Found) || errors.As(err, &noSuchKey) {
|
||||||
return false, nil
|
return false, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -229,18 +233,6 @@ func (c *Client) HeadObject(ctx context.Context, key string) (bool, error) {
|
|||||||
return true, nil
|
return true, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsNotFound reports whether err indicates that an object does not exist.
|
|
||||||
// Head and Get requests surface a missing object as different SDK types,
|
|
||||||
// so both are checked here.
|
|
||||||
func IsNotFound(err error) bool {
|
|
||||||
var (
|
|
||||||
notFound *s3types.NotFound
|
|
||||||
noSuchKey *s3types.NoSuchKey
|
|
||||||
)
|
|
||||||
|
|
||||||
return errors.As(err, ¬Found) || errors.As(err, &noSuchKey)
|
|
||||||
}
|
|
||||||
|
|
||||||
// ObjectInfo contains information about an S3 object.
|
// ObjectInfo contains information about an S3 object.
|
||||||
// It is used by ListObjectsStream to return object metadata
|
// It is used by ListObjectsStream to return object metadata
|
||||||
// along with any errors encountered during listing.
|
// along with any errors encountered during listing.
|
||||||
|
|||||||
+54
-79
@@ -46,18 +46,31 @@ func (f *FileStorer) SetFilesystem(fs afero.Fs) {
|
|||||||
// storage base path.
|
// storage base path.
|
||||||
const storageDirPerm = 0o755
|
const storageDirPerm = 0o755
|
||||||
|
|
||||||
// tempSuffix marks a partially written object. writeAtomic streams into a
|
|
||||||
// temp file carrying this suffix and only renames it onto the real key once
|
|
||||||
// the whole object is on disk, so an interrupted write can never leave a
|
|
||||||
// truncated object at the key a later run would Stat and trust as a complete
|
|
||||||
// blob. List and ListStream skip these files, so a leftover from an
|
|
||||||
// interrupted write is never listed or trusted as a blob; it is otherwise
|
|
||||||
// harmless and is overwritten when the same key is written again.
|
|
||||||
const tempSuffix = ".partial"
|
|
||||||
|
|
||||||
// Put stores data at the specified key.
|
// Put stores data at the specified key.
|
||||||
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
|
func (f *FileStorer) Put(_ context.Context, key string, data io.Reader) error {
|
||||||
return f.writeAtomic(key, data, nil)
|
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
|
||||||
}
|
}
|
||||||
|
|
||||||
// PutWithProgress stores data with progress reporting.
|
// PutWithProgress stores data with progress reporting.
|
||||||
@@ -65,7 +78,35 @@ func (f *FileStorer) PutWithProgress(
|
|||||||
_ context.Context, key string, data io.Reader,
|
_ context.Context, key string, data io.Reader,
|
||||||
_ int64, progress ProgressCallback,
|
_ int64, progress ProgressCallback,
|
||||||
) error {
|
) error {
|
||||||
return f.writeAtomic(key, data, progress)
|
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
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get retrieves data from the specified key.
|
// Get retrieves data from the specified key.
|
||||||
@@ -147,7 +188,7 @@ func (f *FileStorer) List(ctx context.Context, prefix string) ([]string, error)
|
|||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
|
|
||||||
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
if !info.IsDir() {
|
||||||
// Convert back to key (relative path from basePath)
|
// Convert back to key (relative path from basePath)
|
||||||
relPath, err := filepath.Rel(f.basePath, path)
|
relPath, err := filepath.Rel(f.basePath, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -204,7 +245,7 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec
|
|||||||
return nil //nolint:nilerr // continue walking despite errors
|
return nil //nolint:nilerr // continue walking despite errors
|
||||||
}
|
}
|
||||||
|
|
||||||
if !info.IsDir() && !strings.HasSuffix(info.Name(), tempSuffix) {
|
if !info.IsDir() {
|
||||||
relPath, err := filepath.Rel(f.basePath, path)
|
relPath, err := filepath.Rel(f.basePath, path)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
|
ch <- ObjectInfo{Err: fmt.Errorf("computing relative path: %w", err)}
|
||||||
@@ -234,72 +275,6 @@ func (f *FileStorer) Info() Info {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// writeAtomic streams data into a temp file in the destination directory,
|
|
||||||
// fsyncs it, and renames it onto the final key. The key therefore appears
|
|
||||||
// only once the whole object has been durably written; a failure part-way
|
|
||||||
// leaves a temp file (removed here on the failing path) rather than a
|
|
||||||
// truncated object at the key.
|
|
||||||
func (f *FileStorer) writeAtomic(
|
|
||||||
key string, data io.Reader, progress ProgressCallback,
|
|
||||||
) error {
|
|
||||||
path := f.fullPath(key)
|
|
||||||
dir := filepath.Dir(path)
|
|
||||||
|
|
||||||
err := f.fs.MkdirAll(dir, storageDirPerm)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating directories: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
tmp, err := afero.TempFile(f.fs, dir, filepath.Base(path)+"-*"+tempSuffix)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("creating temp file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
tmpPath := tmp.Name()
|
|
||||||
|
|
||||||
// Remove the temp file unless the rename below claims it. On the success
|
|
||||||
// path renamed is true, so the deferred Close and Remove are harmless
|
|
||||||
// no-ops on a name that no longer exists.
|
|
||||||
renamed := false
|
|
||||||
|
|
||||||
defer func() {
|
|
||||||
_ = tmp.Close()
|
|
||||||
|
|
||||||
if !renamed {
|
|
||||||
_ = f.fs.Remove(tmpPath)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
var w io.Writer = tmp
|
|
||||||
if progress != nil {
|
|
||||||
w = &progressWriter{writer: tmp, callback: progress}
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = io.Copy(w, data)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("writing file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = tmp.Sync()
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("syncing temp file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = tmp.Close()
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("closing temp file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = f.fs.Rename(tmpPath, path)
|
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("renaming temp file: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
renamed = true
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// fullPath returns the full filesystem path for a key.
|
// fullPath returns the full filesystem path for a key.
|
||||||
func (f *FileStorer) fullPath(key string) string {
|
func (f *FileStorer) fullPath(key string) string {
|
||||||
return filepath.Join(f.basePath, key)
|
return filepath.Join(f.basePath, key)
|
||||||
|
|||||||
@@ -1,119 +0,0 @@
|
|||||||
package storage_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"strings"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
"sneak.berlin/go/vaultik/internal/storage"
|
|
||||||
)
|
|
||||||
|
|
||||||
// errStreamInterrupted stands in for an upload cut off mid-stream.
|
|
||||||
var errStreamInterrupted = errors.New("connection reset mid-upload")
|
|
||||||
|
|
||||||
// failingReader yields its data once, then fails.
|
|
||||||
type failingReader struct {
|
|
||||||
data []byte
|
|
||||||
done bool
|
|
||||||
}
|
|
||||||
|
|
||||||
func (r *failingReader) Read(p []byte) (int, error) {
|
|
||||||
if r.done {
|
|
||||||
return 0, errStreamInterrupted
|
|
||||||
}
|
|
||||||
|
|
||||||
n := copy(p, r.data)
|
|
||||||
r.done = true
|
|
||||||
|
|
||||||
return n, nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFileStorer_InterruptedWriteLeavesNoTrustedObject checks that a write
|
|
||||||
// cut off mid-stream leaves nothing at the destination key, so a later run
|
|
||||||
// cannot Stat a truncated object and trust it as a complete blob.
|
|
||||||
func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
f, err := storage.NewFileStorer(t.TempDir())
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("NewFileStorer: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
key := "blobs/aa/bb/aabbccddeeff"
|
|
||||||
|
|
||||||
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
|
|
||||||
if err == nil {
|
|
||||||
t.Fatal("expected the interrupted write to fail, got nil")
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = f.Stat(ctx, key)
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Fatalf("expected key absent after interrupted write, got Stat err %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
keys, err := f.List(ctx, "blobs/")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("List: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(keys) != 0 {
|
|
||||||
t.Fatalf("expected no keys listed after interrupted write, got %v", keys)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// TestFileStorer_ListSkipsPartialFiles checks that a leftover temp file (the
|
|
||||||
// storage layer names them with a ".partial" suffix) is never surfaced as a
|
|
||||||
// key by List or ListStream.
|
|
||||||
func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
base := t.TempDir()
|
|
||||||
|
|
||||||
f, err := storage.NewFileStorer(base)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("NewFileStorer: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
realKey := "blobs/aa/bb/aabbccddeeff"
|
|
||||||
|
|
||||||
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Put: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
// A stray temp file, as an interrupted write would leave behind.
|
|
||||||
leftover := filepath.Join(base, "blobs/aa/bb/aabbccddeeff-123456.partial")
|
|
||||||
|
|
||||||
err = os.WriteFile(leftover, []byte("half"), 0o600)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("writing leftover temp file: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
keys, err := f.List(ctx, "blobs/")
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("List: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(keys) != 1 || keys[0] != realKey {
|
|
||||||
t.Fatalf("List should return only the real key, got %v", keys)
|
|
||||||
}
|
|
||||||
|
|
||||||
var streamed []string
|
|
||||||
|
|
||||||
for obj := range f.ListStream(ctx, "blobs/") {
|
|
||||||
if obj.Err != nil {
|
|
||||||
t.Fatalf("ListStream: %v", obj.Err)
|
|
||||||
}
|
|
||||||
|
|
||||||
streamed = append(streamed, obj.Key)
|
|
||||||
}
|
|
||||||
|
|
||||||
if len(streamed) != 1 || streamed[0] != realKey {
|
|
||||||
t.Fatalf("ListStream should return only the real key, got %v", streamed)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+1
-16
@@ -38,29 +38,14 @@ func (s *S3Storer) PutWithProgress(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Get retrieves data from the specified key.
|
// Get retrieves data from the specified key.
|
||||||
// Returns ErrNotFound if the object does not exist.
|
|
||||||
func (s *S3Storer) Get(ctx context.Context, key string) (io.ReadCloser, error) {
|
func (s *S3Storer) Get(ctx context.Context, key string) (io.ReadCloser, error) {
|
||||||
rc, err := s.client.GetObject(ctx, key)
|
return s.client.GetObject(ctx, key)
|
||||||
if err != nil {
|
|
||||||
if s3.IsNotFound(err) {
|
|
||||||
return nil, fmt.Errorf("get %q: %w", key, ErrNotFound)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
return rc, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Stat returns metadata about an object without retrieving its contents.
|
// Stat returns metadata about an object without retrieving its contents.
|
||||||
// Returns ErrNotFound if the object does not exist.
|
|
||||||
func (s *S3Storer) Stat(ctx context.Context, key string) (*ObjectInfo, error) {
|
func (s *S3Storer) Stat(ctx context.Context, key string) (*ObjectInfo, error) {
|
||||||
info, err := s.client.StatObject(ctx, key)
|
info, err := s.client.StatObject(ctx, key)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if s3.IsNotFound(err) {
|
|
||||||
return nil, fmt.Errorf("stat %q: %w", key, ErrNotFound)
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,59 +0,0 @@
|
|||||||
package storage_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"errors"
|
|
||||||
"net/http/httptest"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
"github.com/johannesboyne/gofakes3"
|
|
||||||
"github.com/johannesboyne/gofakes3/backend/s3mem"
|
|
||||||
|
|
||||||
"sneak.berlin/go/vaultik/internal/s3"
|
|
||||||
"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.
|
|
||||||
//
|
|
||||||
//nolint:paralleltest // shares an in-process S3 server via t.Cleanup
|
|
||||||
func TestS3StorerMissingKeyMapsToErrNotFound(t *testing.T) {
|
|
||||||
const bucket = "test-bucket"
|
|
||||||
|
|
||||||
backend := s3mem.New()
|
|
||||||
|
|
||||||
err := backend.CreateBucket(bucket)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("create bucket: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
srv := httptest.NewServer(gofakes3.New(backend).Server())
|
|
||||||
t.Cleanup(srv.Close)
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
client, err := s3.NewClient(ctx, s3.Config{
|
|
||||||
Endpoint: srv.URL,
|
|
||||||
Bucket: bucket,
|
|
||||||
AccessKeyID: "test",
|
|
||||||
SecretAccessKey: "test",
|
|
||||||
Region: "us-east-1",
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("new client: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
storer := storage.NewS3Storer(client)
|
|
||||||
|
|
||||||
_, err = storer.Get(ctx, "does-not-exist")
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Errorf("Get on missing key: got %v, want ErrNotFound", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
_, err = storer.Stat(ctx, "does-not-exist")
|
|
||||||
if !errors.Is(err, storage.ErrNotFound) {
|
|
||||||
t.Errorf("Stat on missing key: got %v, want ErrNotFound", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,108 +0,0 @@
|
|||||||
package vaultik_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"context"
|
|
||||||
"io"
|
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"testing"
|
|
||||||
|
|
||||||
"github.com/spf13/afero"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
"sneak.berlin/go/vaultik/internal/log"
|
|
||||||
"sneak.berlin/go/vaultik/internal/ui"
|
|
||||||
"sneak.berlin/go/vaultik/internal/vaultik"
|
|
||||||
)
|
|
||||||
|
|
||||||
// TestDeepVerifyAcceptsHealthyAndRejectsCorruptBlob backs up a real
|
|
||||||
// snapshot with the on-disk storage backend, runs deep verification on
|
|
||||||
// it, then flips a byte inside one stored blob and runs deep
|
|
||||||
// verification again. A healthy snapshot must pass; a corrupted blob
|
|
||||||
// must fail. The healthy case is the regression guard: deep
|
|
||||||
// verification used to hash the encrypted blob bytes and compare them
|
|
||||||
// to the blob's ID (the double SHA256 of the plaintext), so it reported
|
|
||||||
// every healthy blob as corrupt.
|
|
||||||
func TestDeepVerifyAcceptsHealthyAndRejectsCorruptBlob(t *testing.T) {
|
|
||||||
log.Initialize(log.Config{})
|
|
||||||
t.Parallel()
|
|
||||||
|
|
||||||
fs := afero.NewOsFs()
|
|
||||||
tempDir := t.TempDir()
|
|
||||||
|
|
||||||
dataDir := filepath.Join(tempDir, "source")
|
|
||||||
storeDir := filepath.Join(tempDir, "remote")
|
|
||||||
dbPath := filepath.Join(tempDir, "index.sqlite")
|
|
||||||
|
|
||||||
chunkSize := int64(64 * 1024)
|
|
||||||
maxBlobSize := int64(512 * 1024)
|
|
||||||
|
|
||||||
// One file large enough to span several chunks within a single blob.
|
|
||||||
require.NoError(t, fs.MkdirAll(dataDir, 0o755))
|
|
||||||
require.NoError(t, afero.WriteFile(fs,
|
|
||||||
filepath.Join(dataDir, "data.bin"),
|
|
||||||
bytesPattern("deep-", int(chunkSize*3)), 0o644))
|
|
||||||
|
|
||||||
ctx := context.Background()
|
|
||||||
|
|
||||||
// runFileStorageBackup writes a real snapshot to storeDir and closes
|
|
||||||
// the source index, so verification runs from remote bytes only.
|
|
||||||
cfg, storer, snapshotID := runFileStorageBackup(
|
|
||||||
ctx, t, fs, dataDir, storeDir, dbPath, chunkSize, maxBlobSize)
|
|
||||||
|
|
||||||
newVerifier := func() *vaultik.Vaultik {
|
|
||||||
v := &vaultik.Vaultik{
|
|
||||||
Config: cfg,
|
|
||||||
Storage: storer,
|
|
||||||
Fs: fs,
|
|
||||||
Stdout: io.Discard,
|
|
||||||
Stderr: io.Discard,
|
|
||||||
UI: ui.NewWithColor(io.Discard, false),
|
|
||||||
}
|
|
||||||
v.SetContext(ctx)
|
|
||||||
|
|
||||||
return v
|
|
||||||
}
|
|
||||||
|
|
||||||
require.NoError(t,
|
|
||||||
newVerifier().RunDeepVerify(snapshotID, &vaultik.VerifyOptions{Deep: true}),
|
|
||||||
"deep verify should pass on a healthy snapshot")
|
|
||||||
|
|
||||||
// Flip a byte inside one blob without changing its length, so the
|
|
||||||
// blob-existence and size checks still pass and verification reaches
|
|
||||||
// the blob-content stage.
|
|
||||||
corruptOneBlob(t, fs, filepath.Join(storeDir, "blobs"))
|
|
||||||
|
|
||||||
require.Error(t,
|
|
||||||
newVerifier().RunDeepVerify(snapshotID, &vaultik.VerifyOptions{Deep: true}),
|
|
||||||
"deep verify should fail on a corrupted blob")
|
|
||||||
}
|
|
||||||
|
|
||||||
// corruptOneBlob flips a middle byte of the first blob file found under
|
|
||||||
// blobsDir, leaving the file length unchanged.
|
|
||||||
func corruptOneBlob(t *testing.T, fs afero.Fs, blobsDir string) {
|
|
||||||
t.Helper()
|
|
||||||
|
|
||||||
var blobPath string
|
|
||||||
|
|
||||||
err := afero.Walk(fs, blobsDir,
|
|
||||||
func(path string, info os.FileInfo, err error) error {
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
|
|
||||||
if blobPath == "" && !info.IsDir() {
|
|
||||||
blobPath = path
|
|
||||||
}
|
|
||||||
|
|
||||||
return nil
|
|
||||||
})
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotEmpty(t, blobPath, "expected at least one blob on disk")
|
|
||||||
|
|
||||||
data, err := afero.ReadFile(fs, blobPath)
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NotEmpty(t, data)
|
|
||||||
|
|
||||||
data[len(data)/2] ^= 0xff
|
|
||||||
require.NoError(t, afero.WriteFile(fs, blobPath, data, 0o644))
|
|
||||||
}
|
|
||||||
@@ -33,9 +33,8 @@ func ubytes(n int64) string {
|
|||||||
var (
|
var (
|
||||||
errMalformedSnapshotID = errors.New(
|
errMalformedSnapshotID = errors.New(
|
||||||
"invalid snapshot ID format: expected hostname_snapshotname_timestamp")
|
"invalid snapshot ID format: expected hostname_snapshotname_timestamp")
|
||||||
errInvalidDuration = errors.New("invalid duration")
|
errInvalidDuration = errors.New("invalid duration")
|
||||||
errUnknownTimeUnit = errors.New("unknown time unit")
|
errUnknownTimeUnit = errors.New("unknown time unit")
|
||||||
errNegativeDuration = errors.New("negative durations are not supported")
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Time-unit lengths used by parseDuration.
|
// Time-unit lengths used by parseDuration.
|
||||||
@@ -139,13 +138,8 @@ func parseSnapshotName(snapshotID string) string {
|
|||||||
|
|
||||||
// parseDuration parses a duration string with support for human-friendly units:
|
// parseDuration parses a duration string with support for human-friendly units:
|
||||||
// d/day/days, w/week/weeks, mo/month/months, y/year/years, plus standard Go
|
// d/day/days, w/week/weeks, mo/month/months, y/year/years, plus standard Go
|
||||||
// duration units. Following Go, m is minutes and mo is months. A bare number,
|
// duration units (h, m, s).
|
||||||
// an unknown unit, and a negative value are all rejected.
|
|
||||||
func parseDuration(s string) (time.Duration, error) {
|
func parseDuration(s string) (time.Duration, error) {
|
||||||
if strings.HasPrefix(strings.TrimSpace(s), "-") {
|
|
||||||
return 0, errNegativeDuration
|
|
||||||
}
|
|
||||||
|
|
||||||
d, err := time.ParseDuration(s)
|
d, err := time.ParseDuration(s)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
return d, nil
|
return d, nil
|
||||||
|
|||||||
@@ -51,32 +51,13 @@ func TestParseDuration(t *testing.T) {
|
|||||||
want time.Duration
|
want time.Duration
|
||||||
err bool
|
err bool
|
||||||
}{
|
}{
|
||||||
// Go units, including the m-is-minutes / mo-is-months distinction
|
|
||||||
// that this parser exists to keep straight.
|
|
||||||
{"10ns", 10 * time.Nanosecond, false},
|
|
||||||
{"10us", 10 * time.Microsecond, false},
|
|
||||||
{"500ms", 500 * time.Millisecond, false},
|
|
||||||
{"30s", 30 * time.Second, false},
|
|
||||||
{"6m", 6 * time.Minute, false},
|
|
||||||
{"1h", time.Hour, false},
|
|
||||||
// Extended calendar units.
|
|
||||||
{"30d", 30 * 24 * time.Hour, false},
|
{"30d", 30 * 24 * time.Hour, false},
|
||||||
{"3days", 3 * 24 * time.Hour, false},
|
|
||||||
{"4w", 4 * 7 * 24 * time.Hour, false},
|
{"4w", 4 * 7 * 24 * time.Hour, false},
|
||||||
{"2weeks", 2 * 7 * 24 * time.Hour, false},
|
{"6mo", 6 * 30 * 24 * time.Hour, false},
|
||||||
{"6mo", 180 * 24 * time.Hour, false},
|
|
||||||
{"1month", 30 * 24 * time.Hour, false},
|
|
||||||
{"1y", 365 * 24 * time.Hour, false},
|
{"1y", 365 * 24 * time.Hour, false},
|
||||||
{"2years", 2 * 365 * 24 * time.Hour, false},
|
|
||||||
// Combined units.
|
|
||||||
{"2w3d", 2*7*24*time.Hour + 3*24*time.Hour, false},
|
{"2w3d", 2*7*24*time.Hour + 3*24*time.Hour, false},
|
||||||
{"1y6mo", 365*24*time.Hour + 180*24*time.Hour, false},
|
{"1h", time.Hour, false},
|
||||||
// Rejected inputs.
|
{"30s", 30 * time.Second, false},
|
||||||
{"6", 0, true}, // bare number, no unit
|
|
||||||
{"5x", 0, true}, // unknown unit
|
|
||||||
{"-5d", 0, true}, // negative, extended unit
|
|
||||||
{"-5h", 0, true}, // negative, Go unit
|
|
||||||
{"", 0, true}, // empty
|
|
||||||
{"garbage", 0, true},
|
{"garbage", 0, true},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,79 +0,0 @@
|
|||||||
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,6 +18,7 @@ import (
|
|||||||
"sneak.berlin/go/vaultik/internal/blobgen"
|
"sneak.berlin/go/vaultik/internal/blobgen"
|
||||||
"sneak.berlin/go/vaultik/internal/database"
|
"sneak.berlin/go/vaultik/internal/database"
|
||||||
"sneak.berlin/go/vaultik/internal/log"
|
"sneak.berlin/go/vaultik/internal/log"
|
||||||
|
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||||
"sneak.berlin/go/vaultik/internal/types"
|
"sneak.berlin/go/vaultik/internal/types"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -576,20 +577,14 @@ func (v *Vaultik) handleRestoreVerification(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// downloadSnapshotDB downloads and decrypts the snapshot metadata
|
// downloadSnapshotDB downloads and decrypts the snapshot metadata
|
||||||
// database. The identifier is resolved to the snapshot's remote key: a
|
// database. The snapshotID is the human ID; we hash it to the remote
|
||||||
// human ID is hashed, and a remote key (or its abbreviation, as printed
|
// key for the storage path.
|
||||||
// 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(
|
func (v *Vaultik) downloadSnapshotDB(
|
||||||
snapshotID string, identity age.Identity,
|
snapshotID string, identity age.Identity,
|
||||||
) (*database.DB, error) {
|
) (*database.DB, error) {
|
||||||
remoteKey, err := v.resolveSnapshotRemoteKey(snapshotID)
|
|
||||||
if err != nil {
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
|
|
||||||
// Download encrypted database from storage
|
// Download encrypted database from storage
|
||||||
dbKey := fmt.Sprintf("metadata/%s/db.zst.age", remoteKey)
|
dbKey := fmt.Sprintf("metadata/%s/db.zst.age",
|
||||||
|
snapshot.RemoteSnapshotKey(snapshotID))
|
||||||
|
|
||||||
reader, err := v.Storage.Get(v.ctx, dbKey)
|
reader, err := v.Storage.Get(v.ctx, dbKey)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -1,167 +0,0 @@
|
|||||||
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,7 +8,6 @@ import (
|
|||||||
"path/filepath"
|
"path/filepath"
|
||||||
"regexp"
|
"regexp"
|
||||||
"sort"
|
"sort"
|
||||||
"strconv"
|
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -670,11 +669,9 @@ func (v *Vaultik) VerifySnapshotWithOptions(
|
|||||||
|
|
||||||
v.printVerifyHeader(snapshotID, opts)
|
v.printVerifyHeader(snapshotID, opts)
|
||||||
|
|
||||||
// Resolve the identifier to the snapshot's remote key and download the
|
// Download and parse manifest. The caller supplies a human
|
||||||
// manifest. A human ID is hashed; a remote key (or its abbreviation,
|
// snapshot ID; we hash it to address remote storage.
|
||||||
// as printed for a remote-only snapshot) is used as-is, so a host with
|
manifest, err := v.downloadManifestByKey(snapshot.RemoteSnapshotKey(snapshotID))
|
||||||
// no local index can verify a snapshot it can only see on the store.
|
|
||||||
manifest, err := v.resolveAndDownloadManifest(snapshotID)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if opts.JSON {
|
if opts.JSON {
|
||||||
result.Status = verifyStatusFailed
|
result.Status = verifyStatusFailed
|
||||||
@@ -1543,17 +1540,12 @@ func (v *Vaultik) outputRemoveJSON(result *RemoveResult) error {
|
|||||||
return encoder.Encode(result)
|
return encoder.Encode(result)
|
||||||
}
|
}
|
||||||
|
|
||||||
// PruneResult contains statistics about the prune operation.
|
// PruneResult contains statistics about the prune operation
|
||||||
// SnapshotsDeleted counts snapshots actually deleted. FilesDeleted,
|
|
||||||
// ChunksDeleted, and BlobsDeleted are derived from before/after row
|
|
||||||
// counts of the local index; each is nil when a count could not be read,
|
|
||||||
// so an unreadable count is reported as unknown rather than silently
|
|
||||||
// as 0.
|
|
||||||
type PruneResult struct {
|
type PruneResult struct {
|
||||||
SnapshotsDeleted int64
|
SnapshotsDeleted int64
|
||||||
FilesDeleted *int64
|
FilesDeleted int64
|
||||||
ChunksDeleted *int64
|
ChunksDeleted int64
|
||||||
BlobsDeleted *int64
|
BlobsDeleted int64
|
||||||
}
|
}
|
||||||
|
|
||||||
// PruneDatabase removes incomplete snapshots and orphaned files, chunks,
|
// PruneDatabase removes incomplete snapshots and orphaned files, chunks,
|
||||||
@@ -1568,7 +1560,7 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
result := &PruneResult{}
|
result := &PruneResult{}
|
||||||
|
|
||||||
// Snapshot counts before deletion of incompletes.
|
// Snapshot counts before deletion of incompletes.
|
||||||
snapshotCountBefore := v.tableCountForReport("snapshots")
|
snapshotCountBefore, _ := v.getTableCount("snapshots")
|
||||||
|
|
||||||
// First, delete any incomplete snapshots
|
// First, delete any incomplete snapshots
|
||||||
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
|
incompleteSnapshots, err := v.Repositories.Snapshots.GetIncompleteSnapshots(v.ctx)
|
||||||
@@ -1583,9 +1575,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Get counts before cleanup for reporting
|
// Get counts before cleanup for reporting
|
||||||
fileCountBefore := v.tableCountForReport("files")
|
fileCountBefore, _ := v.getTableCount("files")
|
||||||
chunkCountBefore := v.tableCountForReport("chunks")
|
chunkCountBefore, _ := v.getTableCount("chunks")
|
||||||
blobCountBefore := v.tableCountForReport("blobs")
|
blobCountBefore, _ := v.getTableCount("blobs")
|
||||||
|
|
||||||
// Run the cleanup
|
// Run the cleanup
|
||||||
err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
|
err = v.SnapshotManager.CleanupOrphanedData(v.ctx)
|
||||||
@@ -1594,83 +1586,36 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Get counts after cleanup
|
// Get counts after cleanup
|
||||||
fileCountAfter := v.tableCountForReport("files")
|
fileCountAfter, _ := v.getTableCount("files")
|
||||||
chunkCountAfter := v.tableCountForReport("chunks")
|
chunkCountAfter, _ := v.getTableCount("chunks")
|
||||||
blobCountAfter := v.tableCountForReport("blobs")
|
blobCountAfter, _ := v.getTableCount("blobs")
|
||||||
|
|
||||||
result.FilesDeleted = countDiff(fileCountBefore, fileCountAfter)
|
result.FilesDeleted = fileCountBefore - fileCountAfter
|
||||||
result.ChunksDeleted = countDiff(chunkCountBefore, chunkCountAfter)
|
result.ChunksDeleted = chunkCountBefore - chunkCountAfter
|
||||||
result.BlobsDeleted = countDiff(blobCountBefore, blobCountAfter)
|
result.BlobsDeleted = blobCountBefore - blobCountAfter
|
||||||
|
|
||||||
log.Info("Local database prune complete",
|
log.Info("Local database prune complete",
|
||||||
"incomplete_snapshots", result.SnapshotsDeleted,
|
"incomplete_snapshots", result.SnapshotsDeleted,
|
||||||
"orphaned_files", countText(result.FilesDeleted),
|
"orphaned_files", result.FilesDeleted,
|
||||||
"orphaned_chunks", countText(result.ChunksDeleted),
|
"orphaned_chunks", result.ChunksDeleted,
|
||||||
"orphaned_blobs", countText(result.BlobsDeleted),
|
"orphaned_blobs", result.BlobsDeleted,
|
||||||
)
|
)
|
||||||
|
|
||||||
// Snapshots remaining after removing the incomplete ones; unknown if
|
snapshotCountAfter := snapshotCountBefore - result.SnapshotsDeleted
|
||||||
// the pre-prune snapshot count could not be read.
|
|
||||||
snapshotsRemain := countDiff(snapshotCountBefore, &result.SnapshotsDeleted)
|
|
||||||
|
|
||||||
v.UI.Completef("Pruned local index database.")
|
v.UI.Completef("Pruned local index database.")
|
||||||
v.UI.Detailf("Incomplete snapshots: %s removed (%s remain).",
|
v.UI.Detailf("Incomplete snapshots: %d removed (%d remain).",
|
||||||
countText(&result.SnapshotsDeleted), countText(snapshotsRemain))
|
result.SnapshotsDeleted, snapshotCountAfter)
|
||||||
v.UI.Detailf("Orphaned files: %s removed (%s remain).",
|
v.UI.Detailf("Orphaned files: %d removed (%d remain).",
|
||||||
countText(result.FilesDeleted), countText(fileCountAfter))
|
result.FilesDeleted, fileCountAfter)
|
||||||
v.UI.Detailf("Orphaned chunks: %s removed (%s remain).",
|
v.UI.Detailf("Orphaned chunks: %d removed (%d remain).",
|
||||||
countText(result.ChunksDeleted), countText(chunkCountAfter))
|
result.ChunksDeleted, chunkCountAfter)
|
||||||
v.UI.Detailf("Orphaned blobs: %s removed (%s remain).",
|
v.UI.Detailf("Orphaned blobs: %d removed (%d remain).",
|
||||||
countText(result.BlobsDeleted), countText(blobCountAfter))
|
result.BlobsDeleted, blobCountAfter)
|
||||||
|
|
||||||
return result, nil
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// countUnknown is what a count reads as when its query could not be run,
|
|
||||||
// distinct from "0", which means the table really was empty.
|
|
||||||
const countUnknown = "unknown"
|
|
||||||
|
|
||||||
// tableCountForReport returns the row count of a table for the prune
|
|
||||||
// summary, or nil if the count could not be read. A read failure is
|
|
||||||
// logged at warn — visible even under --json, which routes warnings to
|
|
||||||
// stderr — and then rendered as unknown rather than silently becoming 0,
|
|
||||||
// so a broken query is a visible failure instead of a plausible wrong
|
|
||||||
// number.
|
|
||||||
func (v *Vaultik) tableCountForReport(tableName string) *int64 {
|
|
||||||
count, err := v.getTableCount(tableName)
|
|
||||||
if err != nil {
|
|
||||||
log.Warn("could not read table row count for prune summary",
|
|
||||||
"table", tableName, "error", err)
|
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
return &count
|
|
||||||
}
|
|
||||||
|
|
||||||
// countDiff returns before-after, or nil if either count is unknown so
|
|
||||||
// that an unreadable count does not collapse into a plausible delta.
|
|
||||||
func countDiff(before, after *int64) *int64 {
|
|
||||||
if before == nil || after == nil {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
diff := *before - *after
|
|
||||||
|
|
||||||
return &diff
|
|
||||||
}
|
|
||||||
|
|
||||||
// countText renders a count that may be unknown: nil (the read failed)
|
|
||||||
// becomes "unknown", never "0", so a reader can tell an empty table from
|
|
||||||
// one that could not be queried.
|
|
||||||
func countText(count *int64) string {
|
|
||||||
if count == nil {
|
|
||||||
return countUnknown
|
|
||||||
}
|
|
||||||
|
|
||||||
return strconv.FormatInt(*count, 10)
|
|
||||||
}
|
|
||||||
|
|
||||||
// validTableNameRe matches table names containing only lowercase
|
// validTableNameRe matches table names containing only lowercase
|
||||||
// alphanumeric characters and underscores.
|
// alphanumeric characters and underscores.
|
||||||
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
|
var validTableNameRe = regexp.MustCompile(`^[a-z0-9_]+$`)
|
||||||
|
|||||||
@@ -1,101 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
+23
-37
@@ -138,15 +138,8 @@ func (v *Vaultik) RunDeepVerify(snapshotID string, opts *VerifyOptions) error {
|
|||||||
func (v *Vaultik) loadVerificationData(
|
func (v *Vaultik) loadVerificationData(
|
||||||
snapshotID string, opts *VerifyOptions, result *VerifyResult,
|
snapshotID string, opts *VerifyOptions, result *VerifyResult,
|
||||||
) (*snapshot.Manifest, *tempDB, []snapshot.BlobInfo, error) {
|
) (*snapshot.Manifest, *tempDB, []snapshot.BlobInfo, error) {
|
||||||
// Resolve the identifier to the snapshot's remote key. A human ID is
|
// All remote paths use the hashed key derived from the human ID.
|
||||||
// hashed; a remote key (or its abbreviation, as printed for a
|
remoteKey := snapshot.RemoteSnapshotKey(snapshotID)
|
||||||
// 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
|
// Download manifest. downloadManifestByKey is the single reader for
|
||||||
// remote manifests; see its doc comment.
|
// remote manifests; see its doc comment.
|
||||||
@@ -193,7 +186,7 @@ func (v *Vaultik) loadVerificationData(
|
|||||||
fmt.Errorf("failed to decrypt database: %w", err))
|
fmt.Errorf("failed to decrypt database: %w", err))
|
||||||
}
|
}
|
||||||
|
|
||||||
dbBlobs, err := v.getBlobsFromDatabase(tdb.DB)
|
dbBlobs, err := v.getBlobsFromDatabase(snapshotID, tdb.DB)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = tdb.Close()
|
_ = tdb.Close()
|
||||||
|
|
||||||
@@ -351,8 +344,12 @@ func (v *Vaultik) verifyBlob(blobInfo snapshot.BlobInfo, db *sql.DB) error {
|
|||||||
return fmt.Errorf("failed to get decryptor: %w", err)
|
return fmt.Errorf("failed to get decryptor: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Decrypt blob
|
// Hash the encrypted blob data as it streams through to decryption
|
||||||
decryptedReader, err := decryptor.DecryptStream(reader)
|
blobHasher := sha256.New()
|
||||||
|
teeReader := io.TeeReader(reader, blobHasher)
|
||||||
|
|
||||||
|
// Decrypt blob (reading through teeReader to hash encrypted data)
|
||||||
|
decryptedReader, err := decryptor.DecryptStream(teeReader)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to decrypt: %w", err)
|
return fmt.Errorf("failed to decrypt: %w", err)
|
||||||
}
|
}
|
||||||
@@ -364,19 +361,12 @@ func (v *Vaultik) verifyBlob(blobInfo snapshot.BlobInfo, db *sql.DB) error {
|
|||||||
}
|
}
|
||||||
defer decompressor.Close()
|
defer decompressor.Close()
|
||||||
|
|
||||||
// A blob's hash — its remote name — is the double SHA256 of its
|
chunkCount, err := v.verifyBlobChunks(db, blobInfo.Hash, decompressor)
|
||||||
// decompressed plaintext (see blobgen.Writer.Sum256), not of the
|
|
||||||
// encrypted bytes. Hash the plaintext as chunk verification streams
|
|
||||||
// it, then compare on completion.
|
|
||||||
plaintextHasher := sha256.New()
|
|
||||||
hashedStream := io.TeeReader(decompressor, plaintextHasher)
|
|
||||||
|
|
||||||
chunkCount, err := v.verifyBlobChunks(db, blobInfo.Hash, hashedStream)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
err = v.verifyBlobFinalIntegrity(hashedStream, plaintextHasher, blobInfo.Hash)
|
err = v.verifyBlobFinalIntegrity(decompressor, blobHasher, blobInfo.Hash)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -480,13 +470,14 @@ func (v *Vaultik) verifyBlobChunks(
|
|||||||
}
|
}
|
||||||
|
|
||||||
// verifyBlobFinalIntegrity checks that no trailing data exists in the
|
// verifyBlobFinalIntegrity checks that no trailing data exists in the
|
||||||
// decompressed stream and that the blob hash matches the expected value.
|
// decompressed stream and that the encrypted blob hash matches the
|
||||||
|
// expected value.
|
||||||
func (v *Vaultik) verifyBlobFinalIntegrity(
|
func (v *Vaultik) verifyBlobFinalIntegrity(
|
||||||
plaintext io.Reader, plaintextHasher hash.Hash, expectedHash string,
|
decompressor io.Reader, blobHasher hash.Hash, expectedHash string,
|
||||||
) error {
|
) error {
|
||||||
// Verify no remaining data in blob - if the chunk list is accurate,
|
// Verify no remaining data in blob - if the chunk list is accurate,
|
||||||
// the blob should be fully consumed.
|
// the blob should be fully consumed.
|
||||||
remaining, err := io.Copy(io.Discard, plaintext)
|
remaining, err := io.Copy(io.Discard, decompressor)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to check for remaining blob data: %w", err)
|
return fmt.Errorf("failed to check for remaining blob data: %w", err)
|
||||||
}
|
}
|
||||||
@@ -495,11 +486,8 @@ func (v *Vaultik) verifyBlobFinalIntegrity(
|
|||||||
return fmt.Errorf("%w: %d bytes", errTrailingBlobData, remaining)
|
return fmt.Errorf("%w: %d bytes", errTrailingBlobData, remaining)
|
||||||
}
|
}
|
||||||
|
|
||||||
// The blob hash is the double SHA256 of its plaintext content.
|
// Verify blob hash matches the encrypted data we downloaded
|
||||||
firstHash := plaintextHasher.Sum(nil)
|
calculatedBlobHash := hex.EncodeToString(blobHasher.Sum(nil))
|
||||||
secondHash := sha256.Sum256(firstHash)
|
|
||||||
calculatedBlobHash := hex.EncodeToString(secondHash[:])
|
|
||||||
|
|
||||||
if calculatedBlobHash != expectedHash {
|
if calculatedBlobHash != expectedHash {
|
||||||
return fmt.Errorf("%w: calculated %s, expected %s",
|
return fmt.Errorf("%w: calculated %s, expected %s",
|
||||||
errBlobHashMismatch, calculatedBlobHash, expectedHash)
|
errBlobHashMismatch, calculatedBlobHash, expectedHash)
|
||||||
@@ -508,21 +496,19 @@ func (v *Vaultik) verifyBlobFinalIntegrity(
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// getBlobsFromDatabase gets all blobs for the snapshot from the database.
|
// getBlobsFromDatabase gets all blobs for the snapshot from the database
|
||||||
//
|
func (v *Vaultik) getBlobsFromDatabase(
|
||||||
// The exported per-snapshot database holds exactly one snapshot's data
|
snapshotID string, db *sql.DB,
|
||||||
// (see cleanSnapshotDB), so every row in snapshot_blobs belongs to it.
|
) ([]snapshot.BlobInfo, error) {
|
||||||
// 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 := `
|
query := `
|
||||||
SELECT b.blob_hash, b.compressed_size
|
SELECT b.blob_hash, b.compressed_size
|
||||||
FROM snapshot_blobs sb
|
FROM snapshot_blobs sb
|
||||||
JOIN blobs b ON sb.blob_hash = b.blob_hash
|
JOIN blobs b ON sb.blob_hash = b.blob_hash
|
||||||
|
WHERE sb.snapshot_id = ?
|
||||||
ORDER BY b.blob_hash
|
ORDER BY b.blob_hash
|
||||||
`
|
`
|
||||||
|
|
||||||
rows, err := db.QueryContext(v.ctx, query)
|
rows, err := db.QueryContext(v.ctx, query, snapshotID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("failed to query snapshot blobs: %w", err)
|
return nil, fmt.Errorf("failed to query snapshot blobs: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user