Compare commits

..
1 Commits
Author SHA1 Message Date
sneak c2ce69d258 Correct doc and help sentences that are false about the code (closes #233)
check / check (push) Canceled after 0s
A blob is written whole to a temporary file and uploaded once finished,
not streamed to storage. The README, ARCHITECTURE.md and
config.example.yml now say a backup needs free temporary space for each
blob (twice that for an rclone destination that cannot stream uploads)
and for the metadata export's copies of the local index, in $TMPDIR, or
partly in /var/tmp when TMPDIR is unset. 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 default, the
config search order in the snapshot create help, what snapshot remove
cleans up in the prune help, how the release installs Go, and the
script/release and script/fmt-check comments.

Model: opus-5-5
2026-10-07 15:07:02 +00:00
44 changed files with 267 additions and 1990 deletions
+3 -16
View File
@@ -359,8 +359,7 @@ on the destination in one go, use `vaultik remote nuke --force`.
* `--local-only`: Skip remote cleanup; only touch the local index * `--local-only`: Skip remote cleanup; only touch the local index
* `--dry-run`: Show what would be deleted without deleting * `--dry-run`: Show what would be deleted without deleting
* `--force`: Skip confirmation prompt * `--force`: Skip confirmation prompt
* `--json`: Output result as JSON. Also skips the confirmation prompt, as * `--json`: Output result as JSON
`--force` does.
**`snapshot restore`**: Restore files from a backup snapshot. **`snapshot restore`**: Restore files from a backup snapshot.
* Requires `VAULTIK_AGE_SECRET_KEY` environment variable * Requires `VAULTIK_AGE_SECRET_KEY` environment variable
@@ -384,8 +383,7 @@ manifests — network cost scales with the number of snapshots. `snapshot
create --prune` runs the same cleanup automatically; this is the create --prune` runs the same cleanup automatically; this is the
manual entry point for the same work. manual entry point for the same work.
* `--force`: Skip confirmation prompt * `--force`: Skip confirmation prompt
* `--json`: Output stats as JSON. Also skips the confirmation prompt, as * `--json`: Output stats as JSON
`--force` does.
**`info`**: Display system configuration, storage settings, encryption **`info`**: Display system configuration, storage settings, encryption
recipients, and local database statistics. recipients, and local database statistics.
@@ -398,8 +396,7 @@ key is skipped with a warning and is not printed. If a listed
orphaned blob figures are reported as unknown; `--json` gives them as orphaned blob figures are reported as unknown; `--json` gives them as
`null`, lists the remote key of each unreadable manifest in `null`, lists the remote key of each unreadable manifest in
`unreadable_manifests` and counts the manifests under skipped names in `unreadable_manifests` and counts the manifests under skipped names in
`skipped_manifest_count`. An unreadable manifest also leaves its `skipped_manifest_count`.
snapshot's blob count and blob size unknown, `null` in `--json`.
* `--json`: Output as JSON * `--json`: Output as JSON
**`remote nuke`**: Delete every snapshot's metadata and every blob from the **`remote nuke`**: Delete every snapshot's metadata and every blob from the
@@ -431,16 +428,6 @@ a local or mounted filesystem. Useful for testing or backing up to a NAS.
**Rclone** (`rclone://remote/path`): Uses rclone's 70+ supported cloud **Rclone** (`rclone://remote/path`): Uses rclone's 70+ supported cloud
providers. Requires rclone to be configured separately (`rclone config`). 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 Legacy S3 configuration via `s3.*` fields (endpoint, bucket, prefix, etc.) is
still supported for backward compatibility. `storage_url` takes precedence if still supported for backward compatibility. `storage_url` takes precedence if
both are set. both are set.
-106
View File
@@ -22,102 +22,6 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps # Completed Steps
- 2026-10-08: Made a backup cancelled under `--skip-errors` stop at once
([issue #286](https://git.eeqj.de/sneak/vaultik/issues/286)). After
Ctrl-C or SIGTERM, phase 2 treated the cancellation error from each
remaining file like an unreadable file: it opened the file, printed an
error line for it, counted it as failed and went on to the next one.
The run now stops at the first file after the cancellation, and no
file is counted as failed because of it.
- 2026-10-08: Made `remote nuke` delete the `.partial` files that
uploads cut off part-way leave on the destination store
([issue #281](https://git.eeqj.de/sneak/vaultik/issues/281)). The
file and rclone listings skip such a file, so the command left it in
place and still reported the store empty. It now removes them under
`metadata/` and `blobs/` as its last step.
- 2026-10-08: Made a local index error while recording a directory or
symlink stop a backup under `--skip-errors`
([issue #284](https://git.eeqj.de/sneak/vaultik/issues/284)). Phase 2
only records such an entry in the local index, since it has no data
to open or read, but an error doing so was skipped like an unreadable
file. The snapshot then completed without the entry, and a restore
did not recreate it. The error now stops the backup, as it does
without the flag.
- 2026-10-08: Counted a file that a backup could not store as failed
([issue #280](https://git.eeqj.de/sneak/vaultik/issues/280)). A file
that phase 1 counted and phase 2 could not open, because it was
unreadable under `--skip-errors` or removed in between, was reported
in the summary as unchanged with its bytes as backed up, and the
`snapshots` row's `file_count` and `total_size` included it. The
summary now counts it as failed, and its data total and the row leave
it out.
- 2026-10-08: Made `remote info` report a snapshot's blob count and
blob size as unknown when its manifest cannot be read
([issue #272](https://git.eeqj.de/sneak/vaultik/issues/272)). The
orphan figures were already unknown in that case, but the snapshot's
row still gave 0 blobs and 0 B, in the table and in `--json`. The row
now reads `unknown` and `--json` gives `null`. A directory with no
manifest still shows 0, since the orphan figures count its blobs as
orphaned.
- 2026-10-08: Kept the progress line of a snapshot with more than one
path within 100%
([issue #271](https://git.eeqj.de/sneak/vaultik/issues/271)). The
bytes and files processed counted every path, but the totals they
were divided by held only the current path's, so a second path
smaller than the first showed more than 100% and an ETA of `unknown`.
The totals now add up over the paths scanned so far. The rate behind
the ETA restarts when each path's scan phase ends, so that phase does
not lower the rate while the path is processed; until then the rate
is the previous path's.
- 2026-10-08: Made a second `snapshot create` of one name succeed when
it starts in the same second as the first
([issue #270](https://git.eeqj.de/sneak/vaultik/issues/270)). The
timestamp in a snapshot ID is in whole seconds, so the second run got
the first run's ID and failed with `UNIQUE constraint failed:
snapshots.id`. When the local index already has a snapshot with the
ID, the create now waits a second and takes a new timestamp.
- 2026-10-07: Documented that `--json` skips the confirmation prompt of
`snapshot remove` and `prune`
([issue #268](https://git.eeqj.de/sneak/vaultik/issues/268)). Both
commands delete without asking under `--json`, since a prompt on stdout
would break the JSON document, but the help and the README described
only `--force` as skipping it. The `--json` help of both commands and
their README entries now say so.
- 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 - 2026-10-07: Corrected documentation, help text and comments that were
false about the code false about the code
([issue #233](https://git.eeqj.de/sneak/vaultik/issues/233)). A blob ([issue #233](https://git.eeqj.de/sneak/vaultik/issues/233)). A blob
@@ -134,16 +38,6 @@ the tag exists and is exercised; what is left is merging `next` to
what `snapshot remove` cleans up, how the release gets its Go what `snapshot remove` cleans up, how the release gets its Go
toolchain, and the `script/release` and `script/fmt-check` comments. 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 - 2026-10-07: Made two messages say only what is true
([issue #240](https://git.eeqj.de/sneak/vaultik/issues/240)). A config ([issue #240](https://git.eeqj.de/sneak/vaultik/issues/240)). A config
file that others can read was warned about as containing S3 file that others can read was warned about as containing S3
+14 -39
View File
@@ -206,10 +206,6 @@ func RunApp(ctx context.Context, app *fx.App) error {
// RunOperation through cobra to Entry. // RunOperation through cobra to Entry.
var errReported = errors.New("operation failed") 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 // 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 // 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 // within the goroutine. An os.Exit there skipped main's deferred
@@ -224,21 +220,17 @@ var errInterrupted = errors.New("interrupted before the command finished")
// interrupt OnStop cancels op and waits for the goroutine to return, so // interrupt OnStop cancels op and waits for the goroutine to return, so
// op's cleanup (removing decrypted scratch files) runs before the // op's cleanup (removing decrypted scratch files) runs before the
// process exits; the wait is bounded by shutdownTimeout. report is // process exits; the wait is bounded by shutdownTimeout. report is
// called with a failure so the caller can show it to the user before // called with a non-canceled failure so the caller can show it to the
// it becomes errReported. // user before it becomes errReported. A context cancellation is the
// // interrupt path, not a failure: it is neither reported nor counted as
// The run counts as interrupted unless op returned, without an // one.
// interrupt having cancelled it, before RunWithApp returned. An
// interrupted op is not reported, whatever it returned; RunOperation
// returns errInterrupted instead.
func RunOperation( func RunOperation(
ctx context.Context, opts AppOptions, ctx context.Context, opts AppOptions,
op func(v *vaultik.Vaultik) error, report func(err error), op func(v *vaultik.Vaultik) error, report func(err error),
) error { ) error {
var ( var (
mu sync.Mutex mu sync.Mutex
finished bool // op returned before any interrupt cancelled it failed bool
failed bool // op finished with an error
) )
opts.Invokes = append(opts.Invokes, opts.Invokes = append(opts.Invokes,
@@ -249,21 +241,11 @@ func RunOperation(
OnStart: func(_ context.Context) error { OnStart: func(_ context.Context) error {
stop = v.StartOperation(func() { stop = v.StartOperation(func() {
err := op(v) err := op(v)
if err != nil && !errors.Is(err, context.Canceled) {
// Only stop, called from OnStop below, cancels the report(err)
// 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() mu.Lock()
finished = true failed = true
failed = err != nil
mu.Unlock() mu.Unlock()
} }
@@ -296,23 +278,16 @@ func RunOperation(
return err return err
} }
// RunWithApp returns only after the app was asked to stop, either by // The goroutine sets failed before triggering the shutdown that lets
// an interrupt or by the goroutine's Shutdown call. When op finished // RunWithApp return, so the write is in place by the time we read it.
// 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() mu.Lock()
defer mu.Unlock() defer mu.Unlock()
switch { if failed {
case !finished:
return errInterrupted
case failed:
return errReported return errReported
default:
return nil
} }
return nil
} }
// runVaultikApp runs the standard single-operation command lifecycle // runVaultikApp runs the standard single-operation command lifecycle
+4 -15
View File
@@ -15,11 +15,6 @@ import (
// the startup banner. // the startup banner.
const shortCommitLen = 12 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. // Entry is the main entry point for the CLI application.
// It prints the startup banner to stderr (unless a banner-suppressing // It prints the startup banner to stderr (unless a banner-suppressing
// flag is present in os.Args — see bannerSuppressedInArgs), executes the // flag is present in os.Args — see bannerSuppressedInArgs), executes the
@@ -28,10 +23,9 @@ const exitCodeInterrupted = 130
// The banner goes to stderr because stdout carries only the output the // 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. // user asked for, such as a completion script or a `config get` value.
// //
// It returns the process exit code (0 on success, 130 when interrupted, // It returns the process exit code (0 on success, 1 on error) rather
// 1 on any other error) rather than calling os.Exit, so that main's // than calling os.Exit, so that main's deferred profile writers run
// deferred profile writers run before the process ends. See run in // before the process ends. See run in cmd/vaultik/main.go.
// cmd/vaultik/main.go.
func Entry() int { func Entry() int {
emitStartupBanner(os.Args[1:], os.Stderr) emitStartupBanner(os.Args[1:], os.Stderr)
@@ -45,16 +39,11 @@ func Entry() int {
// document instead); errReported says so. Printing it again // document instead); errReported says so. Printing it again
// here would double the error line. // here would double the error line.
// Every other error — bad arguments, a config that would not // Every other error — bad arguments, a config that would not
// load, an interrupt — reaches Entry unreported, so it is shown // load — reaches Entry unreported, so it is shown here.
// here.
if !errors.Is(err, errReported) { if !errors.Is(err, errReported) {
ReportErrorf("%s", err.Error()) ReportErrorf("%s", err.Error())
} }
if errors.Is(err, errInterrupted) {
return exitCodeInterrupted
}
return 1 return 1
} }
-203
View File
@@ -1,203 +0,0 @@
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
}
-24
View File
@@ -1,24 +0,0 @@
package cli //nolint:testpackage // exercises the unexported command constructor
import (
"testing"
"github.com/spf13/cobra"
"github.com/stretchr/testify/assert"
)
// TestJSONHelpSaysConfirmationPromptIsSkipped checks that the --json help
// of `snapshot remove` and `prune` says the flag skips the confirmation
// prompt. Both delete without asking under --json, since a prompt on
// stdout would break the JSON document.
func TestJSONHelpSaysConfirmationPromptIsSkipped(t *testing.T) {
t.Parallel()
for _, cmd := range []*cobra.Command{
newSnapshotRemoveCommand(),
NewPruneCommand(),
} {
assert.Contains(t, cmd.Flags().Lookup("json").Usage,
"skips the confirmation prompt", cmd.Name())
}
}
+1 -2
View File
@@ -57,8 +57,7 @@ referenced.`,
} }
cmd.Flags().BoolVar(&opts.Force, "force", false, "Skip confirmation prompt") cmd.Flags().BoolVar(&opts.Force, "force", false, "Skip confirmation prompt")
cmd.Flags().BoolVar(&opts.JSON, "json", false, cmd.Flags().BoolVar(&opts.JSON, "json", false, "Output pruning stats as JSON")
"Output pruning stats as JSON; skips the confirmation prompt, as --force does")
return cmd return cmd
} }
+1 -2
View File
@@ -280,8 +280,7 @@ nuke --force' — it is the single supported entry point for that.`,
cmd.Flags().BoolVarP(&opts.Force, "force", "f", false, "Skip confirmation prompt") cmd.Flags().BoolVarP(&opts.Force, "force", "f", false, "Skip confirmation prompt")
cmd.Flags().BoolVar(&opts.DryRun, "dry-run", false, cmd.Flags().BoolVar(&opts.DryRun, "dry-run", false,
"Show what would be removed without removing") "Show what would be removed without removing")
cmd.Flags().BoolVar(&opts.JSON, "json", false, cmd.Flags().BoolVar(&opts.JSON, "json", false, "Output result as JSON")
"Output result as JSON; skips the confirmation prompt, as --force does")
cmd.Flags().BoolVar(&opts.LocalOnly, "local-only", false, cmd.Flags().BoolVar(&opts.LocalOnly, "local-only", false,
"Skip remote cleanup; only touch the local index") "Skip remote cleanup; only touch the local index")
+21 -29
View File
@@ -3,7 +3,6 @@ package database
import ( import (
"context" "context"
"database/sql"
"fmt" "fmt"
"strings" "strings"
"testing" "testing"
@@ -368,7 +367,7 @@ func verifyBlobNullUploadTS(
} }
// createLargeDatasetFiles creates fileCount files and adds every other // createLargeDatasetFiles creates fileCount files and adds every other
// one to the snapshot, in one transaction as a backup writes them. // one to the snapshot.
func createLargeDatasetFiles( func createLargeDatasetFiles(
t *testing.T, t *testing.T,
repos *Repositories, repos *Repositories,
@@ -377,38 +376,31 @@ func createLargeDatasetFiles(
) { ) {
t.Helper() t.Helper()
ctx := context.Background()
start := time.Now() start := time.Now()
err := repos.WithTx(context.Background(), for i := range fileCount {
func(ctx context.Context, tx *sql.Tx) error { file := &File{
for i := range fileCount { Path: types.FilePath(fmt.Sprintf("/large/file%05d.txt", i)),
file := &File{ MTime: time.Now(),
Path: types.FilePath(fmt.Sprintf("/large/file%05d.txt", i)), Size: int64(i * 1024),
MTime: time.Now(), Mode: 0644,
Size: int64(i * 1024), UID: uint32(1000 + (i % 10)),
Mode: 0644, GID: uint32(1000 + (i % 10)),
UID: uint32(1000 + (i % 10)), }
GID: uint32(1000 + (i % 10)),
}
err := repos.Files.Create(ctx, tx, file) err := repos.Files.Create(ctx, nil, file)
if err != nil { if err != nil {
return fmt.Errorf("creating file %d: %w", i, err) t.Fatalf("failed to create file %d: %v", i, err)
} }
// Add half to snapshot // Add half to snapshot
if i%2 == 0 { if i%2 == 0 {
err = repos.Snapshots.AddFileByID(ctx, tx, snapshotID, file.ID) err = repos.Snapshots.AddFileByID(ctx, nil, snapshotID, file.ID)
if err != nil { if err != nil {
return err t.Fatal(err)
}
}
} }
}
return nil
})
if err != nil {
t.Fatal(err)
} }
t.Logf("Created %d files in %v", fileCount, time.Since(start)) t.Logf("Created %d files in %v", fileCount, time.Since(start))
+53 -69
View File
@@ -56,27 +56,23 @@ const (
// ProgressStats holds atomic counters for progress tracking // ProgressStats holds atomic counters for progress tracking
type ProgressStats struct { type ProgressStats struct {
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped) FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
FilesProcessed atomic.Int64 // Files actually processed in phase 2 FilesProcessed atomic.Int64 // Files actually processed in phase 2
FilesSkipped atomic.Int64 // Files skipped due to no changes FilesSkipped atomic.Int64 // Files skipped due to no changes
BytesScanned atomic.Int64 // Bytes from new/changed files only BytesScanned atomic.Int64 // Bytes from new/changed files only
BytesSkipped atomic.Int64 // Bytes from unchanged files BytesSkipped atomic.Int64 // Bytes from unchanged files
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation) BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
ChunksCreated atomic.Int64 ChunksCreated atomic.Int64
BlobsCreated atomic.Int64 BlobsCreated atomic.Int64
BlobsUploaded atomic.Int64 BlobsUploaded atomic.Int64
BytesUploaded atomic.Int64 BytesUploaded atomic.Int64
CurrentFile atomic.Value // stores string CurrentFile atomic.Value // stores string
TotalSize atomic.Int64 // Size to process in the paths scanned so far TotalSize atomic.Int64 // Total size to process (set after scan phase)
TotalFiles atomic.Int64 // Files to process in the paths scanned so far TotalFiles atomic.Int64 // Total files to process in phase 2
StartTime time.Time ProcessStartTime atomic.Value // stores time.Time when processing starts
mu sync.RWMutex StartTime time.Time
lastDetailTime time.Time mu sync.RWMutex
lastDetailTime time.Time
// Guarded by mu: when the last scan phase ended, and
// BytesProcessed at that moment.
processStartTime time.Time
processStartBytes int64
// Upload tracking // Upload tracking
CurrentUpload atomic.Value // stores *UploadInfo CurrentUpload atomic.Value // stores *UploadInfo
@@ -152,36 +148,10 @@ func (pr *ProgressReporter) GetStats() *ProgressStats {
return pr.stats return pr.stats
} }
// AddTotalSize adds the size one path of the snapshot has to process to // SetTotalSize sets the total size to process (after scan phase)
// the total, once that path's scan phase is done, and starts measuring func (pr *ProgressReporter) SetTotalSize(size int64) {
// the processing rate again. The processed counts run across every path, pr.stats.TotalSize.Store(size)
// so the total does too. Restarting the rate keeps a path's scan phase, pr.stats.ProcessStartTime.Store(time.Now().UTC())
// which processes nothing, out of the rate while that path is processed.
// During a later path's scan phase the rate is still the previous
// path's, and falls as that scan goes on.
func (pr *ProgressReporter) AddTotalSize(size int64) {
pr.stats.TotalSize.Add(size)
pr.stats.mu.Lock()
defer pr.stats.mu.Unlock()
pr.stats.processStartTime = time.Now().UTC()
pr.stats.processStartBytes = pr.stats.BytesProcessed.Load()
}
// processRate returns the bytes processed per second since the last scan
// phase ended, or 0 before the first path's scan phase is done.
func (s *ProgressStats) processRate() float64 {
s.mu.RLock()
defer s.mu.RUnlock()
if s.processStartTime.IsZero() {
return 0
}
processed := s.BytesProcessed.Load() - s.processStartBytes
return float64(processed) / time.Since(s.processStartTime).Seconds()
} }
// Helper functions // Helper functions
@@ -387,12 +357,19 @@ func (pr *ProgressReporter) printSummaryStatus() {
// Calculate ETA if we have total size and are processing // Calculate ETA if we have total size and are processing
etaStr := "" etaStr := ""
processRate := pr.stats.processRate() if totalSize > 0 && bytesProcessed > 0 {
if totalSize > 0 && processRate > 0 { processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
remainingBytes := totalSize - bytesProcessed if ok && !processStart.IsZero() {
remainingSeconds := float64(remainingBytes) / processRate processElapsed := time.Since(processStart)
eta := time.Duration(remainingSeconds * float64(time.Second))
etaStr = " | ETA: " + formatDuration(eta) rate := float64(bytesProcessed) / processElapsed.Seconds()
if rate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / rate
eta := time.Duration(remainingSeconds * float64(time.Second))
etaStr = " | ETA: " + formatDuration(eta)
}
}
} }
rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds() rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds()
@@ -444,18 +421,25 @@ func (pr *ProgressReporter) printDetailedStatus() {
log.Info("Elapsed time", "duration", formatDuration(elapsed)) log.Info("Elapsed time", "duration", formatDuration(elapsed))
// Calculate and show ETA if we have data // Calculate and show ETA if we have data
processRate := pr.stats.processRate() if totalSize > 0 && bytesProcessed > 0 {
if totalSize > 0 && processRate > 0 { processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
remainingBytes := totalSize - bytesProcessed if ok && !processStart.IsZero() {
remainingSeconds := float64(remainingBytes) / processRate processElapsed := time.Since(processStart)
eta := time.Duration(remainingSeconds * float64(time.Second))
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale processRate := float64(bytesProcessed) / processElapsed.Seconds()
log.Info("Overall progress", if processRate > 0 {
"percent", fmt.Sprintf("%.1f%%", percentComplete), remainingBytes := totalSize - bytesProcessed
"processed", humanize.Bytes(safeUint64(bytesProcessed)), remainingSeconds := float64(remainingBytes) / processRate
"total", humanize.Bytes(safeUint64(totalSize)), eta := time.Duration(remainingSeconds * float64(time.Second))
"rate", humanize.Bytes(uint64(processRate))+"/s", percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
"eta", formatDuration(eta)) log.Info("Overall progress",
"percent", fmt.Sprintf("%.1f%%", percentComplete),
"processed", humanize.Bytes(safeUint64(bytesProcessed)),
"total", humanize.Bytes(safeUint64(totalSize)),
"rate", humanize.Bytes(uint64(processRate))+"/s",
"eta", formatDuration(eta))
}
}
} }
log.Info("Files processed", log.Info("Files processed",
-54
View File
@@ -1,54 +0,0 @@
//nolint:testpackage // exercises the unexported processRate helper
package snapshot
import (
"testing"
"time"
)
// TestProcessRateNotLoweredByLaterScanPhase feeds the reporter two paths
// that each take 10 seconds to process, the second after a 30-second scan
// phase. That scan phase processes nothing, so the rate while the second
// path is processed must be that path's own.
func TestProcessRateNotLoweredByLaterScanPhase(t *testing.T) {
t.Parallel()
const (
pathSize = 1000
processingTime = 10 * time.Second
scanTime = 30 * time.Second
)
// Never started, but Stop releases its tickers and signal handler.
progress := NewProgressReporter()
defer progress.Stop()
stats := progress.GetStats()
// Moves the processing start time back instead of sleeping.
elapse := func(d time.Duration) {
stats.mu.Lock()
defer stats.mu.Unlock()
stats.processStartTime = stats.processStartTime.Add(-d)
}
progress.AddTotalSize(pathSize)
stats.BytesProcessed.Add(pathSize)
elapse(processingTime)
elapse(scanTime)
progress.AddTotalSize(pathSize)
stats.BytesProcessed.Add(pathSize)
elapse(processingTime)
want := pathSize / processingTime.Seconds()
got := stats.processRate()
// The test's own run time adds to the elapsed time, so got is a hair
// under want.
if got > want || got < want*0.99 {
t.Errorf("rate while the second path is processed is %.1f bytes/s, "+
"want %.1f", got, want)
}
}
-94
View File
@@ -1,94 +0,0 @@
package snapshot_test
import (
"context"
"path/filepath"
"strings"
"testing"
"github.com/spf13/afero"
"sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/snapshot"
)
// TestProgressPercentWithSecondPathSmaller backs up two paths with one
// scanner, as a snapshot with two paths does, the second path smaller
// than the first. The progress line divides the bytes processed by the
// total size, so both must count both paths to stay within 100%.
func TestProgressPercentWithSecondPathSmaller(t *testing.T) {
t.Parallel()
fs := afero.NewMemMapFs()
files := map[string]string{
"/large/one.txt": strings.Repeat("1", 4000),
"/large/two.txt": strings.Repeat("2", 4000),
"/small/three.txt": strings.Repeat("3", 1000),
}
for path, content := range files {
err := fs.MkdirAll(filepath.Dir(path), 0755)
if err != nil {
t.Fatalf("mkdir: %v", err)
}
err = afero.WriteFile(fs, path, []byte(content), 0644)
if err != nil {
t.Fatalf("write %s: %v", path, err)
}
}
db, err := database.NewTestDB()
if err != nil {
t.Fatalf("create test db: %v", err)
}
t.Cleanup(func() {
cerr := db.Close()
if cerr != nil {
t.Errorf("close db: %v", cerr)
}
})
repos := database.NewRepositories(db)
scanner := snapshot.NewScanner(snapshot.ScannerConfig{
FS: fs,
ChunkSize: int64(1024 * 16),
Repositories: repos,
MaxBlobSize: int64(1024 * 1024),
CompressionLevel: 3,
AgeRecipients: []string{testAgePublicKey},
EnableProgress: true,
})
// Never started, but Stop releases its tickers and signal handler.
progress := scanner.GetProgress()
defer progress.Stop()
ctx := context.Background()
snapshotID := "test-snapshot-progress"
createTestSnapshotRecord(ctx, t, repos, snapshotID)
for _, path := range []string{"/large", "/small"} {
_, err := scanner.Scan(ctx, path, snapshotID)
if err != nil {
t.Fatalf("scanning %s: %v", path, err)
}
}
stats := progress.GetStats()
// Directories count toward the total size but produce no chunks, so
// the percentage ends just under 100%.
percent := float64(stats.BytesProcessed.Load()) /
float64(stats.TotalSize.Load()) * 100
if percent > 100 {
t.Errorf("progress after both paths is %.1f%%, want at most 100%%",
percent)
}
if stats.FilesProcessed.Load() != stats.TotalFiles.Load() {
t.Errorf("progress after both paths is %d of %d files, want all",
stats.FilesProcessed.Load(), stats.TotalFiles.Load())
}
}
+33 -67
View File
@@ -133,13 +133,10 @@ type ScannerConfig struct {
// ScanResult contains the results of a scan operation. Files and bytes // ScanResult contains the results of a scan operation. Files and bytes
// are counted per file: BytesScanned is the size of the new and changed // are counted per file: BytesScanned is the size of the new and changed
// files, BytesSkipped that of the unchanged ones. FilesFailed counts the // files, BytesSkipped that of the unchanged ones.
// new and changed files that phase 2 could not store; FilesScanned
// includes them and BytesScanned does not.
type ScanResult struct { type ScanResult struct {
FilesScanned int FilesScanned int
FilesSkipped int FilesSkipped int
FilesFailed int
FilesDeleted int FilesDeleted int
BytesScanned int64 BytesScanned int64
BytesSkipped int64 BytesSkipped int64
@@ -410,8 +407,8 @@ func (s *Scanner) summarizeScanPhase(
} }
if s.progress != nil { if s.progress != nil {
s.progress.AddTotalSize(totalSizeToProcess) s.progress.SetTotalSize(totalSizeToProcess)
s.progress.GetStats().TotalFiles.Add(int64(len(filesToProcess))) s.progress.GetStats().TotalFiles.Store(int64(len(filesToProcess)))
} }
log.Info("Phase 1 complete", log.Info("Phase 1 complete",
@@ -892,10 +889,9 @@ func (s *Scanner) scanPhase(
} }
// Handle symlinks and directories // Handle symlinks and directories
handled, err := s.recordSpecialEntry( if handled := s.recordSpecialEntry(
filePath, info, existingFiles, collector, result) filePath, info, existingFiles, collector, result); handled {
if handled { return nil
return err
} }
// Skip other non-regular files (devices, sockets, etc.) // Skip other non-regular files (devices, sockets, etc.)
@@ -937,25 +933,22 @@ func (s *Scanner) scanPhase(
} }
// recordSpecialEntry records symlinks and directories (which have no // recordSpecialEntry records symlinks and directories (which have no
// data to chunk) and reports whether it handled the entry. For a symlink // data to chunk) and reports whether it handled the entry.
// whose target cannot be read it returns handleWalkError's result.
func (s *Scanner) recordSpecialEntry( func (s *Scanner) recordSpecialEntry(
filePath string, info os.FileInfo, filePath string, info os.FileInfo,
existingFiles map[string]struct{}, existingFiles map[string]struct{},
collector *scanCollector, result *ScanResult, collector *scanCollector, result *ScanResult,
) (bool, error) { ) bool {
// Handle symlinks // Handle symlinks
if info.Mode()&os.ModeSymlink != 0 { if info.Mode()&os.ModeSymlink != 0 {
file, err := s.buildSymlinkEntry(filePath, info) file := s.buildSymlinkEntry(filePath, info)
if err != nil { if file != nil {
return true, s.handleWalkError(filePath, err) existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
} }
existingFiles[filePath] = struct{}{} return true
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
return true, nil
} }
// Handle directories (record for permission/ownership preservation // Handle directories (record for permission/ownership preservation
@@ -965,10 +958,10 @@ func (s *Scanner) recordSpecialEntry(
existingFiles[filePath] = struct{}{} existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file) collector.addToProcess(filePath, info, file)
return true, nil return true
} }
return false, nil return false
} }
// handleWalkError deals with a filesystem error surfaced by the walk: // handleWalkError deals with a filesystem error surfaced by the walk:
@@ -1118,12 +1111,13 @@ func (s *Scanner) printScanProgressLine(
} }
// buildSymlinkEntry creates a File record for a symlink. // buildSymlinkEntry creates a File record for a symlink.
func (s *Scanner) buildSymlinkEntry( // Returns nil if the link target cannot be read.
path string, info os.FileInfo, func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.File {
) (*database.File, error) {
target, err := os.Readlink(path) target, err := os.Readlink(path)
if err != nil { if err != nil {
return nil, err log.Debug("Cannot read symlink target", "path", path, "error", err)
return nil
} }
var uid, gid uint32 var uid, gid uint32
@@ -1142,7 +1136,7 @@ func (s *Scanner) buildSymlinkEntry(
UID: uid, UID: uid,
GID: gid, GID: gid,
LinkTarget: types.FilePath(target), LinkTarget: types.FilePath(target),
}, nil }
} }
// buildDirectoryEntry creates a File record for a directory. // buildDirectoryEntry creates a File record for a directory.
@@ -1312,13 +1306,6 @@ func (s *Scanner) processPhase(
// Process each file // Process each file
for _, fileToProcess := range filesToProcess { for _, fileToProcess := range filesToProcess {
// Check context cancellation
select {
case <-ctx.Done():
return ctx.Err()
default:
}
// Update progress // Update progress
if s.progress != nil { if s.progress != nil {
s.progress.GetStats().CurrentFile.Store(fileToProcess.Path) s.progress.GetStats().CurrentFile.Store(fileToProcess.Path)
@@ -1355,32 +1342,13 @@ func (s *Scanner) processPhase(
return s.finalizeProcessPhase(ctx, result) return s.finalizeProcessPhase(ctx, result)
} }
// processFileWithErrorHandling records a directory or symlink, or wraps // processFileWithErrorHandling wraps processFileStreaming with error recovery for
// processFileStreaming for a regular file with error recovery for deleted // deleted files and skip-errors mode. Returns (skipped, error).
// files and skip-errors mode. Returns (skipped, error).
func (s *Scanner) processFileWithErrorHandling( func (s *Scanner) processFileWithErrorHandling(
ctx context.Context, fileToProcess *FileToProcess, result *ScanResult, ctx context.Context, fileToProcess *FileToProcess, result *ScanResult,
) (bool, error) { ) (bool, error) {
// A directory or symlink has no data to open or read; it is only
// recorded in the local index. An error recording it stops the run
// even under --skip-errors, or the snapshot would complete without it.
mode := os.FileMode(fileToProcess.File.Mode)
if mode&os.ModeSymlink != 0 || mode.IsDir() {
err := s.recordNonRegularFile(ctx, fileToProcess)
if err != nil {
return false, fmt.Errorf("processing file %s: %w", fileToProcess.Path, err)
}
return false, nil
}
err := s.processFileStreaming(ctx, fileToProcess, result) err := s.processFileStreaming(ctx, fileToProcess, result)
if err != nil { if err != nil {
// A cancelled run stops here even under --skip-errors, rather than
// counting this file as failed and going on to the next one.
if ctx.Err() != nil {
return false, fmt.Errorf("processing file %s: %w", fileToProcess.Path, err)
}
// A packer/database/encryption/upload failure means the chunk's data // A packer/database/encryption/upload failure means the chunk's data
// may not have been stored. Skipping the file would let the snapshot // may not have been stored. Skipping the file would let the snapshot
// record a file whose chunk is in no blob and cannot be restored, so // record a file whose chunk is in no blob and cannot be restored, so
@@ -1394,7 +1362,7 @@ func (s *Scanner) processFileWithErrorHandling(
log.Warn("File was deleted during backup, skipping", log.Warn("File was deleted during backup, skipping",
"path", fileToProcess.Path) "path", fileToProcess.Path)
countFailedFile(fileToProcess, result) result.FilesSkipped++
return true, nil return true, nil
} }
@@ -1405,7 +1373,7 @@ func (s *Scanner) processFileWithErrorHandling(
s.ui.Errorf("Failed to process %s: %v. Skipping (--skip-errors).", s.ui.Errorf("Failed to process %s: %v. Skipping (--skip-errors).",
s.ui.Path(fileToProcess.Path), err) s.ui.Path(fileToProcess.Path), err)
countFailedFile(fileToProcess, result) result.FilesSkipped++
return true, nil return true, nil
} }
@@ -1416,13 +1384,6 @@ func (s *Scanner) processFileWithErrorHandling(
return false, nil return false, nil
} }
// countFailedFile counts a file that phase 2 could not store as failed
// and takes its size back out of BytesScanned, where phase 1 put it.
func countFailedFile(fileToProcess *FileToProcess, result *ScanResult) {
result.FilesFailed++
result.BytesScanned -= fileToProcess.FileInfo.Size()
}
// printProcessingProgress prints a periodic progress line during the process phase, // printProcessingProgress prints a periodic progress line during the process phase,
// showing files processed, bytes transferred, throughput, and ETA // showing files processed, bytes transferred, throughput, and ETA
func (s *Scanner) printProcessingProgress( func (s *Scanner) printProcessingProgress(
@@ -1791,11 +1752,16 @@ func (e *packerError) Error() string { return e.err.Error() }
func (e *packerError) Unwrap() error { return e.err } func (e *packerError) Unwrap() error { return e.err }
// processFileStreaming processes a regular file by streaming chunks directly // processFileStreaming processes a file by streaming chunks directly to the packer
// to the packer
func (s *Scanner) processFileStreaming( func (s *Scanner) processFileStreaming(
ctx context.Context, fileToProcess *FileToProcess, result *ScanResult, ctx context.Context, fileToProcess *FileToProcess, result *ScanResult,
) error { ) error {
// Symlinks and directories have no data to chunk — just record them in the DB.
mode := os.FileMode(fileToProcess.File.Mode)
if mode&os.ModeSymlink != 0 || mode.IsDir() {
return s.recordNonRegularFile(ctx, fileToProcess)
}
file, err := s.fs.Open(fileToProcess.Path) file, err := s.fs.Open(fileToProcess.Path)
if err != nil { if err != nil {
return fmt.Errorf("opening file: %w", wrapPermissionError(fileToProcess.Path, err)) return fmt.Errorf("opening file: %w", wrapPermissionError(fileToProcess.Path, err))
+8 -175
View File
@@ -3,7 +3,6 @@ package snapshot_test
import ( import (
"context" "context"
"errors" "errors"
"io"
"os" "os"
"path/filepath" "path/filepath"
"strings" "strings"
@@ -14,7 +13,6 @@ import (
"github.com/spf13/afero" "github.com/spf13/afero"
"sneak.berlin/go/vaultik/internal/database" "sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/snapshot" "sneak.berlin/go/vaultik/internal/snapshot"
"sneak.berlin/go/vaultik/internal/ui"
) )
// errSimTempFail is the one-time temp-file creation failure blobTempFailFs // errSimTempFail is the one-time temp-file creation failure blobTempFailFs
@@ -82,56 +80,6 @@ func (f *readFailFs) Open(name string) (afero.File, error) {
return file, nil return file, nil
} }
// cancelOnOpenFs cancels the run when the scanner opens the target file to
// back it up, as Ctrl-C partway through processing would, and records each
// file opened after that.
type cancelOnOpenFs struct {
afero.Fs
target string
cancel context.CancelFunc
cancelled bool
openedAfterCancel []string
}
//nolint:ireturn // afero.Fs.Open is defined to return the interface.
func (f *cancelOnOpenFs) Open(name string) (afero.File, error) {
if f.cancelled {
f.openedAfterCancel = append(f.openedAfterCancel, name)
}
if name == f.target {
f.cancel()
f.cancelled = true
}
return f.Fs.Open(name)
}
// 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. // writeSkipErrorTestFile writes one file into fs with a fixed mtime.
func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) { func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
t.Helper() t.Helper()
@@ -154,16 +102,14 @@ func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
} }
} }
// runSkipErrorScan scans source on fs under ctx with the given skip-errors // runSkipErrorScan scans /source on fs with the given skip-errors setting and
// setting, printing user-facing messages to uiw (nil discards them), and
// returns the repositories (for inspection) and the scan error. // returns the repositories (for inspection) and the scan error.
func runSkipErrorScan( func runSkipErrorScan(
ctx context.Context, t *testing.T, fs afero.Fs, source string, t *testing.T, fs afero.Fs, skipErrors bool,
skipErrors bool, uiw *ui.Writer,
) (*database.Repositories, error) { ) (*database.Repositories, error) {
t.Helper() t.Helper()
db, err := database.New(ctx, ":memory:") db, err := database.NewTestDB()
if err != nil { if err != nil {
t.Fatalf("create test db: %v", err) t.Fatalf("create test db: %v", err)
} }
@@ -184,14 +130,14 @@ func runSkipErrorScan(
MaxBlobSize: int64(1024 * 1024), MaxBlobSize: int64(1024 * 1024),
CompressionLevel: 3, CompressionLevel: 3,
AgeRecipients: []string{testAgePublicKey}, AgeRecipients: []string{testAgePublicKey},
UI: uiw,
SkipErrors: skipErrors, SkipErrors: skipErrors,
}) })
ctx := context.Background()
snapshotID := "test-snapshot-skip-errors" snapshotID := "test-snapshot-skip-errors"
createTestSnapshotRecord(ctx, t, repos, snapshotID) createTestSnapshotRecord(ctx, t, repos, snapshotID)
_, err = scanner.Scan(ctx, source, snapshotID) _, err = scanner.Scan(ctx, "/source", snapshotID)
return repos, err return repos, err
} }
@@ -211,7 +157,7 @@ func TestScannerPackingFailureAbortsUnderSkipErrors(t *testing.T) {
writeSkipErrorTestFile(t, fs, "/source/file1.txt", "first file content") writeSkipErrorTestFile(t, fs, "/source/file1.txt", "first file content")
writeSkipErrorTestFile(t, fs, "/source/file2.txt", "second file content") writeSkipErrorTestFile(t, fs, "/source/file2.txt", "second file content")
repos, err := runSkipErrorScan(context.Background(), t, fs, "/source", true, nil) repos, err := runSkipErrorScan(t, fs, true)
if err == nil { if err == nil {
t.Fatal("expected scan to abort on the packer error, got nil") t.Fatal("expected scan to abort on the packer error, got nil")
} }
@@ -238,7 +184,7 @@ func TestScannerReadErrorAbortsWithoutSkipErrors(t *testing.T) {
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target} fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
writeSkipErrorTestFile(t, fs, target, "content that cannot be read") writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
_, err := runSkipErrorScan(context.Background(), t, fs, "/source", false, nil) _, err := runSkipErrorScan(t, fs, false)
if err == nil { if err == nil {
t.Fatal("expected scan to fail on the read error, got nil") t.Fatal("expected scan to fail on the read error, got nil")
} }
@@ -254,7 +200,7 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target} fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
writeSkipErrorTestFile(t, fs, target, "content that cannot be read") writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
repos, err := runSkipErrorScan(context.Background(), t, fs, "/source", true, nil) repos, err := runSkipErrorScan(t, fs, true)
if err != nil { if err != nil {
t.Fatalf("expected scan to complete with --skip-errors, got %v", err) t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
} }
@@ -268,116 +214,3 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
t.Fatalf("expected unreadable file skipped, got %d chunks", len(chunks)) 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(context.Background(), 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(context.Background(), 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)
}
}
// TestScannerCancelStopsSkipErrorsRun checks that a --skip-errors backup
// cancelled partway through processing returns the cancellation error,
// opens no further file, and reports no file as failed.
func TestScannerCancelStopsSkipErrorsRun(t *testing.T) {
t.Parallel()
tests := []struct {
name string
targetContent string
}{
// The cancellation lands while the target is being read.
{name: "while reading a file", targetContent: "first file content"},
// An empty target has no chunks, so the cancellation goes unnoticed
// until the run moves on to the next file.
{name: "between files", targetContent: ""},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
const target = "/source/a.txt"
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
fs := &cancelOnOpenFs{
Fs: afero.NewMemMapFs(), target: target, cancel: cancel,
}
writeSkipErrorTestFile(t, fs, target, tt.targetContent)
writeSkipErrorTestFile(t, fs, "/source/b.txt", "second file content")
writeSkipErrorTestFile(t, fs, "/source/c.txt", "third file content")
uiw := ui.NewWithColor(io.Discard, false)
_, err := runSkipErrorScan(ctx, t, fs, "/source", true, uiw)
if !errors.Is(err, context.Canceled) {
t.Fatalf("expected the cancellation error, got %v", err)
}
if len(fs.openedAfterCancel) != 0 {
t.Fatalf("expected no file opened after the cancellation, got %v",
fs.openedAfterCancel)
}
if uiw.ErrorCount() != 0 {
t.Fatalf("expected no file reported as failed, got %d error lines",
uiw.ErrorCount())
}
})
}
}
+9 -26
View File
@@ -115,37 +115,20 @@ func ShortHostname(hostname string) string {
// CreateSnapshotWithName creates a new snapshot record with an optional // CreateSnapshotWithName creates a new snapshot record with an optional
// snapshot name. The snapshot ID format is: hostname_name_timestamp or // snapshot name. The snapshot ID format is: hostname_name_timestamp or
// hostname_timestamp if name is empty. The timestamp is in whole seconds. // hostname_timestamp if name is empty.
// If the local index already has a snapshot with that ID, from a run of the
// same name that started in the same second, it waits a second and takes a
// new timestamp.
func (sm *SnapshotManager) CreateSnapshotWithName( func (sm *SnapshotManager) CreateSnapshotWithName(
ctx context.Context, hostname, name, version, gitRevision string, ctx context.Context, hostname, name, version, gitRevision string,
) (string, error) { ) (string, error) {
shortHostname := ShortHostname(hostname) shortHostname := ShortHostname(hostname)
// Build snapshot ID with optional name
timestamp := time.Now().UTC().Format("2006-01-02T15:04:05Z")
var snapshotID string var snapshotID string
if name != "" {
for { snapshotID = fmt.Sprintf("%s_%s_%s", shortHostname, name, timestamp)
// Build snapshot ID with optional name } else {
timestamp := time.Now().UTC().Format("2006-01-02T15:04:05Z") snapshotID = fmt.Sprintf("%s_%s", shortHostname, timestamp)
if name != "" {
snapshotID = fmt.Sprintf("%s_%s_%s", shortHostname, name, timestamp)
} else {
snapshotID = fmt.Sprintf("%s_%s", shortHostname, timestamp)
}
existing, err := sm.repos.Snapshots.GetByID(ctx, snapshotID)
if err != nil {
return "", fmt.Errorf("looking up snapshot %s: %w", snapshotID, err)
}
if existing == nil {
break
}
time.Sleep(time.Second)
} }
snapshot := &database.Snapshot{ snapshot := &database.Snapshot{
@@ -892,7 +875,7 @@ func (sm *SnapshotManager) getFileSize(path string) int64 {
// BackupStats contains statistics from a backup operation // BackupStats contains statistics from a backup operation
type BackupStats struct { type BackupStats struct {
FilesScanned int FilesScanned int
TotalSize int64 // Total size of the files in the snapshot TotalSize int64 // Total size of all files examined
ChunksCreated int ChunksCreated int
BlobsCreated int BlobsCreated int
BytesUploaded int64 BytesUploaded int64
-33
View File
@@ -386,36 +386,3 @@ func TestCleanSnapshotDBNonExistentSnapshot(t *testing.T) {
t.Fatalf("unexpected error: %v", err) t.Fatalf("unexpected error: %v", err)
} }
} }
// Two creates of one snapshot name back to back start within one second,
// the resolution of the timestamp in a snapshot ID. See
// https://git.eeqj.de/sneak/vaultik/issues/270.
func TestCreateSnapshotWithNameTwiceBackToBack(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background()
db, err := database.New(ctx, filepath.Join(t.TempDir(), "index.sqlite"))
if err != nil {
t.Fatalf("failed to create database: %v", err)
}
defer func() { _ = db.Close() }()
sm := &SnapshotManager{repos: database.NewRepositories(db)}
first, err := sm.CreateSnapshotWithName(ctx, "test-host", "data", "v", "g")
if err != nil {
t.Fatalf("first create failed: %v", err)
}
second, err := sm.CreateSnapshotWithName(ctx, "test-host", "data", "v", "g")
if err != nil {
t.Fatalf("second create failed: %v", err)
}
if first == second {
t.Fatalf("both creates returned snapshot ID %s", first)
}
}
+2 -3
View File
@@ -14,9 +14,8 @@ import (
// runStorerConformance is the shared Storer contract. Every backend that // runStorerConformance is the shared Storer contract. Every backend that
// can run in-process is expected to pass it: TestFileStorer runs it against // can run in-process is expected to pass it: TestFileStorer runs it against
// file://, TestS3Storer against s3://, TestRcloneStorer against rclone's // file://, TestS3Storer against s3://. A new backend inherits this coverage
// local backend. A new backend inherits this coverage by passing its own // by passing its own constructor, so the contract is defined once.
// constructor, so the contract is defined once.
// //
// It exercises the public Storer interface: round-trip, stat, list with // It exercises the public Storer interface: round-trip, stat, list with
// prefix filtering, overwrite, delete, delete-of-missing, and not-found on // prefix filtering, overwrite, delete, delete-of-missing, and not-found on
@@ -171,11 +171,6 @@ func (f *Storer) ListStream(
return f.inner.ListStream(ctx, prefix) return f.inner.ListStream(ctx, prefix)
} }
// DeletePartialUploads delegates unchanged.
func (f *Storer) DeletePartialUploads(ctx context.Context, prefix string) error {
return f.inner.DeletePartialUploads(ctx, prefix)
}
// Info delegates unchanged. // Info delegates unchanged.
func (f *Storer) Info() storage.Info { func (f *Storer) Info() storage.Info {
return f.inner.Info() return f.inner.Info()
+3 -40
View File
@@ -50,10 +50,9 @@ const storageDirPerm = 0o755
// temp file carrying this suffix and only renames it onto the real key once // 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 // 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 // truncated object at the key a later run would Stat and trust as a complete
// blob. The rclone backend's upload does the same on remotes with a // blob. List and ListStream skip these files, so a leftover from an
// server-side move. List and ListStream skip these files, so a leftover from // interrupted write is never listed or trusted as a blob; it is otherwise
// 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.
// harmless.
const tempSuffix = ".partial" const tempSuffix = ".partial"
// Put stores data at the specified key. // Put stores data at the specified key.
@@ -242,42 +241,6 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec
return ch return ch
} }
// DeletePartialUploads removes every file under prefix whose name ends in
// tempSuffix. A missing prefix has none to remove.
func (f *FileStorer) DeletePartialUploads(ctx context.Context, prefix string) error {
basePath := f.fullPath(prefix)
exists, err := afero.Exists(f.fs, basePath)
if err != nil {
return fmt.Errorf("checking path: %w", err)
}
if !exists {
return nil
}
err = afero.Walk(f.fs, basePath, func(path string, info os.FileInfo, err error) error {
if err != nil {
return err
}
if ctx.Err() != nil {
return ctx.Err()
}
if info.IsDir() || !strings.HasSuffix(info.Name(), tempSuffix) {
return nil
}
return f.fs.Remove(path)
})
if err != nil {
return fmt.Errorf("walking directory: %w", err)
}
return nil
}
// Info returns human-readable storage location information. // Info returns human-readable storage location information.
func (f *FileStorer) Info() Info { func (f *FileStorer) Info() Info {
return Info{ return Info{
+2 -47
View File
@@ -14,9 +14,6 @@ import (
// errStreamInterrupted stands in for an upload cut off mid-stream. // errStreamInterrupted stands in for an upload cut off mid-stream.
var errStreamInterrupted = errors.New("connection reset mid-upload") 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. // failingReader yields its data once, then fails.
type failingReader struct { type failingReader struct {
data []byte data []byte
@@ -46,7 +43,7 @@ func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
} }
ctx := context.Background() ctx := context.Background()
key := testBlobKey key := "blobs/aa/bb/aabbccddeeff"
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil) err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
if err == nil { if err == nil {
@@ -82,7 +79,7 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
} }
ctx := context.Background() ctx := context.Background()
realKey := testBlobKey realKey := "blobs/aa/bb/aabbccddeeff"
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes")) err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
if err != nil { if err != nil {
@@ -120,45 +117,3 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
t.Fatalf("ListStream should return only the real key, got %v", streamed) t.Fatalf("ListStream should return only the real key, got %v", streamed)
} }
} }
// TestFileStorer_DeletePartialUploads checks that a leftover temp file is
// removed and the object at the real key is kept.
func TestFileStorer_DeletePartialUploads(t *testing.T) {
t.Parallel()
base := t.TempDir()
f, err := storage.NewFileStorer(base)
if err != nil {
t.Fatalf("NewFileStorer: %v", err)
}
ctx := context.Background()
err = f.Put(ctx, testBlobKey, strings.NewReader("blob-bytes"))
if err != nil {
t.Fatalf("Put: %v", err)
}
leftover := filepath.Join(base, testBlobKey+"-123456.partial")
err = os.WriteFile(leftover, []byte("half"), 0o600)
if err != nil {
t.Fatalf("writing leftover temp file: %v", err)
}
err = f.DeletePartialUploads(ctx, "blobs/")
if err != nil {
t.Fatalf("DeletePartialUploads: %v", err)
}
_, err = os.Stat(leftover)
if !os.IsNotExist(err) {
t.Errorf("leftover temp file was not removed: %v", err)
}
_, err = f.Stat(ctx, testBlobKey)
if err != nil {
t.Errorf("Stat of the real key: %v", err)
}
}
+15 -89
View File
@@ -3,7 +3,6 @@ package storage
import ( import (
"bytes" "bytes"
"context" "context"
"crypto/rand"
"errors" "errors"
"fmt" "fmt"
"io" "io"
@@ -69,7 +68,14 @@ func (r *RcloneStorer) Put(ctx context.Context, key string, data io.Reader) erro
return fmt.Errorf("reading data: %w", err) return fmt.Errorf("reading data: %w", err)
} }
return r.upload(ctx, key, bytes.NewReader(buf)) // 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
} }
// PutWithProgress stores data with progress reporting. // PutWithProgress stores data with progress reporting.
@@ -83,7 +89,13 @@ func (r *RcloneStorer) PutWithProgress(
callback: progress, callback: progress,
} }
return r.upload(ctx, key, pr) // 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
} }
// Get retrieves data from the specified key. // Get retrieves data from the specified key.
@@ -161,10 +173,6 @@ func (r *RcloneStorer) List(ctx context.Context, prefix string) ([]string, error
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) { err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
key := obj.Remote() key := obj.Remote()
if strings.HasSuffix(key, tempSuffix) {
return
}
if prefix == "" || strings.HasPrefix(key, prefix) { if prefix == "" || strings.HasPrefix(key, prefix) {
keys = append(keys, key) keys = append(keys, key)
} }
@@ -194,10 +202,6 @@ func (r *RcloneStorer) ListStream(
} }
key := obj.Remote() key := obj.Remote()
if strings.HasSuffix(key, tempSuffix) {
return
}
if prefix == "" || strings.HasPrefix(key, prefix) { if prefix == "" || strings.HasPrefix(key, prefix) {
ch <- ObjectInfo{ ch <- ObjectInfo{
Key: key, Key: key,
@@ -213,31 +217,6 @@ func (r *RcloneStorer) ListStream(
return ch return ch
} }
// DeletePartialUploads removes every object under prefix whose name ends
// in tempSuffix.
func (r *RcloneStorer) DeletePartialUploads(ctx context.Context, prefix string) error {
var partial []fs.Object
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
key := obj.Remote()
if strings.HasPrefix(key, prefix) && strings.HasSuffix(key, tempSuffix) {
partial = append(partial, obj)
}
})
if err != nil {
return fmt.Errorf("listing objects: %w", err)
}
for _, obj := range partial {
err = obj.Remove(ctx)
if err != nil {
return fmt.Errorf("removing object: %w", err)
}
}
return nil
}
// Info returns human-readable storage location information. // Info returns human-readable storage location information.
func (r *RcloneStorer) Info() Info { func (r *RcloneStorer) Info() Info {
location := r.remote location := r.remote
@@ -251,59 +230,6 @@ 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. // progressReader wraps an io.Reader to track read progress.
type progressReader struct { type progressReader struct {
reader io.Reader reader io.Reader
+14 -348
View File
@@ -1,47 +1,28 @@
package storage_test package storage_test
import ( import (
"bytes"
"context" "context"
"errors" "errors"
"io"
"os"
"path/filepath"
"strings"
"testing" "testing"
"github.com/rclone/rclone/fs"
"github.com/rclone/rclone/fs/config/configmap"
"sneak.berlin/go/vaultik/internal/storage" "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 // These tests use rclone's ":local:" on-the-fly backend, which addresses the
// local filesystem directly without any configured remote, so they run // local filesystem directly without any configured remote, so construction
// entirely in-process. A remote that needs credentials and network access // runs entirely in-process.
// (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 // TestNewRcloneStorerConstruction checks that a valid remote constructs a
// backend and that Info() reports the shaped "remote:path" location. // backend and that Info() reports the shaped "remote:path" location.
@@ -63,321 +44,6 @@ 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)
}
}
// TestRcloneStorerDeletePartialUploads checks that a temporary file left
// by a killed upload is removed and the object at the real key is kept.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerDeletePartialUploads(t *testing.T) {
dir := t.TempDir()
ctx := context.Background()
s, err := storage.NewRcloneStorer(ctx, ":local", dir)
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
err = s.Put(ctx, testBlobKey, strings.NewReader("blob-bytes"))
if err != nil {
t.Fatalf("Put: %v", err)
}
leftover := filepath.Join(dir, testBlobKey+"-123456.partial")
err = os.WriteFile(leftover, []byte("half"), 0o600)
if err != nil {
t.Fatalf("writing leftover temp file: %v", err)
}
err = s.DeletePartialUploads(ctx, "blobs/")
if err != nil {
t.Fatalf("DeletePartialUploads: %v", err)
}
_, err = os.Stat(leftover)
if !os.IsNotExist(err) {
t.Errorf("leftover temp file was not removed: %v", err)
}
_, err = s.Stat(ctx, testBlobKey)
if err != nil {
t.Errorf("Stat of the real key: %v", err)
}
}
// 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 // TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
// rclone config fails construction with the ErrRemoteNotFound sentinel, // rclone config fails construction with the ErrRemoteNotFound sentinel,
// rather than silently returning a backend pointed nowhere. // rather than silently returning a backend pointed nowhere.
-6
View File
@@ -99,12 +99,6 @@ func (s *S3Storer) ListStream(ctx context.Context, prefix string) <-chan ObjectI
return ch return ch
} }
// DeletePartialUploads has nothing to remove: S3 shows an object only once
// its upload has completed, so an upload cut off part-way leaves none.
func (s *S3Storer) DeletePartialUploads(_ context.Context, _ string) error {
return nil
}
// Info returns human-readable storage location information. // Info returns human-readable storage location information.
func (s *S3Storer) Info() Info { func (s *S3Storer) Info() Info {
return Info{ return Info{
-6
View File
@@ -71,12 +71,6 @@ type Storer interface {
// If an error occurs during listing, the final item will have Err set. // If an error occurs during listing, the final item will have Err set.
ListStream(ctx context.Context, prefix string) <-chan ObjectInfo ListStream(ctx context.Context, prefix string) <-chan ObjectInfo
// DeletePartialUploads removes every object under prefix that an
// upload cut off part-way left under a temporary name ending in
// `.partial`. The file and rclone backends' List and ListStream skip
// such an object.
DeletePartialUploads(ctx context.Context, prefix string) error
// Info returns human-readable storage location information. // Info returns human-readable storage location information.
Info() Info Info() Info
} }
@@ -130,9 +130,10 @@ func assertThirdSnapshotRestores(
// up, and that snapshot is removed. The first snapshot keeps the file row, // 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 // which now lists the appended content's chunks, while removal drops the
// blob that held them. // blob that held them.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupAfterRemovingNewestSnapshotRestoresChangedFile(t *testing.T) { func TestBackupAfterRemovingNewestSnapshotRestoresChangedFile(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -177,9 +178,10 @@ func TestBackupAfterRemovingNewestSnapshotRestoresChangedFile(t *testing.T) {
// The next run's prune drops that incomplete snapshot and its blob, while // 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 // the first snapshot keeps the file row, which now lists the appended
// content's chunks. // content's chunks.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupAfterInterruptedRunRestoresChangedFile(t *testing.T) { func TestBackupAfterInterruptedRunRestoresChangedFile(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -257,10 +257,6 @@ func (s *stubLister) List(_ context.Context, _ string) ([]string, error) {
return nil, errStubUnused return nil, errStubUnused
} }
func (s *stubLister) DeletePartialUploads(_ context.Context, _ string) error {
return errStubUnused
}
func (s *stubLister) Info() storage.Info { func (s *stubLister) Info() storage.Info {
return storage.Info{} return storage.Info{}
} }
+22 -15
View File
@@ -38,9 +38,11 @@ import (
// (https://git.eeqj.de/sneak/vaultik/issues/130) and is not re-tested // (https://git.eeqj.de/sneak/vaultik/issues/130) and is not re-tested
// here; these tests target the layers above the backend. // here; these tests target the layers above the backend.
// //
// log.Initialize replaces the package-global logger that a running // The tests run serially, not with t.Parallel: each calls
// backup or restore reads, so each test calls it before t.Parallel, // log.Initialize, which replaces the package-global logger, and a
// while no parallel test is running yet. // 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.
const ( const (
faultChunkSize = int64(64 * 1024) faultChunkSize = int64(64 * 1024)
@@ -163,19 +165,17 @@ func newReaderVaultik(
// Scenario 3: a stored blob's bytes are flipped before restore reads // 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 // them. Restore must fail loudly, and no file must be left on the
// restore target holding corrupt content. // restore target holding corrupt content.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestRestoreRejectsCorruptBlob(t *testing.T) { func TestRestoreRejectsCorruptBlob(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
assertRestoreRejectsDamagedBlob(t, faultstore.GetCorrupt, "corrupt") assertRestoreRejectsDamagedBlob(t, faultstore.GetCorrupt, "corrupt")
} }
// Scenario 4: a stored blob is truncated before restore reads it. Same // Scenario 4: a stored blob is truncated before restore reads it. Same
// contract as the corrupt case. // contract as the corrupt case.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestRestoreRejectsTruncatedBlob(t *testing.T) { func TestRestoreRejectsTruncatedBlob(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
assertRestoreRejectsDamagedBlob(t, faultstore.GetTruncate, "truncated") assertRestoreRejectsDamagedBlob(t, faultstore.GetTruncate, "truncated")
} }
@@ -188,6 +188,7 @@ func assertRestoreRejectsDamagedBlob(
t *testing.T, fault faultstore.GetFault, name string, t *testing.T, fault faultstore.GetFault, name string,
) { ) {
t.Helper() t.Helper()
log.Initialize(log.Config{})
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -231,9 +232,10 @@ func assertRestoreRejectsDamagedBlob(
// Scenario 6: the backend accepts blob uploads and reports success but // Scenario 6: the backend accepts blob uploads and reports success but
// stores nothing. verify --deep must catch it. // stores nothing. verify --deep must catch it.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestDeepVerifyCatchesLyingBackend(t *testing.T) { func TestDeepVerifyCatchesLyingBackend(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -283,9 +285,10 @@ func TestDeepVerifyCatchesLyingBackend(t *testing.T) {
// Scenario 1a: a blob upload fails partway through. The interrupted run // Scenario 1a: a blob upload fails partway through. The interrupted run
// must not record the blob as uploaded, must not reference it from the // must not record the blob as uploaded, must not reference it from the
// snapshot, and must leave no blob object at the destination. // snapshot, and must leave no blob object at the destination.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestInterruptedBlobUploadRecordsNoUploadedBlob(t *testing.T) { func TestInterruptedBlobUploadRecordsNoUploadedBlob(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -358,9 +361,10 @@ func TestInterruptedBlobUploadRecordsNoUploadedBlob(t *testing.T) {
// chunks in a blob that was actually uploaded, so the retry re-chunks and // 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 // re-uploads the affected data instead of silently referencing data that
// never reached storage. // never reached storage.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) { func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -422,9 +426,10 @@ func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) {
// covered by TestBackupCompletesOnlyAfterMetadataExport // covered by TestBackupCompletesOnlyAfterMetadataExport
// (https://git.eeqj.de/sneak/vaultik/issues/177); this test exercises the // (https://git.eeqj.de/sneak/vaultik/issues/177); this test exercises the
// lower-level export path in isolation. // lower-level export path in isolation.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupSurvivesMetadataExportInterruption(t *testing.T) { func TestBackupSurvivesMetadataExportInterruption(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -504,9 +509,10 @@ func TestBackupSurvivesMetadataExportInterruption(t *testing.T) {
// destination. Rerunning the backup must then prune the incomplete // destination. Rerunning the backup must then prune the incomplete
// snapshot, produce a snapshot whose destination metadata and local index // snapshot, produce a snapshot whose destination metadata and local index
// agree, and restore. See https://git.eeqj.de/sneak/vaultik/issues/177. // 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) { func TestBackupCompletesOnlyAfterMetadataExport(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
@@ -679,9 +685,10 @@ func faultScannerFactory(
// Scenario 5: the restore target runs out of space mid-file. Restore // 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 // must fail with an out-of-space error, and must not leave a truncated
// file at the target path presenting as a complete restore. // file at the target path presenting as a complete restore.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestRestoreReportsDiskFull(t *testing.T) { func TestRestoreReportsDiskFull(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
osFS := afero.NewOsFs() osFS := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
+6 -24
View File
@@ -181,11 +181,8 @@ type SnapshotMetadataInfo struct {
ManifestSize int64 `json:"manifest_size"` ManifestSize int64 `json:"manifest_size"`
DatabaseSize int64 `json:"database_size"` DatabaseSize int64 `json:"database_size"`
TotalSize int64 `json:"total_size"` TotalSize int64 `json:"total_size"`
BlobCount int `json:"blob_count"`
// Both stay nil (null in the JSON) when the snapshot's manifest was BlobsSize int64 `json:"blobs_size"`
// listed but could not be read.
BlobCount *int `json:"blob_count"`
BlobsSize *int64 `json:"blobs_size"`
// Set when the listing holds this snapshot's manifest.json.zst. A // Set when the listing holds this snapshot's manifest.json.zst. A
// backup interrupted before its manifest upload leaves a directory // backup interrupted before its manifest upload leaves a directory
@@ -383,10 +380,6 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
for _, snapshotID := range snapshotIDs { for _, snapshotID := range snapshotIDs {
info := snapshotMetadata[snapshotID] info := snapshotMetadata[snapshotID]
if !info.hasManifest { if !info.hasManifest {
// The orphan figures count this directory's blobs as
// orphaned, so it references none.
info.BlobCount, info.BlobsSize = new(int), new(int64)
continue continue
} }
@@ -402,7 +395,7 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
continue continue
} }
blobCount := manifest.BlobCount info.BlobCount = manifest.BlobCount
var blobsSize int64 var blobsSize int64
@@ -411,8 +404,7 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
blobsSize += blob.CompressedSize blobsSize += blob.CompressedSize
} }
info.BlobCount = &blobCount info.BlobsSize = blobsSize
info.BlobsSize = &blobsSize
} }
return referencedBlobs, unreadable return referencedBlobs, unreadable
@@ -524,23 +516,13 @@ func (v *Vaultik) printRemoteInfoTable(result *RemoteInfoResult) {
v.stdoutf("%s", separator) v.stdoutf("%s", separator)
for _, info := range result.Snapshots { for _, info := range result.Snapshots {
blobCount := unknownText
if info.BlobCount != nil {
blobCount = humanize.Comma(int64(*info.BlobCount))
}
blobsSize := unknownText
if info.BlobsSize != nil {
blobsSize = ubytes(*info.BlobsSize)
}
v.stdoutf(rowFormat, v.stdoutf(rowFormat,
truncateString(info.SnapshotID, snapshotIDColWidth), truncateString(info.SnapshotID, snapshotIDColWidth),
ubytes(info.ManifestSize), ubytes(info.ManifestSize),
ubytes(info.DatabaseSize), ubytes(info.DatabaseSize),
ubytes(info.TotalSize), ubytes(info.TotalSize),
blobCount, humanize.Comma(int64(info.BlobCount)),
blobsSize, ubytes(info.BlobsSize),
) )
} }
+11 -16
View File
@@ -155,12 +155,6 @@ func (m *MockStorer) ListStream(
return ch return ch
} }
// DeletePartialUploads has nothing to remove: Put stores each object
// under its key at once.
func (m *MockStorer) DeletePartialUploads(_ context.Context, _ string) error {
return nil
}
func (m *MockStorer) Info() storage.Info { func (m *MockStorer) Info() storage.Info {
return storage.Info{ return storage.Info{
Type: "mock", Type: "mock",
@@ -934,17 +928,17 @@ func setupDedupBackupEnv(
} }
} }
// runDedupSnapshot creates a snapshot with the given name, scans dataDir // runDedupSnapshot creates a "dedup" snapshot, scans dataDir into it,
// into it, completes it, and exports its metadata, returning the snapshot // completes it, and exports its metadata, returning the snapshot ID and
// ID and scan result. // scan result.
func runDedupSnapshot( func runDedupSnapshot(
ctx context.Context, t *testing.T, ctx context.Context, t *testing.T,
sm *snapshot.SnapshotManager, scanner *snapshot.Scanner, sm *snapshot.SnapshotManager, scanner *snapshot.Scanner,
hostname, name, dataDir, dbPath string, hostname, dataDir, dbPath string,
) (string, *snapshot.ScanResult) { ) (string, *snapshot.ScanResult) {
t.Helper() t.Helper()
id, err := sm.CreateSnapshotWithName(ctx, hostname, name, "v", "g") id, err := sm.CreateSnapshotWithName(ctx, hostname, "dedup", "v", "g")
require.NoError(t, err) require.NoError(t, err)
result, err := scanner.Scan(ctx, dataDir, id) result, err := scanner.Scan(ctx, dataDir, id)
@@ -986,15 +980,16 @@ func TestDedupOnlySnapshotRestores(t *testing.T) {
// First snapshot — uploads all blobs. // First snapshot — uploads all blobs.
_, r1 := runDedupSnapshot(ctx, t, sm, makeScanner(), _, r1 := runDedupSnapshot(ctx, t, sm, makeScanner(),
cfg.Hostname, "first", dataDir, dbPath) cfg.Hostname, dataDir, dbPath)
require.Positive(t, r1.BlobsCreated, require.Positive(t, r1.BlobsCreated,
"first snapshot should upload at least one blob") "first snapshot should upload at least one blob")
// Second snapshot — same data, every chunk dedups. Its own name gives // Second snapshot — same data, every chunk dedups. Sleep past the
// it a different snapshot ID without waiting for the one-second // second-precision timestamp so the snapshot IDs differ.
// timestamp in the ID to tick over. time.Sleep(1100 * time.Millisecond)
id2, r2 := runDedupSnapshot(ctx, t, sm, makeScanner(), id2, r2 := runDedupSnapshot(ctx, t, sm, makeScanner(),
cfg.Hostname, "second", dataDir, dbPath) cfg.Hostname, dataDir, dbPath)
require.Equal(t, 0, r2.BlobsCreated, require.Equal(t, 0, r2.BlobsCreated,
"second snapshot should upload zero new blobs (fully dedup'd)") "second snapshot should upload zero new blobs (fully dedup'd)")
+10 -5
View File
@@ -78,9 +78,10 @@ func backUpThenUnplug(
// TestFirstBackupCreatesDestinationDirectory checks that a first backup // TestFirstBackupCreatesDestinationDirectory checks that a first backup
// to a destination directory that does not exist yet creates it, and // to a destination directory that does not exist yet creates it, and
// that the destination can be listed afterwards. // that the destination can be listed afterwards.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestFirstBackupCreatesDestinationDirectory(t *testing.T) { func TestFirstBackupCreatesDestinationDirectory(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
storeDir := filepath.Join(t.TempDir(), "volume", "backup") storeDir := filepath.Join(t.TempDir(), "volume", "backup")
@@ -95,9 +96,10 @@ func TestFirstBackupCreatesDestinationDirectory(t *testing.T) {
// TestListSnapshotsWarnsWhenDestinationMissing checks that snapshot list // TestListSnapshotsWarnsWhenDestinationMissing checks that snapshot list
// warns and shows the local index alone, without reporting the local // warns and shows the local index alone, without reporting the local
// snapshot as missing from the destination. // snapshot as missing from the destination.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestListSnapshotsWarnsWhenDestinationMissing(t *testing.T) { func TestListSnapshotsWarnsWhenDestinationMissing(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
v, repos, out := backUpThenUnplug(ctx, t) v, repos, out := backUpThenUnplug(ctx, t)
@@ -114,9 +116,10 @@ func TestListSnapshotsWarnsWhenDestinationMissing(t *testing.T) {
// TestRemoveSnapshotWarnsWhenDestinationMissing checks that snapshot // TestRemoveSnapshotWarnsWhenDestinationMissing checks that snapshot
// remove warns that the metadata could not be removed from the // remove warns that the metadata could not be removed from the
// destination, instead of reporting that it was. // destination, instead of reporting that it was.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestRemoveSnapshotWarnsWhenDestinationMissing(t *testing.T) { func TestRemoveSnapshotWarnsWhenDestinationMissing(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
v, repos, out := backUpThenUnplug(ctx, t) v, repos, out := backUpThenUnplug(ctx, t)
@@ -135,9 +138,10 @@ func TestRemoveSnapshotWarnsWhenDestinationMissing(t *testing.T) {
// TestPruneKeepsLocalRecordsWhenDestinationMissing checks that prune // TestPruneKeepsLocalRecordsWhenDestinationMissing checks that prune
// fails on a destination it cannot list and deletes no local snapshot // fails on a destination it cannot list and deletes no local snapshot
// record. // record.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestPruneKeepsLocalRecordsWhenDestinationMissing(t *testing.T) { func TestPruneKeepsLocalRecordsWhenDestinationMissing(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
v, repos, _ := backUpThenUnplug(ctx, t) v, repos, _ := backUpThenUnplug(ctx, t)
@@ -154,9 +158,10 @@ func TestPruneKeepsLocalRecordsWhenDestinationMissing(t *testing.T) {
// TestPurgeSaysListingFailedOnceWhenDestinationMissing checks that // TestPurgeSaysListingFailedOnceWhenDestinationMissing checks that
// snapshot purge fails on a destination it cannot list, with an error // snapshot purge fails on a destination it cannot list, with an error
// that says "listing remote snapshots" once. // that says "listing remote snapshots" once.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestPurgeSaysListingFailedOnceWhenDestinationMissing(t *testing.T) { func TestPurgeSaysListingFailedOnceWhenDestinationMissing(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
v, _, _ := backUpThenUnplug(ctx, t) v, _, _ := backUpThenUnplug(ctx, t)
-64
View File
@@ -1,64 +0,0 @@
package vaultik_test
import (
"context"
"io/fs"
"os"
"path/filepath"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/snapshot"
)
// TestNukeRemoteLeavesNoFiles checks that remote nuke leaves no file
// under a file:// destination, including the `.partial` files uploads
// killed part-way leave next to a snapshot's metadata and next to where a
// blob would have been.
func TestNukeRemoteLeavesNoFiles(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background()
storeDir := filepath.Join(t.TempDir(), "store")
v, repos, _ := backUpToFileDestination(ctx, t, storeDir)
snapshots, err := repos.Snapshots.ListRecent(ctx, listRecentTestLimit)
require.NoError(t, err)
require.Len(t, snapshots, 1)
snapshotKey := snapshot.RemoteSnapshotKey(snapshots[0].ID.String())
hash := testBlobHashA
leftovers := []string{
filepath.Join(storeDir, "metadata", snapshotKey,
"db.zst.age-123456.partial"),
filepath.Join(storeDir, "blobs", hash[:2], hash[2:4],
hash+"-123456.partial"),
}
for _, leftover := range leftovers {
require.NoError(t, os.MkdirAll(filepath.Dir(leftover), 0o750))
require.NoError(t, os.WriteFile(leftover, []byte("half an upload"), 0o600))
}
require.NoError(t, v.NukeRemote(true))
var files []string
err = filepath.WalkDir(storeDir,
func(path string, entry fs.DirEntry, err error) error {
if err != nil {
return err
}
if !entry.IsDir() {
files = append(files, path)
}
return nil
})
require.NoError(t, err)
assert.Empty(t, files)
}
+2 -12
View File
@@ -24,9 +24,8 @@ var errNukeRequiresForce = errors.New(
const metadataDirName = "metadata" const metadataDirName = "metadata"
// NukeRemote deletes every snapshot's metadata and every blob from remote // NukeRemote deletes every snapshot's metadata and every blob from remote
// storage, along with any object an upload cut off part-way left under a // storage. After this returns successfully the bucket prefix is empty and
// temporary `.partial` name. After this returns successfully the bucket // the next backup starts from scratch.
// prefix is empty and the next backup starts from scratch.
// //
// Refuses to run unless force is true. The caller is responsible for // Refuses to run unless force is true. The caller is responsible for
// confirming with the user. // confirming with the user.
@@ -49,15 +48,6 @@ func (v *Vaultik) NukeRemote(force bool) error {
return fmt.Errorf("pruning blobs: %w", err) return fmt.Errorf("pruning blobs: %w", err)
} }
// The file and rclone listings skip `.partial` objects, so the two
// steps above never delete them.
for _, prefix := range []string{"metadata/", "blobs/"} {
err = v.Storage.DeletePartialUploads(v.ctx, prefix)
if err != nil {
return fmt.Errorf("deleting partial uploads: %w", err)
}
}
v.UI.Completef("Backup destination store is now empty.") v.UI.Completef("Backup destination store is now empty.")
return nil return nil
+4 -3
View File
@@ -14,9 +14,10 @@ import (
// the discarded-error bug: getTableCount for a table its query cannot // 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 // 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. // 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) { func TestTableCountForReportSurfacesReadFailure(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
@@ -42,7 +43,7 @@ func TestTableCountForReportSurfacesReadFailure(t *testing.T) {
assert.Nil(t, missing, "a failed read is unknown, not a count") assert.Nil(t, missing, "a failed read is unknown, not a count")
// The rendered count for a failed read must say unknown, never 0. // The rendered count for a failed read must say unknown, never 0.
assert.Equal(t, unknownText, countText(missing)) assert.Equal(t, countUnknown, countText(missing))
assert.NotEqual(t, "0", countText(missing)) assert.NotEqual(t, "0", countText(missing))
} }
@@ -57,7 +58,7 @@ func TestCountTextDistinguishesEmptyFromUnknown(t *testing.T) {
assert.Equal(t, "0", countText(&zero)) assert.Equal(t, "0", countText(&zero))
assert.Equal(t, "7", countText(&seven)) assert.Equal(t, "7", countText(&seven))
assert.Equal(t, unknownText, countText(nil)) assert.Equal(t, countUnknown, countText(nil))
} }
// TestCountDiffUnknownWhenEitherSideUnknown checks that a delta computed // TestCountDiffUnknownWhenEitherSideUnknown checks that a delta computed
-40
View File
@@ -1,40 +0,0 @@
package vaultik_test
import (
"encoding/json"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/vaultik"
)
// TestPruneBlobs_JSONDeletesWithoutAsking checks that prune with --json
// and without --force deletes an unreferenced blob without the
// confirmation prompt. Stdin is empty, so a prompt would read no answer
// and cancel, and its text would come before the JSON document.
func TestPruneBlobs_JSONDeletesWithoutAsking(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
env := newListEnv(t)
// The store holds no manifest, so nothing references this blob.
addBlob(t, env.store.testStorer, testBlobHashA)
err := env.v.PruneBlobs(&vaultik.PruneOptions{JSON: true})
require.NoError(t, err)
blobKey := "blobs/" + testBlobHashA[:2] + "/" + testBlobHashA[2:4] +
"/" + testBlobHashA
assert.False(t, env.store.hasKey(blobKey),
"the unreferenced blob must be deleted")
var result vaultik.PruneBlobsResult
require.NoError(t, json.Unmarshal(env.stdout.Bytes(), &result),
"stdout must hold only the JSON document, got:\n%s",
env.stdout.String())
assert.Equal(t, 1, result.BlobsDeleted)
}
-69
View File
@@ -4,7 +4,6 @@ import (
"bytes" "bytes"
"context" "context"
"encoding/json" "encoding/json"
"strings"
"testing" "testing"
"time" "time"
@@ -70,74 +69,6 @@ func TestRemoteInfo_UnreadableManifestLeavesOrphansUnknown(t *testing.T) {
assert.Equal(t, []any{unreadableKey}, doc["unreadable_manifests"]) assert.Equal(t, []any{unreadableKey}, doc["unreadable_manifests"])
} }
// TestRemoteInfo_UnreadableManifestLeavesSnapshotBlobsUnknown checks
// that the row of a snapshot whose manifest cannot be read gives its
// blob count and blob size as unknown in the table and as null in
// --json, not as 0. A directory without a manifest still shows 0: the
// orphan figures count its blobs as orphaned, so it references none.
func TestRemoteInfo_UnreadableManifestLeavesSnapshotBlobsUnknown(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
env := newListEnv(t)
readableKey := env.addRemote(t, listRemoteID,
time.Date(2026, 3, 2, 0, 0, 0, 0, time.UTC))
unreadableKey := snapshot.RemoteSnapshotKey(listLocalID)
require.NoError(t, env.store.Put(context.Background(),
"metadata/"+unreadableKey+"/manifest.json.zst",
bytes.NewReader([]byte("not a valid manifest"))))
noManifestKey := snapshot.RemoteSnapshotKey("testhost_home_2026-03-03T10:00:00Z")
require.NoError(t, env.store.Put(context.Background(),
"metadata/"+noManifestKey+"/db.zst.age",
bytes.NewReader([]byte("not a valid database"))))
require.NoError(t, env.v.RemoteInfo(false))
// The table truncates the remote key, so a row is found by a prefix.
wantUnknown := map[string]int{readableKey: 0, unreadableKey: 2, noManifestKey: 0}
for key, want := range wantUnknown {
var row string
for line := range strings.SplitSeq(env.stdout.String(), "\n") {
if strings.HasPrefix(line, key[:16]) {
row = line
}
}
require.NotEmpty(t, row, "no table row for %s", key)
assert.Equal(t, want, strings.Count(row, "unknown"), "row: %q", row)
}
env.stdout.Reset()
require.NoError(t, env.v.RemoteInfo(true))
var doc struct {
Snapshots []map[string]any `json:"snapshots"`
}
require.NoError(t, json.Unmarshal(env.stdout.Bytes(), &doc))
require.Len(t, doc.Snapshots, len(wantUnknown))
for _, entry := range doc.Snapshots {
switch entry["snapshot_id"] {
case readableKey:
assert.InDelta(t, 1, entry["blob_count"], 0)
assert.InDelta(t, fiveMegabytes, entry["blobs_size"], 0)
case noManifestKey:
assert.InDelta(t, 0, entry["blob_count"], 0)
assert.InDelta(t, 0, entry["blobs_size"], 0)
default:
assert.Equal(t, unreadableKey, entry["snapshot_id"])
assert.Contains(t, entry, "blob_count")
assert.Nil(t, entry["blob_count"])
assert.Contains(t, entry, "blobs_size")
assert.Nil(t, entry["blobs_size"])
}
}
}
// TestRemoteInfo_SkipsNonConformingMetadataName checks that a directory // TestRemoteInfo_SkipsNonConformingMetadataName checks that a directory
// under metadata/ whose name is not a remote key is left out of the // under metadata/ whose name is not a remote key is left out of the
// report, and that the orphan figures are unknown when it holds a // report, and that the orphan figures are unknown when it holds a
-6
View File
@@ -125,12 +125,6 @@ func (s *testStorer) ListStream(
return ch return ch
} }
// DeletePartialUploads has nothing to remove: Put stores each object
// under its key at once.
func (s *testStorer) DeletePartialUploads(_ context.Context, _ string) error {
return nil
}
func (s *testStorer) Info() storage.Info { func (s *testStorer) Info() storage.Info {
return storage.Info{ return storage.Info{
Type: testLabel, Type: testLabel,
+2 -1
View File
@@ -164,9 +164,10 @@ func scratchEntries(t *testing.T, dir string) []string {
// restore while a blob download is in progress. The download fails only // restore while a blob download is in progress. The download fails only
// because of the cancel, so Restore must return context.Canceled without // because of the cancel, so Restore must return context.Canceled without
// reporting the file that needs the blob as failed. // reporting the file that needs the blob as failed.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestRestoreSkipErrorsCancelDuringBlobDownload(t *testing.T) { func TestRestoreSkipErrorsCancelDuringBlobDownload(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
+5 -6
View File
@@ -45,10 +45,9 @@ type missingBlobBackup struct {
// after one blob of a two-blob snapshot was deleted. Every file stored in // 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 // that blob must be reported as failed and left absent, every other file
// must be restored intact, and Restore must still return an error. // 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) { func TestRestoreSkipErrorsSkipsFilesOfMissingBlob(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
backup := backupThenDeleteOneBlob(ctx, t) backup := backupThenDeleteOneBlob(ctx, t)
@@ -86,10 +85,9 @@ func TestRestoreSkipErrorsSkipsFilesOfMissingBlob(t *testing.T) {
// TestRestoreMissingBlobAbortsWithoutSkipErrors checks that a deleted blob // TestRestoreMissingBlobAbortsWithoutSkipErrors checks that a deleted blob
// still ends the restore with an error when SkipErrors is not set. // 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) { func TestRestoreMissingBlobAbortsWithoutSkipErrors(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background() ctx := context.Background()
backup := backupThenDeleteOneBlob(ctx, t) backup := backupThenDeleteOneBlob(ctx, t)
@@ -109,6 +107,7 @@ func backupThenDeleteOneBlob(
ctx context.Context, t *testing.T, ctx context.Context, t *testing.T,
) *missingBlobBackup { ) *missingBlobBackup {
t.Helper() t.Helper()
log.Initialize(log.Config{})
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
+2 -1
View File
@@ -18,9 +18,10 @@ import (
// A file rewritten with its size unchanged and a new mtime in the same // 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 // second as the mtime the index holds must still be backed up. See
// https://git.eeqj.de/sneak/vaultik/issues/226. // https://git.eeqj.de/sneak/vaultik/issues/226.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupOfSameSecondRewriteRestoresNewContent(t *testing.T) { func TestBackupOfSameSecondRewriteRestoresNewContent(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
+6 -16
View File
@@ -185,7 +185,6 @@ type snapshotStats struct {
totalBlobs int totalBlobs int
totalBytesSkipped int64 totalBytesSkipped int64
totalFilesSkipped int totalFilesSkipped int
totalFilesFailed int
totalFilesDeleted int totalFilesDeleted int
totalBytesDeleted int64 totalBytesDeleted int64
totalBytesUploaded int64 totalBytesUploaded int64
@@ -316,7 +315,6 @@ func (v *Vaultik) scanAllDirectories(
stats.totalChunks += result.ChunksCreated stats.totalChunks += result.ChunksCreated
stats.totalBlobs += result.BlobsCreated stats.totalBlobs += result.BlobsCreated
stats.totalFilesSkipped += result.FilesSkipped stats.totalFilesSkipped += result.FilesSkipped
stats.totalFilesFailed += result.FilesFailed
stats.totalBytesSkipped += result.BytesSkipped stats.totalBytesSkipped += result.BytesSkipped
stats.totalFilesDeleted += result.FilesDeleted stats.totalFilesDeleted += result.FilesDeleted
stats.totalBytesDeleted += result.BytesDeleted stats.totalBytesDeleted += result.BytesDeleted
@@ -328,7 +326,6 @@ func (v *Vaultik) scanAllDirectories(
"path", dir, "path", dir,
"files", result.FilesScanned, "files", result.FilesScanned,
"files_skipped", result.FilesSkipped, "files_skipped", result.FilesSkipped,
"files_failed", result.FilesFailed,
"bytes", result.BytesScanned, "bytes", result.BytesScanned,
"bytes_skipped", result.BytesSkipped, "bytes_skipped", result.BytesSkipped,
"chunks", result.ChunksCreated, "chunks", result.ChunksCreated,
@@ -362,11 +359,9 @@ func (v *Vaultik) finalizeSnapshotMetadata(
return fmt.Errorf("getting snapshot blob sizes: %w", err) return fmt.Errorf("getting snapshot blob sizes: %w", err)
} }
// file_count and total_size leave out the files that could not be
// stored; stats.totalBytes already does.
extStats := snapshot.ExtendedBackupStats{ extStats := snapshot.ExtendedBackupStats{
BackupStats: snapshot.BackupStats{ BackupStats: snapshot.BackupStats{
FilesScanned: stats.totalFiles - stats.totalFilesFailed, FilesScanned: stats.totalFiles,
TotalSize: stats.totalBytes + stats.totalBytesSkipped, TotalSize: stats.totalBytes + stats.totalBytesSkipped,
ChunksCreated: stats.totalChunks, ChunksCreated: stats.totalChunks,
BlobsCreated: stats.totalBlobs, BlobsCreated: stats.totalBlobs,
@@ -414,8 +409,7 @@ func (v *Vaultik) printSnapshotSummary(
snapshotID string, startTime time.Time, stats *snapshotStats, snapshotID string, startTime time.Time, stats *snapshotStats,
) { ) {
snapshotDuration := time.Since(startTime) snapshotDuration := time.Since(startTime)
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped - totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
stats.totalFilesFailed
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
var compressionRatio float64 var compressionRatio float64
@@ -432,10 +426,6 @@ func (v *Vaultik) printSnapshotSummary(
v.UI.Count(stats.totalFiles), v.UI.Count(stats.totalFiles),
v.UI.Count(totalFilesChanged), v.UI.Count(totalFilesChanged),
v.UI.Count(stats.totalFilesSkipped)) v.UI.Count(stats.totalFilesSkipped))
if stats.totalFilesFailed > 0 {
filesMsg += fmt.Sprintf(", %s failed", v.UI.Count(stats.totalFilesFailed))
}
if stats.totalFilesDeleted > 0 { if stats.totalFilesDeleted > 0 {
filesMsg += fmt.Sprintf(", %s deleted", v.UI.Count(stats.totalFilesDeleted)) filesMsg += fmt.Sprintf(", %s deleted", v.UI.Count(stats.totalFilesDeleted))
} }
@@ -1766,9 +1756,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
return result, nil return result, nil
} }
// unknownText is what a count or size reads as when it could not be // countUnknown is what a count reads as when its query could not be run,
// determined, distinct from "0", which is a real zero. // distinct from "0", which means the table really was empty.
const unknownText = "unknown" const countUnknown = "unknown"
// tableCountForReport returns the row count of a table for the prune // tableCountForReport returns the row count of a table for the prune
// summary, or nil if the count could not be read. A read failure is // summary, or nil if the count could not be read. A read failure is
@@ -1805,7 +1795,7 @@ func countDiff(before, after *int64) *int64 {
// one that could not be queried. // one that could not be queried.
func countText(count *int64) string { func countText(count *int64) string {
if count == nil { if count == nil {
return unknownText return countUnknown
} }
return strconv.FormatInt(*count, 10) return strconv.FormatInt(*count, 10)
@@ -1,123 +0,0 @@
package vaultik_test
import (
"bytes"
"context"
"fmt"
"os"
"path/filepath"
"testing"
"github.com/spf13/afero"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/config"
"sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/storage"
"sneak.berlin/go/vaultik/internal/ui"
"sneak.berlin/go/vaultik/internal/vaultik"
)
// openFailFs is the real filesystem, except that opening path fails
// with err. Phase 1 of a backup only lstats a file, so it still counts
// path; phase 2 is the first to open it.
type openFailFs struct {
afero.OsFs
path string
err error
}
//nolint:ireturn // afero.Fs.Open is defined to return the interface.
func (f *openFailFs) Open(name string) (afero.File, error) {
if name == f.path {
return nil, &os.PathError{Op: "open", Path: name, Err: f.err}
}
return f.OsFs.Open(name)
}
// A file that phase 2 cannot open is reported as failed, not as
// unchanged, and neither the summary's data total nor the snapshots row
// counts it. See https://git.eeqj.de/sneak/vaultik/issues/280.
func TestSnapshotSummaryCountsFileNotStoredAsFailed(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
tests := []struct {
name string
openErr error
skipErrors bool
}{
// What a normal user gets opening a file with mode 000.
{name: "unopenable under skip-errors",
openErr: os.ErrPermission, skipErrors: true},
// What opening a file removed after phase 1 gives.
{name: "removed between the phases",
openErr: os.ErrNotExist, skipErrors: false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
ctx := context.Background()
// The scan walks the source path with symlinks resolved, so
// failedPath must be spelled the same way to match.
tempDir, err := filepath.EvalSymlinks(t.TempDir())
require.NoError(t, err)
srcDir := filepath.Join(tempDir, "src")
failedPath := filepath.Join(srcDir, "failed.txt")
storedContent := []byte("this file is backed up")
storedSize := int64(len(storedContent))
fs := &openFailFs{path: failedPath, err: tt.openErr}
require.NoError(t, fs.MkdirAll(srcDir, 0o755))
require.NoError(t, afero.WriteFile(fs,
filepath.Join(srcDir, "stored.txt"), storedContent, 0o644))
require.NoError(t, afero.WriteFile(fs,
failedPath, []byte("this file cannot be opened"), 0o644))
cfg := faultTestConfig()
cfg.IndexPath = filepath.Join(tempDir, "index.sqlite")
cfg.Snapshots = map[string]config.SnapshotConfig{
"src": {Paths: []string{srcDir}},
}
store, err := storage.NewFileStorer(filepath.Join(tempDir, "remote"))
require.NoError(t, err)
db, err := database.New(ctx, cfg.IndexPath)
require.NoError(t, err)
t.Cleanup(func() { _ = db.Close() })
repos := database.NewRepositories(db)
out := &bytes.Buffer{}
v := newBackupVaultik(ctx, cfg, store, repos, db, fs)
v.UI = ui.NewWithColor(out, false)
require.NoError(t, v.CreateSnapshot(&vaultik.SnapshotCreateOptions{
SkipErrors: tt.skipErrors,
Snapshots: []string{"src"},
}))
summary := out.String()
assert.Contains(t, summary,
"Files: 2 examined, 1 backed up, 0 unchanged, 1 failed.")
assert.Contains(t, summary,
fmt.Sprintf("Data: %s total (%s backed up).",
v.UI.Size(storedSize), v.UI.Size(storedSize)))
snap, err := repos.Snapshots.GetByID(ctx,
localSnapshotID(ctx, t, repos, "src"))
require.NoError(t, err)
require.NotNil(t, snap)
assert.Equal(t, int64(1), snap.FileCount)
assert.Equal(t, storedSize, snap.TotalSize)
})
}
}
@@ -1,75 +0,0 @@
package vaultik_test
import (
"context"
"fmt"
"path/filepath"
"testing"
"github.com/spf13/afero"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"sneak.berlin/go/vaultik/internal/config"
"sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/log"
"sneak.berlin/go/vaultik/internal/storage"
"sneak.berlin/go/vaultik/internal/vaultik"
)
// --skip-errors skips only a file that cannot be opened or read. A local
// index error while a directory is recorded stops the backup, and the
// snapshot is not recorded as complete. See
// https://git.eeqj.de/sneak/vaultik/issues/284.
func TestSkipErrorsBackupStopsOnIndexErrorRecordingDirectory(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
ctx := context.Background()
// The scan walks the source path with symlinks resolved, so dirPath
// must be spelled the same way to match the row the scan inserts.
tempDir, err := filepath.EvalSymlinks(t.TempDir())
require.NoError(t, err)
srcDir := filepath.Join(tempDir, "src")
dirPath := filepath.Join(srcDir, "dir")
fs := afero.NewOsFs()
require.NoError(t, fs.MkdirAll(dirPath, 0o755))
require.NoError(t, afero.WriteFile(fs,
filepath.Join(dirPath, "file.txt"), []byte("file content"), 0o644))
cfg := faultTestConfig()
cfg.IndexPath = filepath.Join(tempDir, "index.sqlite")
cfg.Snapshots = map[string]config.SnapshotConfig{
"tree": {Paths: []string{srcDir}},
}
store, err := storage.NewFileStorer(filepath.Join(tempDir, "remote"))
require.NoError(t, err)
db, err := database.New(ctx, cfg.IndexPath)
require.NoError(t, err)
t.Cleanup(func() { _ = db.Close() })
// The local index refuses the directory's files row.
_, err = db.Conn().ExecContext(ctx, fmt.Sprintf(`
CREATE TRIGGER refuse_directory BEFORE INSERT ON files
WHEN NEW.path = '%s'
BEGIN SELECT RAISE(ABORT, 'simulated index error'); END`, dirPath))
require.NoError(t, err)
repos := database.NewRepositories(db)
v := newBackupVaultik(ctx, cfg, store, repos, db, fs)
err = v.CreateSnapshot(&vaultik.SnapshotCreateOptions{
SkipErrors: true,
Snapshots: []string{"tree"},
})
require.ErrorContains(t, err, "simulated index error")
snapshots, err := repos.Snapshots.ListRecent(ctx, listRecentTestLimit)
require.NoError(t, err)
require.Len(t, snapshots, 1)
assert.Nil(t, snapshots[0].CompletedAt)
}
+8 -9
View File
@@ -54,6 +54,8 @@ type summaryEnv struct {
func newSummaryEnv(t *testing.T) *summaryEnv { func newSummaryEnv(t *testing.T) *summaryEnv {
t.Helper() t.Helper()
log.Initialize(log.Config{})
fs := afero.NewOsFs() fs := afero.NewOsFs()
tempDir := t.TempDir() tempDir := t.TempDir()
srcDir := filepath.Join(tempDir, "src") srcDir := filepath.Join(tempDir, "src")
@@ -191,10 +193,9 @@ func (e *summaryEnv) dataLine(total, backedUp int64) string {
// A first backup stores copy.bin's chunks while backing up a.bin, so // 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 // copy.bin's chunks are deduplicated within the run. Each file and byte
// is still counted once. // is still counted once.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestSnapshotSummaryFirstRun(t *testing.T) { func TestSnapshotSummaryFirstRun(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
env := newSummaryEnv(t) env := newSummaryEnv(t)
summary := env.backUp(t, "first", false) summary := env.backUp(t, "first", false)
@@ -218,10 +219,9 @@ func TestSnapshotSummaryFirstRun(t *testing.T) {
// An incremental backup where a.bin's mtime changed but its content did // 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 // not: a.bin is backed up again and every one of its chunks is already
// stored. // stored.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) { func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
env := newSummaryEnv(t) env := newSummaryEnv(t)
env.backUp(t, "first", false) env.backUp(t, "first", false)
@@ -251,10 +251,9 @@ func TestSnapshotSummaryIncrementalRunWithDeduplicatedChunks(t *testing.T) {
// Under --cron the progress reporter is off; the upload figures must // Under --cron the progress reporter is off; the upload figures must
// still reach the summary and the snapshots row. The snapshot has two // still reach the summary and the snapshots row. The snapshot has two
// paths, each backed up by its own scan. // paths, each backed up by its own scan.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestSnapshotSummaryCronRunRecordsUploads(t *testing.T) { func TestSnapshotSummaryCronRunRecordsUploads(t *testing.T) {
log.Initialize(log.Config{})
t.Parallel()
env := newSummaryEnv(t) env := newSummaryEnv(t)
summary := env.backUp(t, "split", true) summary := env.backUp(t, "split", true)
+2 -1
View File
@@ -18,9 +18,10 @@ import (
// A backup without --cron runs the progress reporter while one scanner // A backup without --cron runs the progress reporter while one scanner
// scans each path of the snapshot in turn. See // scans each path of the snapshot in turn. See
// https://git.eeqj.de/sneak/vaultik/issues/253. // https://git.eeqj.de/sneak/vaultik/issues/253.
//
//nolint:paralleltest // installs the global logger via log.Initialize
func TestBackupWithoutCronOfTwoPathSnapshotRestoresBothPaths(t *testing.T) { func TestBackupWithoutCronOfTwoPathSnapshotRestoresBothPaths(t *testing.T) {
log.Initialize(log.Config{}) log.Initialize(log.Config{})
t.Parallel()
const snapshotName = "data" const snapshotName = "data"