Compare commits
7
Commits
98b981203e
..
next
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e459a66099 | ||
|
|
32c4a46b49 | ||
|
|
b5389a62b5 | ||
|
|
d87202fb70 | ||
|
|
e161343eac | ||
|
|
d53202eb86 | ||
|
|
8b22ae8d42 |
+7
-3
@@ -54,10 +54,12 @@ The database tracks five primary entities and their relationships:
|
||||
|
||||
#### File (`database.File`)
|
||||
Represents a file, directory, or symlink in the backup system. Stores metadata needed for restoration:
|
||||
- Path, source_path (for restore path stripping), mtime
|
||||
- Path, mtime
|
||||
- Size, mode, ownership (uid, gid)
|
||||
- Symlink target (if applicable)
|
||||
|
||||
It also stores `source_path`, the source directory the scan found it under, made absolute and with symlinks resolved. Restore does not read it.
|
||||
|
||||
#### Chunk (`database.Chunk`)
|
||||
A content-addressed unit of data. Files are split into variable-size chunks using the FastCDC algorithm:
|
||||
- `ChunkHash`: SHA256 hash of chunk content (primary key)
|
||||
@@ -82,9 +84,11 @@ The final storage unit uploaded to S3. Contains many compressed and encrypted ch
|
||||
Blob creation process:
|
||||
1. Chunks are accumulated (up to MaxBlobSize, typically 10GB)
|
||||
2. As each chunk is added, its uncompressed bytes are fed to a running SHA-256
|
||||
3. Concurrently, the same bytes are compressed with zstd, then encrypted with age (recipients configured in config), and streamed to storage
|
||||
3. Concurrently, the same bytes are compressed with zstd, then encrypted with age (recipients configured in config), and written to a temporary file
|
||||
4. On finalize, the blob's name is the double SHA-256 of the uncompressed contents — `hex(SHA256(SHA256(...)))` — not a hash of the compressed, encrypted bytes
|
||||
5. Uploaded to `blobs/{hash[0:2]}/{hash[2:4]}/{hash}`
|
||||
5. The finished file is uploaded to `blobs/{hash[0:2]}/{hash[2:4]}/{hash}` and then deleted
|
||||
|
||||
A backup needs free temporary space, because each blob is written whole to a temporary file before it is uploaded (up to about `blob_size_limit`; an rclone destination that cannot stream uploads needs about twice that) and the metadata export writes copies of the local index. Temporary files go to `$TMPDIR` (default `/tmp`); with `TMPDIR` unset, SQLite writes one of those copies to `/var/tmp`.
|
||||
|
||||
#### BlobChunk (`database.BlobChunk`)
|
||||
Maps chunks to their position within blobs:
|
||||
|
||||
@@ -428,6 +428,16 @@ a local or mounted filesystem. Useful for testing or backing up to a NAS.
|
||||
**Rclone** (`rclone://remote/path`): Uses rclone's 70+ supported cloud
|
||||
providers. Requires rclone to be configured separately (`rclone config`).
|
||||
|
||||
An upload cut off part-way leaves nothing under the object's name on S3, which
|
||||
shows an object only once its upload has completed, and on the local filesystem
|
||||
backend, which writes a temporary file and renames it into place. Rclone
|
||||
remotes with a server-side move (local and sftp among them) are written under a
|
||||
temporary name ending in `.partial` and moved into place. Rclone remotes without
|
||||
one are written in place: where such a remote shows a file while it is still
|
||||
being written, a killed upload can leave a truncated object under its name,
|
||||
which a later backup takes for complete. A leftover `.partial` file is ignored
|
||||
and can be deleted.
|
||||
|
||||
Legacy S3 configuration via `s3.*` fields (endpoint, bucket, prefix, etc.) is
|
||||
still supported for backward compatibility. `storage_url` takes precedence if
|
||||
both are set.
|
||||
@@ -546,7 +556,7 @@ complete annotated example also lives in
|
||||
| `s3.*` | | Legacy S3 configuration (endpoint, bucket, credentials) |
|
||||
| `exclude` | | Global exclude patterns (applied to all snapshots) |
|
||||
| `chunk_size` | `10MB` | Average chunk size for content-defined chunking |
|
||||
| `blob_size_limit` | `10GB` | Maximum blob size before splitting. Must be at least four times `chunk_size` (the largest chunk the chunker can emit), otherwise a single-chunk blob could exceed the limit |
|
||||
| `blob_size_limit` | `10GB` | Maximum blob size before splitting. Must be at least four times `chunk_size` (the largest chunk the chunker can emit), otherwise a single-chunk blob could exceed the limit. A backup needs free temporary space, because each blob is written whole to a temporary file before it is uploaded (up to about `blob_size_limit`; an rclone destination that cannot stream uploads needs about twice that) and the metadata export writes copies of the local index. Temporary files go to `$TMPDIR` (default `/tmp`); with `TMPDIR` unset, SQLite writes one of those copies to `/var/tmp` |
|
||||
| `compression_level` | `3` | zstd compression level (1-19) |
|
||||
| `hostname` | system hostname | Hostname used in snapshot IDs |
|
||||
| `index_path` | platform data dir | Local SQLite index path |
|
||||
@@ -805,7 +815,8 @@ them. We provide:
|
||||
in the `Dockerfile` together.
|
||||
* `script/release` — cross-compile and publish the release artifacts
|
||||
with the pinned `goreleaser`. Refuses a `goreleaser` on `PATH` whose
|
||||
version is not the pinned one, on the same reasoning as `script/lint`.
|
||||
version is not the pinned one, because a different version would build
|
||||
a different release from the same tag.
|
||||
* `script/release-snapshot` — the same build with no publishing and no
|
||||
tagging, into `./dist`
|
||||
* `script/test` — run the test suite by building the `test` phase of
|
||||
@@ -915,14 +926,14 @@ It is passed to `goreleaser` as `GITEA_TOKEN`. The runner's automatic
|
||||
token is deliberately not used: it is not guaranteed to carry release
|
||||
write access.
|
||||
|
||||
The Go toolchain that compiles the released binaries comes from an
|
||||
`actions/setup-go` step pinned by commit sha, reading its version from
|
||||
`go.mod` (currently `1.26.1`, the same version the `Dockerfile` builder
|
||||
stage pins by digest). `goreleaser` shells out to `go` for every
|
||||
The Go toolchain that compiles the released binaries is installed by
|
||||
`script/install-go`, which downloads the version named by `go.mod`
|
||||
(currently `1.26.1`, the same version the `Dockerfile` builder stage
|
||||
pins by digest) and refuses the archive unless its sha256 matches the
|
||||
value committed in the script. `goreleaser` shells out to `go` for every
|
||||
cross-compile, so without that step the release would either fail
|
||||
outright or ship binaries built by whatever unpinned toolchain the
|
||||
runner happened to carry — the one unpinned thing in an otherwise
|
||||
hash-pinned release path.
|
||||
runner happened to carry.
|
||||
|
||||
To rehearse the whole build without publishing or tagging anything:
|
||||
|
||||
|
||||
@@ -22,14 +22,79 @@ the tag exists and is exercised; what is left is merging `next` to
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-10-07: Made a symlink whose target cannot be read stop the backup
|
||||
([issue #269](https://git.eeqj.de/sneak/vaultik/issues/269)). It was
|
||||
left out of the snapshot with only a debug log line, even without
|
||||
`--skip-errors`. It now aborts the run, or with `--skip-errors` is
|
||||
skipped with the `Failed to access` error line that any other entry
|
||||
the scan cannot read gets.
|
||||
|
||||
- 2026-10-07: Stopped a killed rclone upload from leaving a truncated
|
||||
object under its key
|
||||
([issue #266](https://git.eeqj.de/sneak/vaultik/issues/266)). The
|
||||
rclone backend wrote each object straight to its key, so killing an
|
||||
upload to a local or sftp remote left a truncated object there; the
|
||||
next backup found the key with `Stat`, skipped the upload and recorded
|
||||
a snapshot that could not be restored. On a remote with a server-side
|
||||
move, an object is now written under a temporary name ending in
|
||||
`.partial` and moved into place, and listings skip such names. Remotes
|
||||
without one are still written in place.
|
||||
|
||||
- 2026-10-07: Made an interrupted command exit 130 and say so
|
||||
([issue #267](https://git.eeqj.de/sneak/vaultik/issues/267)). Ctrl-C
|
||||
or SIGTERM during `snapshot create`, `snapshot restore` or `snapshot
|
||||
verify` exited 0 with no error line (`snapshot verify --json` exited
|
||||
1, also without one), so a `--cron` run that never finished looked
|
||||
like a success. A command stopped by either signal now exits 130 and
|
||||
prints `interrupted before the command finished` on stderr, under
|
||||
`--cron` and `--json` too.
|
||||
|
||||
- 2026-10-07: Corrected documentation, help text and comments that were
|
||||
false about the code
|
||||
([issue #233](https://git.eeqj.de/sneak/vaultik/issues/233)). A blob
|
||||
is not streamed to storage. The README, `ARCHITECTURE.md` and
|
||||
`config.example.yml` now say a backup needs free temporary space,
|
||||
because each blob is written whole to a temporary file before it is
|
||||
uploaded (up to about `blob_size_limit`; an rclone destination that
|
||||
cannot stream uploads needs about twice that) and the metadata export
|
||||
writes copies of the local index. Temporary files go to `$TMPDIR`
|
||||
(default `/tmp`); with `TMPDIR` unset, SQLite writes one of those
|
||||
copies to `/var/tmp`. Also corrected: the snapshot ID format, what restore
|
||||
reads and how incomplete snapshots are removed in `docs/DATAMODEL.md`,
|
||||
what `source_path` holds, the `index_path` and config file defaults,
|
||||
what `snapshot remove` cleans up, how the release gets its Go
|
||||
toolchain, and the `script/release` and `script/fmt-check` comments.
|
||||
|
||||
- 2026-10-07: Cut the time the `internal/vaultik` and `internal/database`
|
||||
tests take ([issue #235](https://git.eeqj.de/sneak/vaultik/issues/235)).
|
||||
Most of the `internal/vaultik` time went to 24 tests that ran one at a
|
||||
time only because they call `log.Initialize`; they now call it before
|
||||
`t.Parallel()`, as the package's other tests do. `TestLargeDatasets`
|
||||
committed each of its 1,500 inserts on its own and now makes them in
|
||||
one transaction. `TestDedupOnlySnapshotRestores` gives its second
|
||||
backup its own snapshot name instead of sleeping past the one-second
|
||||
timestamp in the snapshot ID.
|
||||
|
||||
- 2026-10-07: Made two messages say only what is true
|
||||
([issue #240](https://git.eeqj.de/sneak/vaultik/issues/240)). A config
|
||||
file that others can read was warned about as containing S3
|
||||
credentials even when it set none, as a `file://` config does. The
|
||||
warning now says the file may contain S3 credentials only when
|
||||
`s3.access_key_id` or `s3.secret_access_key` is set, since either may
|
||||
come from a `${...}` reference rather than the file, and otherwise
|
||||
says the file is readable by others. `snapshot purge` against a
|
||||
destination store it could not list gave an error with
|
||||
`listing remote snapshots:` in it twice; the prefix now appears once.
|
||||
|
||||
- 2026-10-07: Made `s3.part_size` set the multipart upload part size
|
||||
([issue #232](https://git.eeqj.de/sneak/vaultik/issues/232)). It was
|
||||
loaded and defaulted but never passed to the S3 client, whose uploader
|
||||
used a fixed 10MiB part. It now reaches the uploader for `storage_url`
|
||||
and for the `s3.*` fields, and a part size S3 refuses, below 5MiB or
|
||||
above 5GiB, fails at config load. The docs gave the default as `5MB`,
|
||||
which the config file reads as 5,000,000 bytes, below the minimum; they
|
||||
now say `5MiB`.
|
||||
above 5GiB, `0` included, fails at config load. A blob too large for
|
||||
S3's limit of 10,000 parts at the configured size is uploaded in larger
|
||||
parts. The docs gave the default as `5MB`, which the config file reads
|
||||
as 5,000,000 bytes, below the minimum; they now say `5MiB`.
|
||||
|
||||
- 2026-10-07: Made per-name retention work when the hostname contains `_`
|
||||
([issue #230](https://git.eeqj.de/sneak/vaultik/issues/230)). A
|
||||
|
||||
+11
-2
@@ -288,14 +288,17 @@ storage_url: "rclone://myremote/path/to/backups"
|
||||
#
|
||||
# # Part size for multipart uploads
|
||||
# # Minimum 5MiB, maximum 5GiB; affects memory usage during upload
|
||||
# # A blob too large for 10,000 parts of this size gets larger parts
|
||||
# # Supports: 10MB, 16MiB, 100MiB, etc. (5MB is below the minimum)
|
||||
# # Default: 5MiB
|
||||
# #part_size: 5MiB
|
||||
|
||||
# Path to local SQLite index database
|
||||
# This database tracks file state for incremental backups
|
||||
# Default: /var/lib/vaultik/index.sqlite
|
||||
#index_path: /var/lib/vaultik/index.sqlite
|
||||
# Default: the platform data directory, e.g.
|
||||
# macOS: ~/Library/Application Support/vaultik/index.sqlite
|
||||
# Linux: ~/.local/share/vaultik/index.sqlite
|
||||
#index_path: /path/to/index.sqlite
|
||||
|
||||
# Average chunk size for content-defined chunking
|
||||
# Smaller chunks = better deduplication but more metadata
|
||||
@@ -310,6 +313,12 @@ storage_url: "rclone://myremote/path/to/backups"
|
||||
# Chunking uses no secret (the FastCDC parameters are fixed and public). At a
|
||||
# large limit a blob holds hundreds of chunks, so individual chunk lengths are
|
||||
# not visible in its size; lowering the limit toward chunk_size exposes them.
|
||||
# A backup needs free temporary space, because each blob is written whole to
|
||||
# a temporary file before it is uploaded (up to about blob_size_limit; an
|
||||
# rclone destination that cannot stream uploads needs about twice that) and
|
||||
# the metadata export writes copies of the local index. Temporary files go to
|
||||
# $TMPDIR (default /tmp); with TMPDIR unset, SQLite writes one of those copies
|
||||
# to /var/tmp.
|
||||
# Supports: 1GB, 10G, 500MB, 1GiB, etc.
|
||||
# Default: 10GB
|
||||
#blob_size_limit: 10GB
|
||||
|
||||
+5
-6
@@ -111,7 +111,7 @@ Maps chunks to the blobs that contain them.
|
||||
Tracks backup snapshots.
|
||||
|
||||
**Columns:**
|
||||
- `id` (TEXT PRIMARY KEY) - Snapshot ID (format: hostname-YYYYMMDD-HHMMSSZ)
|
||||
- `id` (TEXT PRIMARY KEY) - Snapshot ID (format: `hostname_name_timestamp`, e.g. `server1_home_2025-06-01T12:00:00Z`: the hostname up to its first `.`, the snapshot name, and an RFC 3339 UTC timestamp)
|
||||
- `hostname` (TEXT) - Hostname where backup was created
|
||||
- `vaultik_version` (TEXT) - Version of Vaultik used
|
||||
- `vaultik_git_revision` (TEXT) - Git revision of Vaultik used
|
||||
@@ -218,8 +218,8 @@ The `{remote-key}` directory name is a one-way hash of the human snapshot ID, so
|
||||
### 4. Restore Process
|
||||
|
||||
The restore process doesn't use the local database. Instead:
|
||||
1. Downloads snapshot metadata from S3
|
||||
2. Downloads required blobs based on manifest
|
||||
1. Downloads and decrypts the snapshot's metadata database (`db.zst.age`) from S3
|
||||
2. Downloads the blobs holding the chunks of the files being restored, found through that database's `blob_chunks` table; the manifest is not read
|
||||
3. Reconstructs files from decrypted and decompressed chunks
|
||||
|
||||
### 5. Pruning
|
||||
@@ -232,9 +232,8 @@ The restore process doesn't use the local database. Instead:
|
||||
|
||||
Before each backup:
|
||||
1. Query incomplete snapshots (where `completed_at IS NULL`)
|
||||
2. Check if metadata exists in S3
|
||||
3. If no metadata, delete snapshot and all associations
|
||||
4. Clean up orphaned files, chunks, and blobs
|
||||
2. Delete each one and all its associations, without checking S3 for its metadata
|
||||
3. Clean up orphaned files, chunks, and blobs
|
||||
|
||||
## Repository Pattern
|
||||
|
||||
|
||||
+39
-14
@@ -206,6 +206,10 @@ func RunApp(ctx context.Context, app *fx.App) error {
|
||||
// RunOperation through cobra to Entry.
|
||||
var errReported = errors.New("operation failed")
|
||||
|
||||
// errInterrupted marks an operation that SIGINT or SIGTERM stopped
|
||||
// before it finished. Entry shows it and returns exitCodeInterrupted.
|
||||
var errInterrupted = errors.New("interrupted before the command finished")
|
||||
|
||||
// RunOperation runs op against the Vaultik instance inside the fx app
|
||||
// and turns a failure into a returned error rather than an os.Exit from
|
||||
// within the goroutine. An os.Exit there skipped main's deferred
|
||||
@@ -220,17 +224,21 @@ var errReported = errors.New("operation failed")
|
||||
// interrupt OnStop cancels op and waits for the goroutine to return, so
|
||||
// op's cleanup (removing decrypted scratch files) runs before the
|
||||
// process exits; the wait is bounded by shutdownTimeout. report is
|
||||
// called with a non-canceled failure so the caller can show it to the
|
||||
// user before it becomes errReported. A context cancellation is the
|
||||
// interrupt path, not a failure: it is neither reported nor counted as
|
||||
// one.
|
||||
// called with a failure so the caller can show it to the user before
|
||||
// it becomes errReported.
|
||||
//
|
||||
// The run counts as interrupted unless op returned, without an
|
||||
// interrupt having cancelled it, before RunWithApp returned. An
|
||||
// interrupted op is not reported, whatever it returned; RunOperation
|
||||
// returns errInterrupted instead.
|
||||
func RunOperation(
|
||||
ctx context.Context, opts AppOptions,
|
||||
op func(v *vaultik.Vaultik) error, report func(err error),
|
||||
) error {
|
||||
var (
|
||||
mu sync.Mutex
|
||||
failed bool
|
||||
mu sync.Mutex
|
||||
finished bool // op returned before any interrupt cancelled it
|
||||
failed bool // op finished with an error
|
||||
)
|
||||
|
||||
opts.Invokes = append(opts.Invokes,
|
||||
@@ -241,11 +249,21 @@ func RunOperation(
|
||||
OnStart: func(_ context.Context) error {
|
||||
stop = v.StartOperation(func() {
|
||||
err := op(v)
|
||||
if err != nil && !errors.Is(err, context.Canceled) {
|
||||
report(err)
|
||||
|
||||
// Only stop, called from OnStop below, cancels the
|
||||
// Vaultik context, so a live context means no
|
||||
// interrupt cancelled op. Check the context, not
|
||||
// err: an interrupted op need not return
|
||||
// context.Canceled (`snapshot verify --json`
|
||||
// returns a verification failure).
|
||||
if v.Context().Err() == nil {
|
||||
if err != nil {
|
||||
report(err)
|
||||
}
|
||||
|
||||
mu.Lock()
|
||||
failed = true
|
||||
finished = true
|
||||
failed = err != nil
|
||||
mu.Unlock()
|
||||
}
|
||||
|
||||
@@ -278,16 +296,23 @@ func RunOperation(
|
||||
return err
|
||||
}
|
||||
|
||||
// The goroutine sets failed before triggering the shutdown that lets
|
||||
// RunWithApp return, so the write is in place by the time we read it.
|
||||
// RunWithApp returns only after the app was asked to stop, either by
|
||||
// an interrupt or by the goroutine's Shutdown call. When op finished
|
||||
// without being cancelled, the goroutine set finished before that
|
||||
// call. So if finished is unset here, an interrupt stopped the app,
|
||||
// and op either returned after it was cancelled or is still running
|
||||
// because the shutdown timed out.
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
|
||||
if failed {
|
||||
switch {
|
||||
case !finished:
|
||||
return errInterrupted
|
||||
case failed:
|
||||
return errReported
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// runVaultikApp runs the standard single-operation command lifecycle
|
||||
|
||||
+15
-4
@@ -15,6 +15,11 @@ import (
|
||||
// the startup banner.
|
||||
const shortCommitLen = 12
|
||||
|
||||
// exitCodeInterrupted is the exit status of a command that SIGINT or
|
||||
// SIGTERM stopped. It is 128 plus SIGINT's number, 2, which is what a
|
||||
// shell reports for a command stopped by Ctrl-C.
|
||||
const exitCodeInterrupted = 130
|
||||
|
||||
// Entry is the main entry point for the CLI application.
|
||||
// It prints the startup banner to stderr (unless a banner-suppressing
|
||||
// flag is present in os.Args — see bannerSuppressedInArgs), executes the
|
||||
@@ -23,9 +28,10 @@ const shortCommitLen = 12
|
||||
// The banner goes to stderr because stdout carries only the output the
|
||||
// user asked for, such as a completion script or a `config get` value.
|
||||
//
|
||||
// It returns the process exit code (0 on success, 1 on error) rather
|
||||
// than calling os.Exit, so that main's deferred profile writers run
|
||||
// before the process ends. See run in cmd/vaultik/main.go.
|
||||
// It returns the process exit code (0 on success, 130 when interrupted,
|
||||
// 1 on any other error) rather than calling os.Exit, so that main's
|
||||
// deferred profile writers run before the process ends. See run in
|
||||
// cmd/vaultik/main.go.
|
||||
func Entry() int {
|
||||
emitStartupBanner(os.Args[1:], os.Stderr)
|
||||
|
||||
@@ -39,11 +45,16 @@ func Entry() int {
|
||||
// document instead); errReported says so. Printing it again
|
||||
// here would double the error line.
|
||||
// Every other error — bad arguments, a config that would not
|
||||
// load — reaches Entry unreported, so it is shown here.
|
||||
// load, an interrupt — reaches Entry unreported, so it is shown
|
||||
// here.
|
||||
if !errors.Is(err, errReported) {
|
||||
ReportErrorf("%s", err.Error())
|
||||
}
|
||||
|
||||
if errors.Is(err, errInterrupted) {
|
||||
return exitCodeInterrupted
|
||||
}
|
||||
|
||||
return 1
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
package cli //nolint:testpackage // shares runEntry and the argument constants
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/adrg/xdg"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// stalledStoreConfig is hermeticConfig with an s3:// destination store
|
||||
// in place of the file:// one. The server behind it accepts any
|
||||
// credentials.
|
||||
const stalledStoreConfig = `age_recipients:
|
||||
- age1278m9q7dp3chsh2dcy82qk27v047zywyvtxwnj4cvt0z65jw6a7q5dqhfj
|
||||
snapshots:
|
||||
test:
|
||||
paths:
|
||||
- %s
|
||||
storage_url: s3://bucket?endpoint=%s&ssl=false
|
||||
s3:
|
||||
access_key_id: key
|
||||
secret_access_key: secret
|
||||
index_path: %s
|
||||
hostname: test-host
|
||||
`
|
||||
|
||||
// interruptRepeat is how often interruptOnFirstRequest sends SIGINT.
|
||||
const interruptRepeat = 50 * time.Millisecond
|
||||
|
||||
// TestEntryInterruptedRun sends SIGINT to the test process while a
|
||||
// command waits on the destination store, and checks that Entry returns
|
||||
// 130 and prints one line on stderr saying the run was interrupted. The
|
||||
// store is a local HTTP server that holds every request open, so the
|
||||
// command is always mid-operation when the signal arrives. The two
|
||||
// cases cover --cron and --json, which silence other output.
|
||||
//
|
||||
// Not parallel: it signals the process and replaces os.Args, os.Stdout,
|
||||
// os.Stderr and the xdg globals.
|
||||
//
|
||||
//nolint:paralleltest // signals the process and replaces process globals
|
||||
func TestEntryInterruptedRun(t *testing.T) {
|
||||
for _, testCase := range []struct {
|
||||
name string
|
||||
args []string
|
||||
}{
|
||||
{
|
||||
name: "snapshot create --cron",
|
||||
args: []string{cmdSnapshot, cmdCreate, "--cron"},
|
||||
},
|
||||
{
|
||||
name: "snapshot verify --json",
|
||||
args: []string{cmdSnapshot, cmdVerify, someSnapshotID, flagJSON},
|
||||
},
|
||||
} {
|
||||
t.Run(testCase.name, func(t *testing.T) {
|
||||
endpoint, requestArrived := startStalledStore(t)
|
||||
configPath := writeStalledStoreConfig(t, endpoint)
|
||||
interruptOnFirstRequest(t, requestArrived)
|
||||
|
||||
code, _, stderr := runEntry(t,
|
||||
append([]string{flagConfig, configPath}, testCase.args...)...)
|
||||
|
||||
assert.Equal(t, 130, code)
|
||||
assert.Equal(t, 1,
|
||||
strings.Count(stderr, errInterrupted.Error()), stderr)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// interruptOnFirstRequest sends SIGINT to the test process every
|
||||
// interruptRepeat, from the first request to the destination store until
|
||||
// the test ends. One signal is not enough: the command can reach the
|
||||
// store before fx has started catching signals. The test catches SIGINT
|
||||
// too, so that a signal fx is not catching does not kill the test
|
||||
// binary.
|
||||
func interruptOnFirstRequest(t *testing.T, requestArrived <-chan struct{}) {
|
||||
t.Helper()
|
||||
|
||||
self, err := os.FindProcess(os.Getpid())
|
||||
require.NoError(t, err)
|
||||
|
||||
caught := make(chan os.Signal, 1)
|
||||
signal.Notify(caught, os.Interrupt)
|
||||
|
||||
testEnded := make(chan struct{})
|
||||
senderDone := make(chan struct{})
|
||||
|
||||
// Stop catching SIGINT only after the sender has returned. The sender
|
||||
// waits for each SIGINT it sends to arrive on caught; one still on
|
||||
// its way after signal.Stop would kill the test binary.
|
||||
t.Cleanup(func() {
|
||||
close(testEnded)
|
||||
<-senderDone
|
||||
signal.Stop(caught)
|
||||
})
|
||||
|
||||
go func() {
|
||||
defer close(senderDone)
|
||||
|
||||
select {
|
||||
case <-requestArrived:
|
||||
case <-testEnded:
|
||||
return
|
||||
}
|
||||
|
||||
ticker := time.NewTicker(interruptRepeat)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
// Empty caught, so that the receive below waits for this
|
||||
// SIGINT rather than an earlier one.
|
||||
select {
|
||||
case <-caught:
|
||||
default:
|
||||
}
|
||||
|
||||
sendErr := self.Signal(os.Interrupt)
|
||||
if sendErr != nil {
|
||||
t.Errorf("sending SIGINT: %v", sendErr)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
<-caught
|
||||
|
||||
select {
|
||||
case <-testEnded:
|
||||
return
|
||||
case <-ticker.C:
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// startStalledStore starts an HTTP server that never answers: each
|
||||
// request is held until the client gives up on it or the test ends.
|
||||
// It returns the server's host:port and a channel that receives a value
|
||||
// when the first request arrives.
|
||||
func startStalledStore(t *testing.T) (string, <-chan struct{}) {
|
||||
t.Helper()
|
||||
|
||||
requestArrived := make(chan struct{}, 1)
|
||||
release := make(chan struct{})
|
||||
|
||||
server := httptest.NewServer(http.HandlerFunc(
|
||||
func(_ http.ResponseWriter, r *http.Request) {
|
||||
select {
|
||||
case requestArrived <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
|
||||
select {
|
||||
case <-r.Context().Done():
|
||||
case <-release:
|
||||
}
|
||||
}))
|
||||
|
||||
// Cleanups run last-registered first, so release lets any held
|
||||
// request return before Close waits for it.
|
||||
t.Cleanup(server.Close)
|
||||
t.Cleanup(func() { close(release) })
|
||||
|
||||
return server.Listener.Addr().String(), requestArrived
|
||||
}
|
||||
|
||||
// writeStalledStoreConfig writes a config whose destination store is the
|
||||
// server at endpoint and whose snapshot source holds one small file, so
|
||||
// that `snapshot create` has a blob to upload. Returns the config path.
|
||||
func writeStalledStoreConfig(t *testing.T, endpoint string) string {
|
||||
t.Helper()
|
||||
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, "config.yml")
|
||||
sourceDir := filepath.Join(dir, "source")
|
||||
|
||||
require.NoError(t, os.Mkdir(sourceDir, 0o750))
|
||||
require.NoError(t, os.WriteFile(filepath.Join(sourceDir, "file.txt"),
|
||||
[]byte("contents"), 0o600))
|
||||
|
||||
contents := fmt.Sprintf(stalledStoreConfig,
|
||||
sourceDir, endpoint, filepath.Join(dir, "index.sqlite"))
|
||||
|
||||
require.NoError(t,
|
||||
os.WriteFile(configPath, []byte(contents), configFileMode))
|
||||
|
||||
// The PID lock lives under xdg.DataHome, which xdg resolves at
|
||||
// package init; point it at the temp dir so the test neither
|
||||
// touches nor collides with the real one.
|
||||
t.Setenv("XDG_DATA_HOME", filepath.Join(dir, "data"))
|
||||
xdg.Reload()
|
||||
t.Cleanup(xdg.Reload)
|
||||
|
||||
return configPath
|
||||
}
|
||||
@@ -22,9 +22,11 @@ scans every snapshot manifest in the destination store, builds the
|
||||
set of still-referenced blob hashes, and deletes any blob not in that
|
||||
set.
|
||||
|
||||
Snapshot create --prune and snapshot remove run the same cleanup
|
||||
automatically; this command is the manual entry point for the same
|
||||
work (e.g. after a crashed backup or to reclaim storage).`,
|
||||
Snapshot create --prune runs the same cleanup automatically; this
|
||||
command is the manual entry point for the same work (e.g. after a
|
||||
crashed backup or to reclaim storage). Snapshot remove leaves blobs in
|
||||
place; run this command afterwards to delete the ones no longer
|
||||
referenced.`,
|
||||
Args: cobra.NoArgs,
|
||||
RunE: func(cmd *cobra.Command, _ []string) error {
|
||||
// Use unified config resolution
|
||||
|
||||
@@ -66,8 +66,9 @@ func newSnapshotCreateCommand() *cobra.Command {
|
||||
If snapshot names are provided, only those snapshots are created.
|
||||
If no names are provided, all configured snapshots are created.
|
||||
|
||||
Config is located at /etc/vaultik/config.yml by default, but can be overridden by
|
||||
specifying a path using --config or by setting VAULTIK_CONFIG to a path.`,
|
||||
The config is read from the path given by --config or VAULTIK_CONFIG;
|
||||
otherwise from the platform config directory (~/.config/vaultik/config.yml
|
||||
on Linux), then /etc/vaultik/config.yml.`,
|
||||
Args: cobra.ArbitraryArgs,
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
// Pass snapshot names from args
|
||||
|
||||
@@ -248,6 +248,7 @@ func Load(path string) (*Config, error) {
|
||||
ChunkSize: defaultChunkSize,
|
||||
IndexPath: filepath.Join(xdg.DataHome, appName, "index.sqlite"),
|
||||
CompressionLevel: defaultCompressionLevel,
|
||||
S3: S3Config{PartSize: defaultS3PartSize},
|
||||
}
|
||||
|
||||
// Convert smartconfig data to YAML then unmarshal
|
||||
@@ -298,17 +299,13 @@ func Load(path string) (*Config, error) {
|
||||
cfg.S3.Region = "us-east-1"
|
||||
}
|
||||
|
||||
if cfg.S3.PartSize == 0 {
|
||||
cfg.S3.PartSize = defaultS3PartSize
|
||||
}
|
||||
|
||||
// Check config file permissions (warn if world or group readable)
|
||||
//nolint:gosec // G703: config path is operator-supplied by design
|
||||
info, statErr := os.Stat(path)
|
||||
if statErr == nil {
|
||||
mode := info.Mode().Perm()
|
||||
if mode&0044 != 0 { // group or world readable
|
||||
log.Warn("Config file has insecure permissions (contains S3 credentials)",
|
||||
log.Warn(cfg.readableByOthersWarning(),
|
||||
"path", path,
|
||||
"mode", fmt.Sprintf("%04o", mode),
|
||||
"recommendation", "chmod 600 "+path)
|
||||
@@ -421,6 +418,18 @@ func (c *Config) setAgeSecretKey() {
|
||||
}
|
||||
}
|
||||
|
||||
// readableByOthersWarning is the warning Load logs when others can read
|
||||
// the config file. It says "may contain" because the S3 credentials are
|
||||
// seen only after smartconfig has replaced any ${...} reference in the
|
||||
// file with its value, so a set credential need not be in the file.
|
||||
func (c *Config) readableByOthersWarning() string {
|
||||
if c.S3.AccessKeyID != "" || c.S3.SecretAccessKey != "" {
|
||||
return "Config file is readable by others and may contain S3 credentials"
|
||||
}
|
||||
|
||||
return "Config file is readable by others"
|
||||
}
|
||||
|
||||
// validateStorage validates storage configuration.
|
||||
// If StorageURL is set, it takes precedence. S3 URLs require credentials.
|
||||
// File URLs don't require any S3 configuration.
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"testing"
|
||||
|
||||
"sneak.berlin/go/vaultik/internal/chunker"
|
||||
"sneak.berlin/go/vaultik/internal/log"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -288,6 +289,54 @@ func TestValidateS3PartSize(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestLoadS3PartSize checks that a config file without s3.part_size loads
|
||||
// with the 5MiB default, and that an explicit 0 fails at load like any other
|
||||
// part size S3 refuses.
|
||||
func TestLoadS3PartSize(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const withoutPartSize = "snapshots:\n" +
|
||||
" test:\n" +
|
||||
" paths: [/tmp/vaultik-test-source]\n" +
|
||||
"storage_url: file:///tmp/vaultik-test-storage\n"
|
||||
|
||||
writeConfig := func(t *testing.T, text string) string {
|
||||
t.Helper()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "config.yml")
|
||||
|
||||
err := os.WriteFile(path, []byte(text), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("write config: %v", err)
|
||||
}
|
||||
|
||||
return path
|
||||
}
|
||||
|
||||
t.Run("absent loads as 5MiB", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cfg, err := Load(writeConfig(t, withoutPartSize))
|
||||
if err != nil {
|
||||
t.Fatalf("Load() unexpected error: %v", err)
|
||||
}
|
||||
|
||||
if cfg.S3.PartSize != defaultS3PartSize {
|
||||
t.Errorf("s3.part_size = %d, want %d",
|
||||
cfg.S3.PartSize, defaultS3PartSize)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("0 is rejected", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
_, err := Load(writeConfig(t, withoutPartSize+"s3:\n part_size: 0\n"))
|
||||
if !errors.Is(err, errBadS3PartSize) {
|
||||
t.Fatalf("Load() error = %v, want errBadS3PartSize", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// TestValidateAgeRecipients checks that recipients are parsed at config load
|
||||
// (a bad entry fails immediately, not mid-backup) and that no invalid entry —
|
||||
// least of all a pasted secret key — is echoed in the error. An empty list
|
||||
@@ -411,3 +460,125 @@ func TestAgeSecretKeySourceName(t *testing.T) {
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// loadReadableConfig writes configYAML to a file that others can read,
|
||||
// loads it, and returns what the logger wrote to stderr meanwhile. The
|
||||
// logger writes to the os.Stderr it finds when it is initialized, so
|
||||
// os.Stderr is pointed at a file first. Not parallel-safe: os.Stderr and
|
||||
// the logger are process-global.
|
||||
func loadReadableConfig(t *testing.T, configYAML string) string {
|
||||
t.Helper()
|
||||
|
||||
dir := t.TempDir()
|
||||
configPath := filepath.Join(dir, "config.yml")
|
||||
stderrPath := filepath.Join(dir, "stderr")
|
||||
|
||||
err := os.WriteFile(configPath, []byte(configYAML), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("writing config: %v", err)
|
||||
}
|
||||
|
||||
//nolint:gosec // G302: the test needs a config file others can read
|
||||
err = os.Chmod(configPath, 0o644)
|
||||
if err != nil {
|
||||
t.Fatalf("chmod config: %v", err)
|
||||
}
|
||||
|
||||
stderrFile, err := os.Create(stderrPath) //nolint:gosec // G304: test temp path
|
||||
if err != nil {
|
||||
t.Fatalf("creating stderr file: %v", err)
|
||||
}
|
||||
|
||||
previous := os.Stderr
|
||||
os.Stderr = stderrFile
|
||||
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
_, loadErr := Load(configPath)
|
||||
|
||||
os.Stderr = previous
|
||||
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
_ = stderrFile.Close()
|
||||
|
||||
if loadErr != nil {
|
||||
t.Fatalf("Load() error = %v", loadErr)
|
||||
}
|
||||
|
||||
captured, err := os.ReadFile(stderrPath) //nolint:gosec // G304: test temp path
|
||||
if err != nil {
|
||||
t.Fatalf("reading stderr file: %v", err)
|
||||
}
|
||||
|
||||
return string(captured)
|
||||
}
|
||||
|
||||
// TestLoadWarnsReadableConfigWithoutS3Credentials checks that a config
|
||||
// file others can read, holding no S3 credentials, is warned about
|
||||
// without a claim that it holds them.
|
||||
//
|
||||
//nolint:paralleltest // loadReadableConfig replaces os.Stderr
|
||||
func TestLoadWarnsReadableConfigWithoutS3Credentials(t *testing.T) {
|
||||
stderr := loadReadableConfig(t, `
|
||||
storage_url: file:///var/backups/vaultik
|
||||
snapshots:
|
||||
home:
|
||||
paths:
|
||||
- /home
|
||||
`)
|
||||
|
||||
if !strings.Contains(stderr, "Config file is readable by others") {
|
||||
t.Errorf("expected a warning that the file is readable by others, got %q",
|
||||
stderr)
|
||||
}
|
||||
|
||||
if strings.Contains(stderr, "S3 credentials") {
|
||||
t.Errorf("warning names S3 credentials the file does not set: %q", stderr)
|
||||
}
|
||||
}
|
||||
|
||||
// TestLoadWarnsReadableConfigWithS3Credentials checks that a config file
|
||||
// others can read and that sets S3 credentials, as values or as ${ENV:...}
|
||||
// references, is warned about as one that may contain them.
|
||||
//
|
||||
//nolint:paralleltest // loadReadableConfig replaces os.Stderr
|
||||
func TestLoadWarnsReadableConfigWithS3Credentials(t *testing.T) {
|
||||
t.Setenv("VAULTIK_TEST_ACCESS_KEY_ID", "test-access-key")
|
||||
t.Setenv("VAULTIK_TEST_SECRET_ACCESS_KEY", "test-secret-key")
|
||||
|
||||
configs := map[string]string{
|
||||
"values": `
|
||||
storage_url: s3://bucket/prefix?endpoint=s3.example.com
|
||||
s3:
|
||||
access_key_id: test-access-key
|
||||
secret_access_key: test-secret-key
|
||||
snapshots:
|
||||
home:
|
||||
paths:
|
||||
- /home
|
||||
`,
|
||||
"references": `
|
||||
storage_url: s3://bucket/prefix?endpoint=s3.example.com
|
||||
s3:
|
||||
access_key_id: ${ENV:VAULTIK_TEST_ACCESS_KEY_ID}
|
||||
secret_access_key: ${ENV:VAULTIK_TEST_SECRET_ACCESS_KEY}
|
||||
snapshots:
|
||||
home:
|
||||
paths:
|
||||
- /home
|
||||
`,
|
||||
}
|
||||
|
||||
for name, configYAML := range configs {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
stderr := loadReadableConfig(t, configYAML)
|
||||
|
||||
if !strings.Contains(stderr,
|
||||
"Config file is readable by others and may contain S3 credentials") {
|
||||
t.Errorf("expected a warning naming the S3 credentials, got %q",
|
||||
stderr)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,8 +14,7 @@ type File struct {
|
||||
ID types.FileID // UUID primary key
|
||||
Path types.FilePath // Absolute path of the file
|
||||
|
||||
// SourcePath is the source directory this file came from (used for
|
||||
// restore path stripping).
|
||||
// SourcePath is the source directory this file came from.
|
||||
SourcePath types.SourcePath
|
||||
MTime time.Time
|
||||
Size int64
|
||||
|
||||
@@ -3,6 +3,7 @@ package database
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
@@ -367,7 +368,7 @@ func verifyBlobNullUploadTS(
|
||||
}
|
||||
|
||||
// createLargeDatasetFiles creates fileCount files and adds every other
|
||||
// one to the snapshot.
|
||||
// one to the snapshot, in one transaction as a backup writes them.
|
||||
func createLargeDatasetFiles(
|
||||
t *testing.T,
|
||||
repos *Repositories,
|
||||
@@ -376,31 +377,38 @@ func createLargeDatasetFiles(
|
||||
) {
|
||||
t.Helper()
|
||||
|
||||
ctx := context.Background()
|
||||
start := time.Now()
|
||||
|
||||
for i := range fileCount {
|
||||
file := &File{
|
||||
Path: types.FilePath(fmt.Sprintf("/large/file%05d.txt", i)),
|
||||
MTime: time.Now(),
|
||||
Size: int64(i * 1024),
|
||||
Mode: 0644,
|
||||
UID: uint32(1000 + (i % 10)),
|
||||
GID: uint32(1000 + (i % 10)),
|
||||
}
|
||||
err := repos.WithTx(context.Background(),
|
||||
func(ctx context.Context, tx *sql.Tx) error {
|
||||
for i := range fileCount {
|
||||
file := &File{
|
||||
Path: types.FilePath(fmt.Sprintf("/large/file%05d.txt", i)),
|
||||
MTime: time.Now(),
|
||||
Size: int64(i * 1024),
|
||||
Mode: 0644,
|
||||
UID: uint32(1000 + (i % 10)),
|
||||
GID: uint32(1000 + (i % 10)),
|
||||
}
|
||||
|
||||
err := repos.Files.Create(ctx, nil, file)
|
||||
if err != nil {
|
||||
t.Fatalf("failed to create file %d: %v", i, err)
|
||||
}
|
||||
err := repos.Files.Create(ctx, tx, file)
|
||||
if err != nil {
|
||||
return fmt.Errorf("creating file %d: %w", i, err)
|
||||
}
|
||||
|
||||
// Add half to snapshot
|
||||
if i%2 == 0 {
|
||||
err = repos.Snapshots.AddFileByID(ctx, nil, snapshotID, file.ID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
// Add half to snapshot
|
||||
if i%2 == 0 {
|
||||
err = repos.Snapshots.AddFileByID(ctx, tx, snapshotID, file.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
t.Logf("Created %d files in %v", fileCount, time.Since(start))
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
CREATE TABLE IF NOT EXISTS files (
|
||||
id TEXT PRIMARY KEY, -- UUID
|
||||
path TEXT NOT NULL UNIQUE,
|
||||
source_path TEXT NOT NULL DEFAULT '', -- The source directory this file came from (for restore path stripping)
|
||||
source_path TEXT NOT NULL DEFAULT '', -- The source directory this file came from
|
||||
mtime INTEGER NOT NULL, -- whole seconds since the Unix epoch
|
||||
mtime_nsec INTEGER NOT NULL, -- nanoseconds within that second, 0 to 999999999
|
||||
size INTEGER NOT NULL,
|
||||
|
||||
+17
-1
@@ -44,6 +44,7 @@ type Config struct {
|
||||
SecretAccessKey string
|
||||
Region string
|
||||
// PartSize is the size in bytes of each part of a multipart upload.
|
||||
// An upload too large for S3's limit of 10,000 parts gets larger parts.
|
||||
PartSize int64
|
||||
}
|
||||
|
||||
@@ -130,7 +131,7 @@ func (c *Client) PutObjectWithProgress(
|
||||
|
||||
// Create an uploader with the S3 client
|
||||
uploader := manager.NewUploader(c.s3Client, func(u *manager.Uploader) {
|
||||
u.PartSize = c.partSize
|
||||
u.PartSize = uploadPartSize(c.partSize, size)
|
||||
})
|
||||
|
||||
// Create a progress reader that tracks upload progress
|
||||
@@ -151,6 +152,21 @@ func (c *Client) PutObjectWithProgress(
|
||||
return err
|
||||
}
|
||||
|
||||
// uploadPartSize returns the part size for an upload of size bytes: the
|
||||
// configured part size (the SDK default when zero), raised where needed so
|
||||
// the upload fits in S3's limit of 10,000 parts. The uploader cannot raise
|
||||
// it itself, because it cannot seek the progress reader to learn its size.
|
||||
func uploadPartSize(configured, size int64) int64 {
|
||||
if configured == 0 {
|
||||
configured = manager.DefaultUploadPartSize
|
||||
}
|
||||
|
||||
maxParts := int64(manager.MaxUploadParts)
|
||||
smallestThatFits := (size + maxParts - 1) / maxParts // rounded up
|
||||
|
||||
return max(configured, smallestThatFits)
|
||||
}
|
||||
|
||||
// GetObject downloads an object from S3 with the specified key.
|
||||
// The key is automatically prefixed with the configured prefix.
|
||||
// Returns a ReadCloser containing the object data. The caller must
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
package s3
|
||||
|
||||
import "testing"
|
||||
|
||||
// TestUploadPartSize checks that an upload too large for 10,000 parts of the
|
||||
// configured size gets parts just large enough to fit in 10,000.
|
||||
func TestUploadPartSize(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const mib = 1024 * 1024
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
configured int64
|
||||
size int64
|
||||
want int64
|
||||
}{
|
||||
{
|
||||
name: "an upload that fits keeps the configured size",
|
||||
configured: 5 * mib,
|
||||
size: 10 * 1024 * mib,
|
||||
want: 5 * mib,
|
||||
},
|
||||
{
|
||||
name: "exactly 10,000 parts keeps the configured size",
|
||||
configured: 6 * mib,
|
||||
size: 10_000 * 6 * mib,
|
||||
want: 6 * mib,
|
||||
},
|
||||
{
|
||||
name: "one byte more than 10,000 parts adds a byte to each",
|
||||
configured: 6 * mib,
|
||||
size: 10_000*6*mib + 1,
|
||||
want: 6*mib + 1,
|
||||
},
|
||||
{
|
||||
name: "zero means the SDK default of 5MiB",
|
||||
configured: 0,
|
||||
size: 1,
|
||||
want: 5 * mib,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
got := uploadPartSize(tt.configured, tt.size)
|
||||
if got != tt.want {
|
||||
t.Errorf("uploadPartSize(%d, %d) = %d, want %d",
|
||||
tt.configured, tt.size, got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -58,8 +58,8 @@ type Scanner struct {
|
||||
compressionLevel int
|
||||
ageRecipient string
|
||||
snapshotID string // Current snapshot being processed
|
||||
// currentSourcePath is the source directory being scanned (used for
|
||||
// restore path stripping).
|
||||
// currentSourcePath is the source directory being scanned, stored with
|
||||
// each file record.
|
||||
currentSourcePath string
|
||||
exclude []string // Glob patterns for files/directories to exclude
|
||||
compiledExclude []compiledPattern // Compiled glob patterns
|
||||
@@ -211,7 +211,6 @@ func (s *Scanner) Scan(
|
||||
ctx context.Context, path string, snapshotID string,
|
||||
) (*ScanResult, error) {
|
||||
s.snapshotID = snapshotID
|
||||
// Store source path for file records (used during restore)
|
||||
s.currentSourcePath = path
|
||||
result := &ScanResult{
|
||||
StartTime: time.Now().UTC(),
|
||||
@@ -890,9 +889,10 @@ func (s *Scanner) scanPhase(
|
||||
}
|
||||
|
||||
// Handle symlinks and directories
|
||||
if handled := s.recordSpecialEntry(
|
||||
filePath, info, existingFiles, collector, result); handled {
|
||||
return nil
|
||||
handled, err := s.recordSpecialEntry(
|
||||
filePath, info, existingFiles, collector, result)
|
||||
if handled {
|
||||
return err
|
||||
}
|
||||
|
||||
// Skip other non-regular files (devices, sockets, etc.)
|
||||
@@ -934,22 +934,25 @@ func (s *Scanner) scanPhase(
|
||||
}
|
||||
|
||||
// recordSpecialEntry records symlinks and directories (which have no
|
||||
// data to chunk) and reports whether it handled the entry.
|
||||
// data to chunk) and reports whether it handled the entry. For a symlink
|
||||
// whose target cannot be read it returns handleWalkError's result.
|
||||
func (s *Scanner) recordSpecialEntry(
|
||||
filePath string, info os.FileInfo,
|
||||
existingFiles map[string]struct{},
|
||||
collector *scanCollector, result *ScanResult,
|
||||
) bool {
|
||||
) (bool, error) {
|
||||
// Handle symlinks
|
||||
if info.Mode()&os.ModeSymlink != 0 {
|
||||
file := s.buildSymlinkEntry(filePath, info)
|
||||
if file != nil {
|
||||
existingFiles[filePath] = struct{}{}
|
||||
collector.addToProcess(filePath, info, file)
|
||||
s.updateScanEntryStats(result, true, info)
|
||||
file, err := s.buildSymlinkEntry(filePath, info)
|
||||
if err != nil {
|
||||
return true, s.handleWalkError(filePath, err)
|
||||
}
|
||||
|
||||
return true
|
||||
existingFiles[filePath] = struct{}{}
|
||||
collector.addToProcess(filePath, info, file)
|
||||
s.updateScanEntryStats(result, true, info)
|
||||
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// Handle directories (record for permission/ownership preservation
|
||||
@@ -959,10 +962,10 @@ func (s *Scanner) recordSpecialEntry(
|
||||
existingFiles[filePath] = struct{}{}
|
||||
collector.addToProcess(filePath, info, file)
|
||||
|
||||
return true
|
||||
return true, nil
|
||||
}
|
||||
|
||||
return false
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// handleWalkError deals with a filesystem error surfaced by the walk:
|
||||
@@ -1112,13 +1115,12 @@ func (s *Scanner) printScanProgressLine(
|
||||
}
|
||||
|
||||
// buildSymlinkEntry creates a File record for a symlink.
|
||||
// Returns nil if the link target cannot be read.
|
||||
func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.File {
|
||||
func (s *Scanner) buildSymlinkEntry(
|
||||
path string, info os.FileInfo,
|
||||
) (*database.File, error) {
|
||||
target, err := os.Readlink(path)
|
||||
if err != nil {
|
||||
log.Debug("Cannot read symlink target", "path", path, "error", err)
|
||||
|
||||
return nil
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var uid, gid uint32
|
||||
@@ -1137,7 +1139,7 @@ func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.Fil
|
||||
UID: uid,
|
||||
GID: gid,
|
||||
LinkTarget: types.FilePath(target),
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
// buildDirectoryEntry creates a File record for a directory.
|
||||
@@ -1198,9 +1200,8 @@ func (s *Scanner) checkFileInMemory(
|
||||
}
|
||||
|
||||
file := &database.File{
|
||||
ID: fileID,
|
||||
Path: types.FilePath(path),
|
||||
// Store source directory for restore path stripping
|
||||
ID: fileID,
|
||||
Path: types.FilePath(path),
|
||||
SourcePath: types.SourcePath(s.currentSourcePath),
|
||||
MTime: info.ModTime(),
|
||||
Size: info.Size(),
|
||||
|
||||
@@ -3,6 +3,7 @@ package snapshot_test
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
@@ -13,6 +14,7 @@ import (
|
||||
"github.com/spf13/afero"
|
||||
"sneak.berlin/go/vaultik/internal/database"
|
||||
"sneak.berlin/go/vaultik/internal/snapshot"
|
||||
"sneak.berlin/go/vaultik/internal/ui"
|
||||
)
|
||||
|
||||
// errSimTempFail is the one-time temp-file creation failure blobTempFailFs
|
||||
@@ -80,6 +82,30 @@ func (f *readFailFs) Open(name string) (afero.File, error) {
|
||||
return file, nil
|
||||
}
|
||||
|
||||
// linkRemovedAfterLstatFs is the real filesystem, except that the symlink at
|
||||
// target is removed right after the walk lstats it, as happens when a link is
|
||||
// deleted during a backup. The scanner's readlink of it then fails.
|
||||
type linkRemovedAfterLstatFs struct {
|
||||
afero.OsFs
|
||||
|
||||
t *testing.T
|
||||
target string
|
||||
}
|
||||
|
||||
func (f *linkRemovedAfterLstatFs) LstatIfPossible(
|
||||
name string,
|
||||
) (os.FileInfo, bool, error) {
|
||||
info, lstatCalled, err := f.OsFs.LstatIfPossible(name)
|
||||
if err == nil && name == f.target {
|
||||
rmErr := os.Remove(name)
|
||||
if rmErr != nil {
|
||||
f.t.Errorf("removing %s: %v", name, rmErr)
|
||||
}
|
||||
}
|
||||
|
||||
return info, lstatCalled, err
|
||||
}
|
||||
|
||||
// writeSkipErrorTestFile writes one file into fs with a fixed mtime.
|
||||
func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
|
||||
t.Helper()
|
||||
@@ -102,10 +128,11 @@ func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
|
||||
}
|
||||
}
|
||||
|
||||
// runSkipErrorScan scans /source on fs with the given skip-errors setting and
|
||||
// returns the repositories (for inspection) and the scan error.
|
||||
// runSkipErrorScan scans source on fs with the given skip-errors setting,
|
||||
// printing user-facing messages to uiw (nil discards them), and returns the
|
||||
// repositories (for inspection) and the scan error.
|
||||
func runSkipErrorScan(
|
||||
t *testing.T, fs afero.Fs, skipErrors bool,
|
||||
t *testing.T, fs afero.Fs, source string, skipErrors bool, uiw *ui.Writer,
|
||||
) (*database.Repositories, error) {
|
||||
t.Helper()
|
||||
|
||||
@@ -130,6 +157,7 @@ func runSkipErrorScan(
|
||||
MaxBlobSize: int64(1024 * 1024),
|
||||
CompressionLevel: 3,
|
||||
AgeRecipients: []string{testAgePublicKey},
|
||||
UI: uiw,
|
||||
SkipErrors: skipErrors,
|
||||
})
|
||||
|
||||
@@ -137,7 +165,7 @@ func runSkipErrorScan(
|
||||
snapshotID := "test-snapshot-skip-errors"
|
||||
createTestSnapshotRecord(ctx, t, repos, snapshotID)
|
||||
|
||||
_, err = scanner.Scan(ctx, "/source", snapshotID)
|
||||
_, err = scanner.Scan(ctx, source, snapshotID)
|
||||
|
||||
return repos, err
|
||||
}
|
||||
@@ -157,7 +185,7 @@ func TestScannerPackingFailureAbortsUnderSkipErrors(t *testing.T) {
|
||||
writeSkipErrorTestFile(t, fs, "/source/file1.txt", "first file content")
|
||||
writeSkipErrorTestFile(t, fs, "/source/file2.txt", "second file content")
|
||||
|
||||
repos, err := runSkipErrorScan(t, fs, true)
|
||||
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
|
||||
if err == nil {
|
||||
t.Fatal("expected scan to abort on the packer error, got nil")
|
||||
}
|
||||
@@ -184,7 +212,7 @@ func TestScannerReadErrorAbortsWithoutSkipErrors(t *testing.T) {
|
||||
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
|
||||
writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
|
||||
|
||||
_, err := runSkipErrorScan(t, fs, false)
|
||||
_, err := runSkipErrorScan(t, fs, "/source", false, nil)
|
||||
if err == nil {
|
||||
t.Fatal("expected scan to fail on the read error, got nil")
|
||||
}
|
||||
@@ -200,7 +228,7 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
|
||||
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
|
||||
writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
|
||||
|
||||
repos, err := runSkipErrorScan(t, fs, true)
|
||||
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
|
||||
if err != nil {
|
||||
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
|
||||
}
|
||||
@@ -214,3 +242,63 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
|
||||
t.Fatalf("expected unreadable file skipped, got %d chunks", len(chunks))
|
||||
}
|
||||
}
|
||||
|
||||
// writeSymlinkSource creates a source directory on disk holding one symlink
|
||||
// and returns the directory and the symlink's path.
|
||||
func writeSymlinkSource(t *testing.T) (string, string) {
|
||||
t.Helper()
|
||||
|
||||
sourceDir := t.TempDir()
|
||||
linkPath := filepath.Join(sourceDir, "link")
|
||||
|
||||
err := os.Symlink("target.txt", linkPath)
|
||||
if err != nil {
|
||||
t.Fatalf("creating symlink: %v", err)
|
||||
}
|
||||
|
||||
return sourceDir, linkPath
|
||||
}
|
||||
|
||||
// TestScannerUnreadableSymlinkAbortsWithoutSkipErrors checks that a symlink
|
||||
// whose target cannot be read aborts the run when --skip-errors is not set.
|
||||
func TestScannerUnreadableSymlinkAbortsWithoutSkipErrors(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
sourceDir, linkPath := writeSymlinkSource(t)
|
||||
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
|
||||
|
||||
_, err := runSkipErrorScan(t, fs, sourceDir, false, nil)
|
||||
if !errors.Is(err, os.ErrNotExist) {
|
||||
t.Fatalf("expected scan to fail on the removed symlink, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScannerUnreadableSymlinkSkippedWithSkipErrors checks that a symlink
|
||||
// whose target cannot be read is skipped with an error line, and the run
|
||||
// completes, when --skip-errors is set.
|
||||
func TestScannerUnreadableSymlinkSkippedWithSkipErrors(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
sourceDir, linkPath := writeSymlinkSource(t)
|
||||
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
|
||||
uiw := ui.NewWithColor(io.Discard, false)
|
||||
|
||||
repos, err := runSkipErrorScan(t, fs, sourceDir, true, uiw)
|
||||
if err != nil {
|
||||
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
|
||||
}
|
||||
|
||||
if uiw.ErrorCount() != 1 {
|
||||
t.Fatalf("expected one error line for the symlink, got %d",
|
||||
uiw.ErrorCount())
|
||||
}
|
||||
|
||||
file, err := repos.Files.GetByPath(context.Background(), linkPath)
|
||||
if err != nil {
|
||||
t.Fatalf("getting %s: %v", linkPath, err)
|
||||
}
|
||||
|
||||
if file != nil {
|
||||
t.Fatalf("expected %s not to be recorded", linkPath)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,8 +14,9 @@ import (
|
||||
|
||||
// runStorerConformance is the shared Storer contract. Every backend that
|
||||
// can run in-process is expected to pass it: TestFileStorer runs it against
|
||||
// file://, TestS3Storer against s3://. A new backend inherits this coverage
|
||||
// by passing its own constructor, so the contract is defined once.
|
||||
// file://, TestS3Storer against s3://, TestRcloneStorer against rclone's
|
||||
// local backend. A new backend inherits this coverage by passing its own
|
||||
// constructor, so the contract is defined once.
|
||||
//
|
||||
// It exercises the public Storer interface: round-trip, stat, list with
|
||||
// prefix filtering, overwrite, delete, delete-of-missing, and not-found on
|
||||
|
||||
@@ -50,9 +50,10 @@ const storageDirPerm = 0o755
|
||||
// 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.
|
||||
// blob. The rclone backend's upload does the same on remotes with a
|
||||
// server-side move. 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.
|
||||
const tempSuffix = ".partial"
|
||||
|
||||
// Put stores data at the specified key.
|
||||
|
||||
@@ -14,6 +14,9 @@ import (
|
||||
// errStreamInterrupted stands in for an upload cut off mid-stream.
|
||||
var errStreamInterrupted = errors.New("connection reset mid-upload")
|
||||
|
||||
// testBlobKey is a key laid out as a blob's key is.
|
||||
const testBlobKey = "blobs/aa/bb/aabbccddeeff"
|
||||
|
||||
// failingReader yields its data once, then fails.
|
||||
type failingReader struct {
|
||||
data []byte
|
||||
@@ -43,7 +46,7 @@ func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
key := "blobs/aa/bb/aabbccddeeff"
|
||||
key := testBlobKey
|
||||
|
||||
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
|
||||
if err == nil {
|
||||
@@ -79,7 +82,7 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
realKey := "blobs/aa/bb/aabbccddeeff"
|
||||
realKey := testBlobKey
|
||||
|
||||
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
|
||||
if err != nil {
|
||||
|
||||
+64
-15
@@ -3,6 +3,7 @@ package storage
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
@@ -68,14 +69,7 @@ func (r *RcloneStorer) Put(ctx context.Context, key string, data io.Reader) erro
|
||||
return fmt.Errorf("reading data: %w", err)
|
||||
}
|
||||
|
||||
// Upload the object
|
||||
_, err = operations.Rcat(ctx, r.fsys, key,
|
||||
io.NopCloser(bytes.NewReader(buf)), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
return r.upload(ctx, key, bytes.NewReader(buf))
|
||||
}
|
||||
|
||||
// PutWithProgress stores data with progress reporting.
|
||||
@@ -89,13 +83,7 @@ func (r *RcloneStorer) PutWithProgress(
|
||||
callback: progress,
|
||||
}
|
||||
|
||||
// Upload the object
|
||||
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(pr), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
return r.upload(ctx, key, pr)
|
||||
}
|
||||
|
||||
// Get retrieves data from the specified key.
|
||||
@@ -173,6 +161,10 @@ func (r *RcloneStorer) List(ctx context.Context, prefix string) ([]string, error
|
||||
|
||||
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
|
||||
key := obj.Remote()
|
||||
if strings.HasSuffix(key, tempSuffix) {
|
||||
return
|
||||
}
|
||||
|
||||
if prefix == "" || strings.HasPrefix(key, prefix) {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
@@ -202,6 +194,10 @@ func (r *RcloneStorer) ListStream(
|
||||
}
|
||||
|
||||
key := obj.Remote()
|
||||
if strings.HasSuffix(key, tempSuffix) {
|
||||
return
|
||||
}
|
||||
|
||||
if prefix == "" || strings.HasPrefix(key, prefix) {
|
||||
ch <- ObjectInfo{
|
||||
Key: key,
|
||||
@@ -230,6 +226,59 @@ func (r *RcloneStorer) Info() Info {
|
||||
}
|
||||
}
|
||||
|
||||
// upload writes data to key. Where the remote has a server-side move, it
|
||||
// writes under a temporary name ending in tempSuffix and moves the object
|
||||
// onto key once it is complete, so a killed upload cannot leave a truncated
|
||||
// object at key; a remote without one is written in place. List and
|
||||
// ListStream skip a temporary object left behind.
|
||||
//
|
||||
// rclone's own copy does this only where the remote also sets
|
||||
// PartialUploads. That flag is not checked here: hdfs, for one, shows a
|
||||
// file while it is written without setting it.
|
||||
func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error {
|
||||
if r.fsys.Features().Move == nil {
|
||||
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(data), time.Now(), nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
tempKey := key + "-" + rand.Text() + tempSuffix
|
||||
|
||||
obj, err := operations.Rcat(ctx, r.fsys, tempKey, io.NopCloser(data), time.Now(), nil)
|
||||
if err != nil {
|
||||
// Rcat returns the object it wrote when the written data fails its check.
|
||||
if obj != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
}
|
||||
|
||||
return fmt.Errorf("uploading object: %w", err)
|
||||
}
|
||||
|
||||
// On drive, dropbox, onedrive and others the remote's own move does not
|
||||
// replace an object already at key. operations.Move removes the object
|
||||
// it is given first, and copies where the remote refuses the move.
|
||||
existing, err := r.fsys.NewObject(ctx, key)
|
||||
if errors.Is(err, fs.ErrorObjectNotFound) {
|
||||
existing = nil
|
||||
} else if err != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
|
||||
return fmt.Errorf("looking up existing object: %w", err)
|
||||
}
|
||||
|
||||
_, err = operations.Move(ctx, r.fsys, existing, key, obj)
|
||||
if err != nil {
|
||||
_ = obj.Remove(ctx)
|
||||
|
||||
return fmt.Errorf("moving object into place: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// progressReader wraps an io.Reader to track read progress.
|
||||
type progressReader struct {
|
||||
reader io.Reader
|
||||
|
||||
+307
-14
@@ -1,28 +1,47 @@
|
||||
package storage_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/rclone/rclone/fs"
|
||||
"github.com/rclone/rclone/fs/config/configmap"
|
||||
"sneak.berlin/go/vaultik/internal/storage"
|
||||
)
|
||||
|
||||
// The rclone backend is a thin adapter over the rclone library: it turns a
|
||||
// (remote, path) pair into rclone's "remote:path" string, hands it to
|
||||
// rclone, and maps rclone's own results back to the Storer interface. What
|
||||
// can be tested in-process, without a configured remote or network, is that
|
||||
// adapter layer — how the arguments are shaped and how construction errors
|
||||
// are reported. The data-plane operations (Put/Get/List/Delete) are rclone's
|
||||
// own, exercised against a real provider (drive, s3-via-rclone, ...), which
|
||||
// needs a configured remote with credentials and network access and so is
|
||||
// out of reach of a unit test. The shared Storer conformance suite therefore
|
||||
// runs against the in-process file and s3 backends; the rclone backend
|
||||
// inherits that contract once a remote is configured.
|
||||
//
|
||||
// These tests use rclone's ":local:" on-the-fly backend, which addresses the
|
||||
// local filesystem directly without any configured remote, so construction
|
||||
// runs entirely in-process.
|
||||
// local filesystem directly without any configured remote, so they run
|
||||
// entirely in-process. A remote that needs credentials and network access
|
||||
// (drive, s3 via rclone, ...) is out of reach of a unit test.
|
||||
|
||||
// newRcloneStorer builds an rclone backend on rclone's local backend,
|
||||
// rooted at a fresh temp directory.
|
||||
//
|
||||
//nolint:ireturn // conformance runs against the Storer interface by design
|
||||
func newRcloneStorer(t *testing.T) storage.Storer {
|
||||
t.Helper()
|
||||
|
||||
s, err := storage.NewRcloneStorer(context.Background(), ":local", t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("NewRcloneStorer: %v", err)
|
||||
}
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
// TestRcloneStorer runs the shared Storer contract against the rclone
|
||||
// backend.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorer(t *testing.T) {
|
||||
runStorerConformance(t, newRcloneStorer)
|
||||
}
|
||||
|
||||
// TestNewRcloneStorerConstruction checks that a valid remote constructs a
|
||||
// backend and that Info() reports the shaped "remote:path" location.
|
||||
@@ -44,6 +63,280 @@ func TestNewRcloneStorerConstruction(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// unbufferedUploadContext returns a context in which rclone neither reads
|
||||
// ahead of the write nor holds a small upload in memory. Without it, a
|
||||
// progress callback runs before anything is written to the remote.
|
||||
func unbufferedUploadContext() context.Context {
|
||||
ctx, ci := fs.AddConfig(context.Background())
|
||||
ci.BufferSize = 0
|
||||
ci.StreamingUploadCutoff = 0
|
||||
|
||||
return ctx
|
||||
}
|
||||
|
||||
// objectAtKeyDuringUpload uploads data to testBlobKey and reports whether
|
||||
// an object was at the key before the upload finished. It fails the test
|
||||
// unless the key then reads back as the uploaded data.
|
||||
func objectAtKeyDuringUpload(
|
||||
ctx context.Context, t *testing.T, s *storage.RcloneStorer,
|
||||
) bool {
|
||||
t.Helper()
|
||||
|
||||
data := bytes.Repeat([]byte("blob-bytes"), 1000)
|
||||
seen := false
|
||||
|
||||
err := s.PutWithProgress(ctx, testBlobKey, bytes.NewReader(data),
|
||||
int64(len(data)), func(int64) error {
|
||||
_, statErr := s.Stat(ctx, testBlobKey)
|
||||
if statErr == nil {
|
||||
seen = true
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("PutWithProgress: %v", err)
|
||||
}
|
||||
|
||||
got := readObject(ctx, t, s, testBlobKey)
|
||||
if !bytes.Equal(got, data) {
|
||||
t.Errorf("object read back as %d bytes, want the %d uploaded",
|
||||
len(got), len(data))
|
||||
}
|
||||
|
||||
return seen
|
||||
}
|
||||
|
||||
// readObject returns the contents of the object at key.
|
||||
func readObject(
|
||||
ctx context.Context, t *testing.T, s *storage.RcloneStorer, key string,
|
||||
) []byte {
|
||||
t.Helper()
|
||||
|
||||
rc, err := s.Get(ctx, key)
|
||||
if err != nil {
|
||||
t.Fatalf("Get: %v", err)
|
||||
}
|
||||
|
||||
defer func() { _ = rc.Close() }()
|
||||
|
||||
got, err := io.ReadAll(rc)
|
||||
if err != nil {
|
||||
t.Fatalf("reading object: %v", err)
|
||||
}
|
||||
|
||||
return got
|
||||
}
|
||||
|
||||
// TestRcloneStorerObjectAppearsOnlyWhenComplete checks that on a remote
|
||||
// with a server-side move, such as local, nothing is at the key until the
|
||||
// upload has finished, so a killed upload cannot leave a truncated object
|
||||
// there.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerObjectAppearsOnlyWhenComplete(t *testing.T) {
|
||||
ctx := unbufferedUploadContext()
|
||||
|
||||
s, err := storage.NewRcloneStorer(ctx, ":local", t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("NewRcloneStorer: %v", err)
|
||||
}
|
||||
|
||||
if objectAtKeyDuringUpload(ctx, t, s) {
|
||||
t.Error("object was at its key before the upload finished")
|
||||
}
|
||||
}
|
||||
|
||||
// TestRcloneStorerListSkipsPartialFiles checks that a temporary file left
|
||||
// by a killed upload is never listed as a key.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerListSkipsPartialFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
ctx := context.Background()
|
||||
|
||||
s, err := storage.NewRcloneStorer(ctx, ":local", dir)
|
||||
if err != nil {
|
||||
t.Fatalf("NewRcloneStorer: %v", err)
|
||||
}
|
||||
|
||||
realKey := testBlobKey
|
||||
|
||||
err = s.Put(ctx, realKey, strings.NewReader("blob-bytes"))
|
||||
if err != nil {
|
||||
t.Fatalf("Put: %v", err)
|
||||
}
|
||||
|
||||
leftover := filepath.Join(dir, realKey+"-123456.partial")
|
||||
|
||||
err = os.WriteFile(leftover, []byte("half"), 0o600)
|
||||
if err != nil {
|
||||
t.Fatalf("writing leftover temp file: %v", err)
|
||||
}
|
||||
|
||||
keys, err := s.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 s.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)
|
||||
}
|
||||
}
|
||||
|
||||
// newRcloneStorerOnWrappedLocal registers name as rclone's local backend
|
||||
// wrapped by wrap, and builds an rclone backend on it rooted at a fresh
|
||||
// temp directory. wrap changes the features the local backend reports, so
|
||||
// that it behaves like a remote a unit test cannot reach.
|
||||
func newRcloneStorerOnWrappedLocal(
|
||||
ctx context.Context, t *testing.T, name string, wrap func(fs.Fs) fs.Fs,
|
||||
) *storage.RcloneStorer {
|
||||
t.Helper()
|
||||
|
||||
fs.Register(&fs.RegInfo{
|
||||
Name: name,
|
||||
NewFs: func(
|
||||
ctx context.Context, _, root string, _ configmap.Mapper,
|
||||
) (fs.Fs, error) {
|
||||
local, err := fs.NewFs(ctx, ":local:"+root)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return wrap(local), nil
|
||||
},
|
||||
})
|
||||
|
||||
s, err := storage.NewRcloneStorer(ctx, ":"+name, t.TempDir())
|
||||
if err != nil {
|
||||
t.Fatalf("NewRcloneStorer: %v", err)
|
||||
}
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
// withoutPartialUploads is rclone's local backend with the PartialUploads
|
||||
// flag cleared. Like hdfs, it then has a server-side move and shows a file
|
||||
// while it is written, without setting that flag.
|
||||
type withoutPartialUploads struct {
|
||||
fs.Fs
|
||||
}
|
||||
|
||||
func (f *withoutPartialUploads) Features() *fs.Features {
|
||||
features := *f.Fs.Features()
|
||||
features.PartialUploads = false
|
||||
|
||||
return &features
|
||||
}
|
||||
|
||||
// TestRcloneStorerMovesIntoPlaceWithoutPartialUploads checks that on a
|
||||
// remote with a server-side move nothing is at the key until the upload has
|
||||
// finished, even when rclone does not mark the remote as showing partial
|
||||
// uploads.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerMovesIntoPlaceWithoutPartialUploads(t *testing.T) {
|
||||
ctx := unbufferedUploadContext()
|
||||
s := newRcloneStorerOnWrappedLocal(ctx, t, "withoutpartialuploads",
|
||||
func(local fs.Fs) fs.Fs { return &withoutPartialUploads{Fs: local} })
|
||||
|
||||
if objectAtKeyDuringUpload(ctx, t, s) {
|
||||
t.Error("object was at its key before the upload finished")
|
||||
}
|
||||
}
|
||||
|
||||
// withoutMove is rclone's local backend reporting no server-side move.
|
||||
type withoutMove struct {
|
||||
fs.Fs
|
||||
}
|
||||
|
||||
func (f *withoutMove) Features() *fs.Features {
|
||||
features := *f.Fs.Features()
|
||||
features.Move = nil
|
||||
|
||||
return &features
|
||||
}
|
||||
|
||||
// TestRcloneStorerWritesInPlaceWithoutMove checks that on a remote with no
|
||||
// server-side move an object is written straight to its key and reads back
|
||||
// from there.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerWritesInPlaceWithoutMove(t *testing.T) {
|
||||
ctx := unbufferedUploadContext()
|
||||
s := newRcloneStorerOnWrappedLocal(ctx, t, "withoutmove",
|
||||
func(local fs.Fs) fs.Fs { return &withoutMove{Fs: local} })
|
||||
|
||||
if !objectAtKeyDuringUpload(ctx, t, s) {
|
||||
t.Error("object was not at its key while it was uploaded")
|
||||
}
|
||||
}
|
||||
|
||||
var errNameConflict = errors.New("an object with this name already exists")
|
||||
|
||||
// moveRefusesExisting is rclone's local backend with a server-side move
|
||||
// that, like dropbox's or onedrive's, refuses to move onto an existing
|
||||
// object.
|
||||
type moveRefusesExisting struct {
|
||||
fs.Fs
|
||||
}
|
||||
|
||||
func (f *moveRefusesExisting) Features() *fs.Features {
|
||||
features := *f.Fs.Features()
|
||||
features.Move = f.move
|
||||
|
||||
return &features
|
||||
}
|
||||
|
||||
//nolint:ireturn // the signature is rclone's
|
||||
func (f *moveRefusesExisting) move(
|
||||
ctx context.Context, src fs.Object, remote string,
|
||||
) (fs.Object, error) {
|
||||
_, err := f.NewObject(ctx, remote)
|
||||
if err == nil {
|
||||
return nil, errNameConflict
|
||||
}
|
||||
|
||||
return f.Fs.Features().Move(ctx, src, remote)
|
||||
}
|
||||
|
||||
// TestRcloneStorerOverwritesWhereMoveRefusesExisting checks that writing a
|
||||
// key twice replaces the object on a remote whose server-side move will not
|
||||
// replace an existing object.
|
||||
//
|
||||
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
|
||||
func TestRcloneStorerOverwritesWhereMoveRefusesExisting(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
s := newRcloneStorerOnWrappedLocal(ctx, t, "moverefusesexisting",
|
||||
func(local fs.Fs) fs.Fs { return &moveRefusesExisting{Fs: local} })
|
||||
|
||||
for _, content := range []string{"first", "second"} {
|
||||
err := s.Put(ctx, testBlobKey, strings.NewReader(content))
|
||||
if err != nil {
|
||||
t.Fatalf("Put %q: %v", content, err)
|
||||
}
|
||||
}
|
||||
|
||||
got := readObject(ctx, t, s, testBlobKey)
|
||||
if string(got) != "second" {
|
||||
t.Errorf("object = %q, want %q", got, "second")
|
||||
}
|
||||
}
|
||||
|
||||
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
|
||||
// rclone config fails construction with the ErrRemoteNotFound sentinel,
|
||||
// rather than silently returning a backend pointed nowhere.
|
||||
|
||||
@@ -155,8 +155,8 @@ type BlobHash string
|
||||
// FilePath represents an absolute path to a file or directory.
|
||||
type FilePath string
|
||||
|
||||
// SourcePath represents the root directory from which files are backed up.
|
||||
// Used during restore to strip the source prefix from paths.
|
||||
// SourcePath is the source directory a scan found a file under, made
|
||||
// absolute and with symlinks resolved.
|
||||
type SourcePath string
|
||||
|
||||
// Hostname identifies a host machine.
|
||||
|
||||
@@ -130,10 +130,9 @@ func assertThirdSnapshotRestores(
|
||||
// up, and that snapshot is removed. The first snapshot keeps the file row,
|
||||
// which now lists the appended content's chunks, while removal drops the
|
||||
// blob that held them.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupAfterRemovingNewestSnapshotRestoresChangedFile(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -178,10 +177,9 @@ func TestBackupAfterRemovingNewestSnapshotRestoresChangedFile(t *testing.T) {
|
||||
// The next run's prune drops that incomplete snapshot and its blob, while
|
||||
// the first snapshot keeps the file row, which now lists the appended
|
||||
// content's chunks.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupAfterInterruptedRunRestoresChangedFile(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
@@ -38,11 +38,9 @@ import (
|
||||
// (https://git.eeqj.de/sneak/vaultik/issues/130) and is not re-tested
|
||||
// here; these tests target the layers above the backend.
|
||||
//
|
||||
// The tests run serially, not with t.Parallel: each calls
|
||||
// log.Initialize, which replaces the package-global logger, and a
|
||||
// backup or restore running concurrently reads that same logger. Under
|
||||
// -race the two collide. Running one at a time is the same choice
|
||||
// prune_count_test.go already makes for the same reason.
|
||||
// log.Initialize replaces the package-global logger that a running
|
||||
// backup or restore reads, so each test calls it before t.Parallel,
|
||||
// while no parallel test is running yet.
|
||||
|
||||
const (
|
||||
faultChunkSize = int64(64 * 1024)
|
||||
@@ -165,17 +163,19 @@ func newReaderVaultik(
|
||||
// Scenario 3: a stored blob's bytes are flipped before restore reads
|
||||
// them. Restore must fail loudly, and no file must be left on the
|
||||
// restore target holding corrupt content.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreRejectsCorruptBlob(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
assertRestoreRejectsDamagedBlob(t, faultstore.GetCorrupt, "corrupt")
|
||||
}
|
||||
|
||||
// Scenario 4: a stored blob is truncated before restore reads it. Same
|
||||
// contract as the corrupt case.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreRejectsTruncatedBlob(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
assertRestoreRejectsDamagedBlob(t, faultstore.GetTruncate, "truncated")
|
||||
}
|
||||
|
||||
@@ -188,7 +188,6 @@ func assertRestoreRejectsDamagedBlob(
|
||||
t *testing.T, fault faultstore.GetFault, name string,
|
||||
) {
|
||||
t.Helper()
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -232,10 +231,9 @@ func assertRestoreRejectsDamagedBlob(
|
||||
|
||||
// Scenario 6: the backend accepts blob uploads and reports success but
|
||||
// stores nothing. verify --deep must catch it.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestDeepVerifyCatchesLyingBackend(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -285,10 +283,9 @@ func TestDeepVerifyCatchesLyingBackend(t *testing.T) {
|
||||
// Scenario 1a: a blob upload fails partway through. The interrupted run
|
||||
// must not record the blob as uploaded, must not reference it from the
|
||||
// snapshot, and must leave no blob object at the destination.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestInterruptedBlobUploadRecordsNoUploadedBlob(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -361,10 +358,9 @@ func TestInterruptedBlobUploadRecordsNoUploadedBlob(t *testing.T) {
|
||||
// chunks in a blob that was actually uploaded, so the retry re-chunks and
|
||||
// re-uploads the affected data instead of silently referencing data that
|
||||
// never reached storage.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -426,10 +422,9 @@ func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) {
|
||||
// covered by TestBackupCompletesOnlyAfterMetadataExport
|
||||
// (https://git.eeqj.de/sneak/vaultik/issues/177); this test exercises the
|
||||
// lower-level export path in isolation.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupSurvivesMetadataExportInterruption(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -509,10 +504,9 @@ func TestBackupSurvivesMetadataExportInterruption(t *testing.T) {
|
||||
// destination. Rerunning the backup must then prune the incomplete
|
||||
// snapshot, produce a snapshot whose destination metadata and local index
|
||||
// agree, and restore. See https://git.eeqj.de/sneak/vaultik/issues/177.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupCompletesOnlyAfterMetadataExport(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
@@ -685,10 +679,9 @@ func faultScannerFactory(
|
||||
// Scenario 5: the restore target runs out of space mid-file. Restore
|
||||
// must fail with an out-of-space error, and must not leave a truncated
|
||||
// file at the target path presenting as a complete restore.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreReportsDiskFull(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
osFS := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
@@ -928,17 +928,17 @@ func setupDedupBackupEnv(
|
||||
}
|
||||
}
|
||||
|
||||
// runDedupSnapshot creates a "dedup" snapshot, scans dataDir into it,
|
||||
// completes it, and exports its metadata, returning the snapshot ID and
|
||||
// scan result.
|
||||
// runDedupSnapshot creates a snapshot with the given name, scans dataDir
|
||||
// into it, completes it, and exports its metadata, returning the snapshot
|
||||
// ID and scan result.
|
||||
func runDedupSnapshot(
|
||||
ctx context.Context, t *testing.T,
|
||||
sm *snapshot.SnapshotManager, scanner *snapshot.Scanner,
|
||||
hostname, dataDir, dbPath string,
|
||||
hostname, name, dataDir, dbPath string,
|
||||
) (string, *snapshot.ScanResult) {
|
||||
t.Helper()
|
||||
|
||||
id, err := sm.CreateSnapshotWithName(ctx, hostname, "dedup", "v", "g")
|
||||
id, err := sm.CreateSnapshotWithName(ctx, hostname, name, "v", "g")
|
||||
require.NoError(t, err)
|
||||
|
||||
result, err := scanner.Scan(ctx, dataDir, id)
|
||||
@@ -980,16 +980,15 @@ func TestDedupOnlySnapshotRestores(t *testing.T) {
|
||||
|
||||
// First snapshot — uploads all blobs.
|
||||
_, r1 := runDedupSnapshot(ctx, t, sm, makeScanner(),
|
||||
cfg.Hostname, dataDir, dbPath)
|
||||
cfg.Hostname, "first", dataDir, dbPath)
|
||||
require.Positive(t, r1.BlobsCreated,
|
||||
"first snapshot should upload at least one blob")
|
||||
|
||||
// Second snapshot — same data, every chunk dedups. Sleep past the
|
||||
// second-precision timestamp so the snapshot IDs differ.
|
||||
time.Sleep(1100 * time.Millisecond)
|
||||
|
||||
// Second snapshot — same data, every chunk dedups. Its own name gives
|
||||
// it a different snapshot ID without waiting for the one-second
|
||||
// timestamp in the ID to tick over.
|
||||
id2, r2 := runDedupSnapshot(ctx, t, sm, makeScanner(),
|
||||
cfg.Hostname, dataDir, dbPath)
|
||||
cfg.Hostname, "second", dataDir, dbPath)
|
||||
require.Equal(t, 0, r2.BlobsCreated,
|
||||
"second snapshot should upload zero new blobs (fully dedup'd)")
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"io/fs"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/spf13/afero"
|
||||
@@ -77,10 +78,9 @@ func backUpThenUnplug(
|
||||
// TestFirstBackupCreatesDestinationDirectory checks that a first backup
|
||||
// to a destination directory that does not exist yet creates it, and
|
||||
// that the destination can be listed afterwards.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestFirstBackupCreatesDestinationDirectory(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
storeDir := filepath.Join(t.TempDir(), "volume", "backup")
|
||||
@@ -95,10 +95,9 @@ func TestFirstBackupCreatesDestinationDirectory(t *testing.T) {
|
||||
// TestListSnapshotsWarnsWhenDestinationMissing checks that snapshot list
|
||||
// warns and shows the local index alone, without reporting the local
|
||||
// snapshot as missing from the destination.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestListSnapshotsWarnsWhenDestinationMissing(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
v, repos, out := backUpThenUnplug(ctx, t)
|
||||
@@ -115,10 +114,9 @@ func TestListSnapshotsWarnsWhenDestinationMissing(t *testing.T) {
|
||||
// TestRemoveSnapshotWarnsWhenDestinationMissing checks that snapshot
|
||||
// remove warns that the metadata could not be removed from the
|
||||
// destination, instead of reporting that it was.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRemoveSnapshotWarnsWhenDestinationMissing(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
v, repos, out := backUpThenUnplug(ctx, t)
|
||||
@@ -137,10 +135,9 @@ func TestRemoveSnapshotWarnsWhenDestinationMissing(t *testing.T) {
|
||||
// TestPruneKeepsLocalRecordsWhenDestinationMissing checks that prune
|
||||
// fails on a destination it cannot list and deletes no local snapshot
|
||||
// record.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestPruneKeepsLocalRecordsWhenDestinationMissing(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
v, repos, _ := backUpThenUnplug(ctx, t)
|
||||
@@ -153,3 +150,22 @@ func TestPruneKeepsLocalRecordsWhenDestinationMissing(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, snapshots, 1, "prune must delete no local snapshot record")
|
||||
}
|
||||
|
||||
// TestPurgeSaysListingFailedOnceWhenDestinationMissing checks that
|
||||
// snapshot purge fails on a destination it cannot list, with an error
|
||||
// that says "listing remote snapshots" once.
|
||||
func TestPurgeSaysListingFailedOnceWhenDestinationMissing(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
v, _, _ := backUpThenUnplug(ctx, t)
|
||||
|
||||
err := v.PurgeSnapshotsWithOptions(&vaultik.SnapshotPurgeOptions{
|
||||
KeepLatest: true,
|
||||
Force: true,
|
||||
})
|
||||
require.ErrorIs(t, err, fs.ErrNotExist)
|
||||
assert.Equal(t, 1, strings.Count(err.Error(), "listing remote snapshots"),
|
||||
err.Error())
|
||||
}
|
||||
|
||||
@@ -14,10 +14,9 @@ import (
|
||||
// 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{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
|
||||
@@ -164,10 +164,9 @@ func scratchEntries(t *testing.T, dir string) []string {
|
||||
// restore while a blob download is in progress. The download fails only
|
||||
// because of the cancel, so Restore must return context.Canceled without
|
||||
// reporting the file that needs the blob as failed.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreSkipErrorsCancelDuringBlobDownload(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
@@ -45,9 +45,10 @@ type missingBlobBackup struct {
|
||||
// after one blob of a two-blob snapshot was deleted. Every file stored in
|
||||
// that blob must be reported as failed and left absent, every other file
|
||||
// must be restored intact, and Restore must still return an error.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreSkipErrorsSkipsFilesOfMissingBlob(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
backup := backupThenDeleteOneBlob(ctx, t)
|
||||
|
||||
@@ -85,9 +86,10 @@ func TestRestoreSkipErrorsSkipsFilesOfMissingBlob(t *testing.T) {
|
||||
|
||||
// TestRestoreMissingBlobAbortsWithoutSkipErrors checks that a deleted blob
|
||||
// still ends the restore with an error when SkipErrors is not set.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestRestoreMissingBlobAbortsWithoutSkipErrors(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
backup := backupThenDeleteOneBlob(ctx, t)
|
||||
|
||||
@@ -107,7 +109,6 @@ func backupThenDeleteOneBlob(
|
||||
ctx context.Context, t *testing.T,
|
||||
) *missingBlobBackup {
|
||||
t.Helper()
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
@@ -18,10 +18,9 @@ import (
|
||||
// A file rewritten with its size unchanged and a new mtime in the same
|
||||
// second as the mtime the index holds must still be backed up. See
|
||||
// https://git.eeqj.de/sneak/vaultik/issues/226.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupOfSameSecondRewriteRestoresNewContent(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
|
||||
@@ -1052,7 +1052,7 @@ func (v *Vaultik) syncWithRemote() error {
|
||||
// every local snapshot record (issue #160).
|
||||
remoteKeys, err := v.listAllRemoteSnapshotKeys()
|
||||
if err != nil {
|
||||
return fmt.Errorf("listing remote snapshots: %w", err)
|
||||
return err
|
||||
}
|
||||
|
||||
remoteKeySet := make(map[string]bool, len(remoteKeys))
|
||||
|
||||
@@ -54,8 +54,6 @@ type summaryEnv struct {
|
||||
func newSummaryEnv(t *testing.T) *summaryEnv {
|
||||
t.Helper()
|
||||
|
||||
log.Initialize(log.Config{})
|
||||
|
||||
fs := afero.NewOsFs()
|
||||
tempDir := t.TempDir()
|
||||
srcDir := filepath.Join(tempDir, "src")
|
||||
@@ -193,9 +191,10 @@ func (e *summaryEnv) dataLine(total, backedUp int64) string {
|
||||
// A first backup stores copy.bin's chunks while backing up a.bin, so
|
||||
// copy.bin's chunks are deduplicated within the run. Each file and byte
|
||||
// is still counted once.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryFirstRun(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
summary := env.backUp(t, "first", false)
|
||||
@@ -219,9 +218,10 @@ func TestSnapshotSummaryFirstRun(t *testing.T) {
|
||||
// An incremental backup where a.bin's mtime changed but its content did
|
||||
// not: a.bin is backed up again and every one of its chunks is already
|
||||
// stored.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
env.backUp(t, "first", false)
|
||||
@@ -251,9 +251,10 @@ func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) {
|
||||
// Under --cron the progress reporter is off; the upload figures must
|
||||
// still reach the summary and the snapshots row. The snapshot has two
|
||||
// paths, each backed up by its own scan.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestSnapshotSummaryCronRunRecordsUploads(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
env := newSummaryEnv(t)
|
||||
|
||||
summary := env.backUp(t, "split", true)
|
||||
|
||||
@@ -18,10 +18,9 @@ import (
|
||||
// A backup without --cron runs the progress reporter while one scanner
|
||||
// scans each path of the snapshot in turn. See
|
||||
// https://git.eeqj.de/sneak/vaultik/issues/253.
|
||||
//
|
||||
//nolint:paralleltest // installs the global logger via log.Initialize
|
||||
func TestBackupWithoutCronOfTwoPathSnapshotRestoresBothPaths(t *testing.T) {
|
||||
log.Initialize(log.Config{})
|
||||
t.Parallel()
|
||||
|
||||
const snapshotName = "data"
|
||||
|
||||
|
||||
+5
-2
@@ -1,6 +1,9 @@
|
||||
#!/bin/sh
|
||||
# script/fmt-check: check formatting (read-only). Same scope as
|
||||
# script/fmt, but fails instead of writing.
|
||||
# script/fmt-check: check formatting (read-only). Fails instead of
|
||||
# writing. It checks every Go file outside .tool, which is more than
|
||||
# script/fmt formats: `go fmt ./...` skips `testdata` directories and
|
||||
# files and directories whose names start with `.` or `_`. Fix a file
|
||||
# only this reports with `gofmt -w`.
|
||||
set -eu
|
||||
|
||||
ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
||||
|
||||
+3
-4
@@ -21,10 +21,9 @@ goreleaser_version() {
|
||||
head -n 1
|
||||
}
|
||||
|
||||
# Resolve the goreleaser to run, on the same rule script/lint uses for
|
||||
# golangci-lint: a binary on PATH is accepted only when it is exactly
|
||||
# the pinned version, because a differently versioned tool would
|
||||
# produce a differently built release from the same tag. Anything else
|
||||
# Resolve the goreleaser to run. A binary on PATH is accepted only when
|
||||
# it is exactly the pinned version, because a differently versioned tool
|
||||
# would produce a differently built release from the same tag. Anything else
|
||||
# comes from .tool/bin, and a missing one is a loud failure naming the
|
||||
# script that installs it rather than a silent fallback.
|
||||
resolve_goreleaser() {
|
||||
|
||||
+1
-1
@@ -19,7 +19,7 @@ s3:
|
||||
secret_access_key: test-secret-key
|
||||
region: us-east-1
|
||||
use_ssl: true
|
||||
part_size: 5242880 # 5MB
|
||||
part_size: 5242880 # 5MiB
|
||||
index_path: /tmp/vaultik-test.sqlite
|
||||
chunk_size: 10MB
|
||||
blob_size_limit: 10GB
|
||||
|
||||
Reference in New Issue
Block a user