8 Commits
Author SHA1 Message Date
clawbot 79a73fa122 Count a file a backup could not store as failed (closes #280)
check / check (push) Waiting to run
A file that phase 1 of a backup counted and phase 2 could not open,
because it was unreadable under --skip-errors or removed in between,
was added to the unchanged count while its size stayed in
BytesScanned. The summary showed it as unchanged with its bytes backed
up, and the snapshots row's file_count and total_size included it.

The scanner now counts such a file in FilesFailed and takes its size
out of BytesScanned. The summary's files line adds "N failed", and
file_count leaves the file out.

A directory phase 2 cannot record is not counted as failed, since
phase 1 counts no directories; that case has no test.

Model: opus-5-5
2026-10-08 10:12:10 +02:00
clawbot 0d0368df81 Report unknown blob figures for a snapshot whose manifest cannot be read (closes #272)
check / check (push) Waiting to run
When remote info could not read a snapshot's manifest, the orphan
figures were unknown but the snapshot's row still gave 0 blobs and 0 B,
in the table and in --json. The row's blob count and blob size are now
unknown, and null in --json.

A directory with no manifest, as an interrupted backup leaves, still
shows 0: the orphan figures count its blobs as orphaned, so it
references none.

The constant holding the "unknown" text is renamed from countUnknown to
unknownText, since it now also stands for a size.

Model: opus-5-5
2026-10-08 08:46:11 +02:00
clawbot c06c4d201b Keep the progress line of a multi-path snapshot within 100% (closes #271)
check / check (push) Waiting to run
One progress reporter spans every path of a snapshot. Its counts of
files and bytes processed run across all the paths, while the totals
they were divided by were reset to each path's own, 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. SetTotalSize is renamed AddTotalSize because it
now adds.

The issue says the ETA went negative; it was computed negative and
printed as `unknown`.

Model: opus-5-5
2026-10-08 07:29:11 +02:00
clawbot 70f008a21a Document that --json skips the confirmation prompt (closes #268)
check / check (push) Waiting to run
snapshot remove and prune delete without asking under --json, because a
prompt on stdout would break the JSON document. Their --json help and
README entries named only --force as skipping the prompt. Both now say
--json skips it too. One test checks the --json help of both commands;
another runs prune with --json, no --force and empty stdin, and checks
that the unreferenced blob is deleted and stdout holds only the JSON
document.

Model: opus-5-5
2026-10-08 04:59:30 +02:00
clawbot fce253fa39 Give a second snapshot create in the same second its own ID (closes #270)
check / check (push) Waiting to run
The timestamp in a snapshot ID is in whole seconds, so a `snapshot
create` that started in the same second as the previous run of that
snapshot name got the same ID, and inserting its row failed with
`UNIQUE constraint failed: snapshots.id`. CreateSnapshotWithName now
looks the ID up in the local index first and, if it is taken, waits a
second and takes a new timestamp. The ID format is unchanged.

Only the local index is checked. That is where the insert fails, and
the process lock serializes runs, so nothing takes the ID between the
lookup and the insert.

Model: opus-5-5
2026-10-08 03:29:07 +02:00
clawbot e459a66099 Stop a backup on a symlink whose target cannot be read (closes #269)
check / check (push) Waiting to run
When readlink failed, the scanner logged at debug level and left the
symlink out of the snapshot, and the run reported success even without
--skip-errors. The error now goes through the same handling as any
other entry the walk cannot read: the run aborts, or with --skip-errors
the symlink is skipped with the usual "Failed to access" error line.

A symlink removed between the walk's lstat and the readlink also
aborts the run, as an entry that vanishes during the walk already does.

Model: opus-5-5
2026-10-08 00:12:13 +02:00
clawbot 32c4a46b49 Write rclone uploads under a temporary name and move them into place (closes #266)
check / check (push) Waiting to run
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 that the
next backup trusted.

On every remote with a server-side move, an object is now written under
a name ending in `.partial` and moved onto its key with rclone's
operations.Move, which removes an object already at the key first;
drive, dropbox and others will not move onto one. Listings skip
`.partial` names. Rclone's own copy also requires the PartialUploads
flag; this does not, because hdfs shows a file while it is written
without setting it. Remotes without a move are written in place.

Model: opus-5-5
2026-10-07 22:59:26 +02:00
clawbot b5389a62b5 Exit 130 and say so when a command is interrupted (closes #267)
check / check (push) Waiting to run
Ctrl-C or SIGTERM during snapshot create, restore or verify exited 0
with no error line (1, also silent, under snapshot verify --json), so
an unfinished --cron backup looked like a success. RunOperation now
records whether op returned while the Vaultik context was still live;
a run where it had not by the time RunWithApp returned is interrupted,
whatever op returned. It used to look for context.Canceled in op's
error, which verify --json does not return. Entry prints "interrupted
before the command finished" on stderr for it and returns 130, under
--cron and --json too.

SIGTERM also gives 130, as the issue asks, not 143.
The test does not cover an op still running when the 30s shutdown
timeout ends.

Model: opus-5-5
2026-10-07 21:59:27 +02:00
24 changed files with 1203 additions and 176 deletions
+16 -3
View File
@@ -359,7 +359,8 @@ on the destination in one go, use `vaultik remote nuke --force`.
* `--local-only`: Skip remote cleanup; only touch the local index
* `--dry-run`: Show what would be deleted without deleting
* `--force`: Skip confirmation prompt
* `--json`: Output result as JSON
* `--json`: Output result as JSON. Also skips the confirmation prompt, as
`--force` does.
**`snapshot restore`**: Restore files from a backup snapshot.
* Requires `VAULTIK_AGE_SECRET_KEY` environment variable
@@ -383,7 +384,8 @@ manifests — network cost scales with the number of snapshots. `snapshot
create --prune` runs the same cleanup automatically; this is the
manual entry point for the same work.
* `--force`: Skip confirmation prompt
* `--json`: Output stats as JSON
* `--json`: Output stats as JSON. Also skips the confirmation prompt, as
`--force` does.
**`info`**: Display system configuration, storage settings, encryption
recipients, and local database statistics.
@@ -396,7 +398,8 @@ key is skipped with a warning and is not printed. If a listed
orphaned blob figures are reported as unknown; `--json` gives them as
`null`, lists the remote key of each unreadable manifest in
`unreadable_manifests` and counts the manifests under skipped names in
`skipped_manifest_count`.
`skipped_manifest_count`. An unreadable manifest also leaves its
snapshot's blob count and blob size unknown, `null` in `--json`.
* `--json`: Output as JSON
**`remote nuke`**: Delete every snapshot's metadata and every blob from the
@@ -428,6 +431,16 @@ a local or mounted filesystem. Useful for testing or backing up to a NAS.
**Rclone** (`rclone://remote/path`): Uses rclone's 70+ supported cloud
providers. Requires rclone to be configured separately (`rclone config`).
An upload cut off part-way leaves nothing under the object's name on S3, which
shows an object only once its upload has completed, and on the local filesystem
backend, which writes a temporary file and renames it into place. Rclone
remotes with a server-side move (local and sftp among them) are written under a
temporary name ending in `.partial` and moved into place. Rclone remotes without
one are written in place: where such a remote shows a file while it is still
being written, a killed upload can leave a truncated object under its name,
which a later backup takes for complete. A leftover `.partial` file is ignored
and can be deleted.
Legacy S3 configuration via `s3.*` fields (endpoint, bucket, prefix, etc.) is
still supported for backward compatibility. `storage_url` takes precedence if
both are set.
+63
View File
@@ -22,6 +22,69 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps
- 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
+26 -28
View File
@@ -207,7 +207,7 @@ func RunApp(ctx context.Context, app *fx.App) error {
var errReported = errors.New("operation failed")
// errInterrupted marks an operation that SIGINT or SIGTERM stopped
// before it finished. Entry shows it and exits with exitCodeInterrupted.
// 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
@@ -225,16 +225,20 @@ var errInterrupted = errors.New("interrupted before the command finished")
// op's cleanup (removing decrypted scratch files) runs before the
// process exits; the wait is bounded by shutdownTimeout. report is
// called with a failure so the caller can show it to the user before
// it becomes errReported. An interrupted op is not reported, whatever
// it returned; RunOperation returns errInterrupted instead.
// it becomes errReported.
//
// The run counts as interrupted unless op returned, without an
// interrupt having cancelled it, before RunWithApp returned. An
// interrupted op is not reported, whatever it returned; RunOperation
// returns errInterrupted instead.
func RunOperation(
ctx context.Context, opts AppOptions,
op func(v *vaultik.Vaultik) error, report func(err error),
) error {
var (
mu sync.Mutex
failed bool
interrupted bool
mu sync.Mutex
finished bool // op returned before any interrupt cancelled it
failed bool // op finished with an error
)
opts.Invokes = append(opts.Invokes,
@@ -247,21 +251,19 @@ func RunOperation(
err := op(v)
// Only stop, called from OnStop below, cancels the
// Vaultik context; before op has returned, that
// happens only on an interrupt. Check the context,
// not err: an interrupted op need not return
// 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).
switch {
case v.Context().Err() != nil:
mu.Lock()
interrupted = true
mu.Unlock()
case err != nil:
report(err)
if v.Context().Err() == nil {
if err != nil {
report(err)
}
mu.Lock()
failed = true
finished = true
failed = err != nil
mu.Unlock()
}
@@ -282,12 +284,6 @@ func RunOperation(
if !stop(ctx) {
log.Warn("Shutdown timed out before the operation " +
"finished; decrypted temporary files may remain")
// op has not returned, so the app is stopping on
// an interrupt.
mu.Lock()
interrupted = true
mu.Unlock()
}
return nil
@@ -300,15 +296,17 @@ func RunOperation(
return err
}
// The goroutine sets failed or interrupted before triggering the
// shutdown that lets RunWithApp return, so the write is in place by
// the time we read it. If OnStop timed out, the goroutine has not
// got that far, and OnStop has set interrupted itself.
// RunWithApp returns only after the app was asked to stop, either by
// an interrupt or by the goroutine's Shutdown call. When op finished
// without being cancelled, the goroutine set finished before that
// call. So if finished is unset here, an interrupt stopped the app,
// and op either returned after it was cancelled or is still running
// because the shutdown timed out.
mu.Lock()
defer mu.Unlock()
switch {
case interrupted:
case !finished:
return errInterrupted
case failed:
return errReported
+24
View File
@@ -0,0 +1,24 @@
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())
}
}
+2 -1
View File
@@ -57,7 +57,8 @@ referenced.`,
}
cmd.Flags().BoolVar(&opts.Force, "force", false, "Skip confirmation prompt")
cmd.Flags().BoolVar(&opts.JSON, "json", false, "Output pruning stats as JSON")
cmd.Flags().BoolVar(&opts.JSON, "json", false,
"Output pruning stats as JSON; skips the confirmation prompt, as --force does")
return cmd
}
+2 -1
View File
@@ -280,7 +280,8 @@ nuke --force' — it is the single supported entry point for that.`,
cmd.Flags().BoolVarP(&opts.Force, "force", "f", false, "Skip confirmation prompt")
cmd.Flags().BoolVar(&opts.DryRun, "dry-run", false,
"Show what would be removed without removing")
cmd.Flags().BoolVar(&opts.JSON, "json", false, "Output result as JSON")
cmd.Flags().BoolVar(&opts.JSON, "json", false,
"Output result as JSON; skips the confirmation prompt, as --force does")
cmd.Flags().BoolVar(&opts.LocalOnly, "local-only", false,
"Skip remote cleanup; only touch the local index")
+69 -53
View File
@@ -56,23 +56,27 @@ const (
// ProgressStats holds atomic counters for progress tracking
type ProgressStats struct {
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
FilesProcessed atomic.Int64 // Files actually processed in phase 2
FilesSkipped atomic.Int64 // Files skipped due to no changes
BytesScanned atomic.Int64 // Bytes from new/changed files only
BytesSkipped atomic.Int64 // Bytes from unchanged files
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
ChunksCreated atomic.Int64
BlobsCreated atomic.Int64
BlobsUploaded atomic.Int64
BytesUploaded atomic.Int64
CurrentFile atomic.Value // stores string
TotalSize atomic.Int64 // Total size to process (set after scan phase)
TotalFiles atomic.Int64 // Total files to process in phase 2
ProcessStartTime atomic.Value // stores time.Time when processing starts
StartTime time.Time
mu sync.RWMutex
lastDetailTime time.Time
FilesScanned atomic.Int64 // Total files seen during scan (includes skipped)
FilesProcessed atomic.Int64 // Files actually processed in phase 2
FilesSkipped atomic.Int64 // Files skipped due to no changes
BytesScanned atomic.Int64 // Bytes from new/changed files only
BytesSkipped atomic.Int64 // Bytes from unchanged files
BytesProcessed atomic.Int64 // Actual bytes processed (for ETA calculation)
ChunksCreated atomic.Int64
BlobsCreated atomic.Int64
BlobsUploaded atomic.Int64
BytesUploaded atomic.Int64
CurrentFile atomic.Value // stores string
TotalSize atomic.Int64 // Size to process in the paths scanned so far
TotalFiles atomic.Int64 // Files to process in the paths scanned so far
StartTime 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
CurrentUpload atomic.Value // stores *UploadInfo
@@ -148,10 +152,36 @@ func (pr *ProgressReporter) GetStats() *ProgressStats {
return pr.stats
}
// SetTotalSize sets the total size to process (after scan phase)
func (pr *ProgressReporter) SetTotalSize(size int64) {
pr.stats.TotalSize.Store(size)
pr.stats.ProcessStartTime.Store(time.Now().UTC())
// AddTotalSize adds the size one path of the snapshot has to process to
// the total, once that path's scan phase is done, and starts measuring
// the processing rate again. The processed counts run across every path,
// so the total does too. Restarting the rate keeps a path's scan phase,
// 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
@@ -357,19 +387,12 @@ func (pr *ProgressReporter) printSummaryStatus() {
// Calculate ETA if we have total size and are processing
etaStr := ""
if totalSize > 0 && bytesProcessed > 0 {
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
if ok && !processStart.IsZero() {
processElapsed := time.Since(processStart)
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)
}
}
processRate := pr.stats.processRate()
if totalSize > 0 && processRate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / processRate
eta := time.Duration(remainingSeconds * float64(time.Second))
etaStr = " | ETA: " + formatDuration(eta)
}
rate := float64(bytesScanned+bytesSkipped) / elapsed.Seconds()
@@ -421,25 +444,18 @@ func (pr *ProgressReporter) printDetailedStatus() {
log.Info("Elapsed time", "duration", formatDuration(elapsed))
// Calculate and show ETA if we have data
if totalSize > 0 && bytesProcessed > 0 {
processStart, ok := pr.stats.ProcessStartTime.Load().(time.Time)
if ok && !processStart.IsZero() {
processElapsed := time.Since(processStart)
processRate := float64(bytesProcessed) / processElapsed.Seconds()
if processRate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / processRate
eta := time.Duration(remainingSeconds * float64(time.Second))
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
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))
}
}
processRate := pr.stats.processRate()
if totalSize > 0 && processRate > 0 {
remainingBytes := totalSize - bytesProcessed
remainingSeconds := float64(remainingBytes) / processRate
eta := time.Duration(remainingSeconds * float64(time.Second))
percentComplete := float64(bytesProcessed) / float64(totalSize) * percentScale
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",
+54
View File
@@ -0,0 +1,54 @@
//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
@@ -0,0 +1,94 @@
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())
}
}
+42 -24
View File
@@ -133,10 +133,13 @@ type ScannerConfig struct {
// ScanResult contains the results of a scan operation. Files and bytes
// are counted per file: BytesScanned is the size of the new and changed
// files, BytesSkipped that of the unchanged ones.
// files, BytesSkipped that of the unchanged ones. FilesFailed counts the
// new and changed files that phase 2 could not store; FilesScanned
// includes them and BytesScanned does not.
type ScanResult struct {
FilesScanned int
FilesSkipped int
FilesFailed int
FilesDeleted int
BytesScanned int64
BytesSkipped int64
@@ -407,8 +410,8 @@ func (s *Scanner) summarizeScanPhase(
}
if s.progress != nil {
s.progress.SetTotalSize(totalSizeToProcess)
s.progress.GetStats().TotalFiles.Store(int64(len(filesToProcess)))
s.progress.AddTotalSize(totalSizeToProcess)
s.progress.GetStats().TotalFiles.Add(int64(len(filesToProcess)))
}
log.Info("Phase 1 complete",
@@ -889,9 +892,10 @@ func (s *Scanner) scanPhase(
}
// Handle symlinks and directories
if handled := s.recordSpecialEntry(
filePath, info, existingFiles, collector, result); handled {
return nil
handled, err := s.recordSpecialEntry(
filePath, info, existingFiles, collector, result)
if handled {
return err
}
// Skip other non-regular files (devices, sockets, etc.)
@@ -933,22 +937,25 @@ func (s *Scanner) scanPhase(
}
// recordSpecialEntry records symlinks and directories (which have no
// data to chunk) and reports whether it handled the entry.
// data to chunk) and reports whether it handled the entry. For a symlink
// whose target cannot be read it returns handleWalkError's result.
func (s *Scanner) recordSpecialEntry(
filePath string, info os.FileInfo,
existingFiles map[string]struct{},
collector *scanCollector, result *ScanResult,
) bool {
) (bool, error) {
// Handle symlinks
if info.Mode()&os.ModeSymlink != 0 {
file := s.buildSymlinkEntry(filePath, info)
if file != nil {
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
file, err := s.buildSymlinkEntry(filePath, info)
if err != nil {
return true, s.handleWalkError(filePath, err)
}
return true
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
return true, nil
}
// Handle directories (record for permission/ownership preservation
@@ -958,10 +965,10 @@ func (s *Scanner) recordSpecialEntry(
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
return true
return true, nil
}
return false
return false, nil
}
// handleWalkError deals with a filesystem error surfaced by the walk:
@@ -1111,13 +1118,12 @@ func (s *Scanner) printScanProgressLine(
}
// buildSymlinkEntry creates a File record for a symlink.
// Returns nil if the link target cannot be read.
func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.File {
func (s *Scanner) buildSymlinkEntry(
path string, info os.FileInfo,
) (*database.File, error) {
target, err := os.Readlink(path)
if err != nil {
log.Debug("Cannot read symlink target", "path", path, "error", err)
return nil
return nil, err
}
var uid, gid uint32
@@ -1136,7 +1142,7 @@ func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.Fil
UID: uid,
GID: gid,
LinkTarget: types.FilePath(target),
}
}, nil
}
// buildDirectoryEntry creates a File record for a directory.
@@ -1362,7 +1368,7 @@ func (s *Scanner) processFileWithErrorHandling(
log.Warn("File was deleted during backup, skipping",
"path", fileToProcess.Path)
result.FilesSkipped++
countFailedFile(fileToProcess, result)
return true, nil
}
@@ -1373,7 +1379,7 @@ func (s *Scanner) processFileWithErrorHandling(
s.ui.Errorf("Failed to process %s: %v. Skipping (--skip-errors).",
s.ui.Path(fileToProcess.Path), err)
result.FilesSkipped++
countFailedFile(fileToProcess, result)
return true, nil
}
@@ -1384,6 +1390,18 @@ func (s *Scanner) processFileWithErrorHandling(
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.
// Phase 1 counts no directories, so a directory is not counted here.
func countFailedFile(fileToProcess *FileToProcess, result *ScanResult) {
if fileToProcess.FileInfo.IsDir() {
return
}
result.FilesFailed++
result.BytesScanned -= fileToProcess.FileInfo.Size()
}
// printProcessingProgress prints a periodic progress line during the process phase,
// showing files processed, bytes transferred, throughput, and ETA
func (s *Scanner) printProcessingProgress(
+95 -7
View File
@@ -3,6 +3,7 @@ package snapshot_test
import (
"context"
"errors"
"io"
"os"
"path/filepath"
"strings"
@@ -13,6 +14,7 @@ import (
"github.com/spf13/afero"
"sneak.berlin/go/vaultik/internal/database"
"sneak.berlin/go/vaultik/internal/snapshot"
"sneak.berlin/go/vaultik/internal/ui"
)
// errSimTempFail is the one-time temp-file creation failure blobTempFailFs
@@ -80,6 +82,30 @@ func (f *readFailFs) Open(name string) (afero.File, error) {
return file, nil
}
// linkRemovedAfterLstatFs is the real filesystem, except that the symlink at
// target is removed right after the walk lstats it, as happens when a link is
// deleted during a backup. The scanner's readlink of it then fails.
type linkRemovedAfterLstatFs struct {
afero.OsFs
t *testing.T
target string
}
func (f *linkRemovedAfterLstatFs) LstatIfPossible(
name string,
) (os.FileInfo, bool, error) {
info, lstatCalled, err := f.OsFs.LstatIfPossible(name)
if err == nil && name == f.target {
rmErr := os.Remove(name)
if rmErr != nil {
f.t.Errorf("removing %s: %v", name, rmErr)
}
}
return info, lstatCalled, err
}
// writeSkipErrorTestFile writes one file into fs with a fixed mtime.
func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
t.Helper()
@@ -102,10 +128,11 @@ func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
}
}
// runSkipErrorScan scans /source on fs with the given skip-errors setting and
// returns the repositories (for inspection) and the scan error.
// runSkipErrorScan scans source on fs with the given skip-errors setting,
// printing user-facing messages to uiw (nil discards them), and returns the
// repositories (for inspection) and the scan error.
func runSkipErrorScan(
t *testing.T, fs afero.Fs, skipErrors bool,
t *testing.T, fs afero.Fs, source string, skipErrors bool, uiw *ui.Writer,
) (*database.Repositories, error) {
t.Helper()
@@ -130,6 +157,7 @@ func runSkipErrorScan(
MaxBlobSize: int64(1024 * 1024),
CompressionLevel: 3,
AgeRecipients: []string{testAgePublicKey},
UI: uiw,
SkipErrors: skipErrors,
})
@@ -137,7 +165,7 @@ func runSkipErrorScan(
snapshotID := "test-snapshot-skip-errors"
createTestSnapshotRecord(ctx, t, repos, snapshotID)
_, err = scanner.Scan(ctx, "/source", snapshotID)
_, err = scanner.Scan(ctx, source, snapshotID)
return repos, err
}
@@ -157,7 +185,7 @@ func TestScannerPackingFailureAbortsUnderSkipErrors(t *testing.T) {
writeSkipErrorTestFile(t, fs, "/source/file1.txt", "first file content")
writeSkipErrorTestFile(t, fs, "/source/file2.txt", "second file content")
repos, err := runSkipErrorScan(t, fs, true)
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
if err == nil {
t.Fatal("expected scan to abort on the packer error, got nil")
}
@@ -184,7 +212,7 @@ func TestScannerReadErrorAbortsWithoutSkipErrors(t *testing.T) {
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
_, err := runSkipErrorScan(t, fs, false)
_, err := runSkipErrorScan(t, fs, "/source", false, nil)
if err == nil {
t.Fatal("expected scan to fail on the read error, got nil")
}
@@ -200,7 +228,7 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
fs := &readFailFs{Fs: afero.NewMemMapFs(), target: target}
writeSkipErrorTestFile(t, fs, target, "content that cannot be read")
repos, err := runSkipErrorScan(t, fs, true)
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
if err != nil {
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
}
@@ -214,3 +242,63 @@ func TestScannerReadErrorSkippedWithSkipErrors(t *testing.T) {
t.Fatalf("expected unreadable file skipped, got %d chunks", len(chunks))
}
}
// writeSymlinkSource creates a source directory on disk holding one symlink
// and returns the directory and the symlink's path.
func writeSymlinkSource(t *testing.T) (string, string) {
t.Helper()
sourceDir := t.TempDir()
linkPath := filepath.Join(sourceDir, "link")
err := os.Symlink("target.txt", linkPath)
if err != nil {
t.Fatalf("creating symlink: %v", err)
}
return sourceDir, linkPath
}
// TestScannerUnreadableSymlinkAbortsWithoutSkipErrors checks that a symlink
// whose target cannot be read aborts the run when --skip-errors is not set.
func TestScannerUnreadableSymlinkAbortsWithoutSkipErrors(t *testing.T) {
t.Parallel()
sourceDir, linkPath := writeSymlinkSource(t)
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
_, err := runSkipErrorScan(t, fs, sourceDir, false, nil)
if !errors.Is(err, os.ErrNotExist) {
t.Fatalf("expected scan to fail on the removed symlink, got %v", err)
}
}
// TestScannerUnreadableSymlinkSkippedWithSkipErrors checks that a symlink
// whose target cannot be read is skipped with an error line, and the run
// completes, when --skip-errors is set.
func TestScannerUnreadableSymlinkSkippedWithSkipErrors(t *testing.T) {
t.Parallel()
sourceDir, linkPath := writeSymlinkSource(t)
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
uiw := ui.NewWithColor(io.Discard, false)
repos, err := runSkipErrorScan(t, fs, sourceDir, true, uiw)
if err != nil {
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
}
if uiw.ErrorCount() != 1 {
t.Fatalf("expected one error line for the symlink, got %d",
uiw.ErrorCount())
}
file, err := repos.Files.GetByPath(context.Background(), linkPath)
if err != nil {
t.Fatalf("getting %s: %v", linkPath, err)
}
if file != nil {
t.Fatalf("expected %s not to be recorded", linkPath)
}
}
+26 -9
View File
@@ -115,20 +115,37 @@ func ShortHostname(hostname string) string {
// CreateSnapshotWithName creates a new snapshot record with an optional
// snapshot name. The snapshot ID format is: hostname_name_timestamp or
// hostname_timestamp if name is empty.
// hostname_timestamp if name is empty. The timestamp is in whole seconds.
// 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(
ctx context.Context, hostname, name, version, gitRevision string,
) (string, error) {
shortHostname := ShortHostname(hostname)
// Build snapshot ID with optional name
timestamp := time.Now().UTC().Format("2006-01-02T15:04:05Z")
var snapshotID string
if name != "" {
snapshotID = fmt.Sprintf("%s_%s_%s", shortHostname, name, timestamp)
} else {
snapshotID = fmt.Sprintf("%s_%s", shortHostname, timestamp)
for {
// Build snapshot ID with optional name
timestamp := time.Now().UTC().Format("2006-01-02T15:04:05Z")
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{
@@ -875,7 +892,7 @@ func (sm *SnapshotManager) getFileSize(path string) int64 {
// BackupStats contains statistics from a backup operation
type BackupStats struct {
FilesScanned int
TotalSize int64 // Total size of all files examined
TotalSize int64 // Total size of the files in the snapshot
ChunksCreated int
BlobsCreated int
BytesUploaded int64
+33
View File
@@ -386,3 +386,36 @@ func TestCleanSnapshotDBNonExistentSnapshot(t *testing.T) {
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)
}
}
+3 -2
View File
@@ -14,8 +14,9 @@ import (
// runStorerConformance is the shared Storer contract. Every backend that
// can run in-process is expected to pass it: TestFileStorer runs it against
// file://, TestS3Storer against s3://. A new backend inherits this coverage
// by passing its own constructor, so the contract is defined once.
// file://, TestS3Storer against s3://, TestRcloneStorer against rclone's
// local backend. A new backend inherits this coverage by passing its own
// constructor, so the contract is defined once.
//
// It exercises the public Storer interface: round-trip, stat, list with
// prefix filtering, overwrite, delete, delete-of-missing, and not-found on
+4 -3
View File
@@ -50,9 +50,10 @@ const storageDirPerm = 0o755
// temp file carrying this suffix and only renames it onto the real key once
// the whole object is on disk, so an interrupted write can never leave a
// truncated object at the key a later run would Stat and trust as a complete
// blob. List and ListStream skip these files, so a leftover from an
// interrupted write is never listed or trusted as a blob; it is otherwise
// harmless and is overwritten when the same key is written again.
// blob. The rclone backend's upload does the same on remotes with a
// server-side move. List and ListStream skip these files, so a leftover from
// an interrupted write is never listed or trusted as a blob; it is otherwise
// harmless.
const tempSuffix = ".partial"
// Put stores data at the specified key.
+5 -2
View File
@@ -14,6 +14,9 @@ import (
// errStreamInterrupted stands in for an upload cut off mid-stream.
var errStreamInterrupted = errors.New("connection reset mid-upload")
// testBlobKey is a key laid out as a blob's key is.
const testBlobKey = "blobs/aa/bb/aabbccddeeff"
// failingReader yields its data once, then fails.
type failingReader struct {
data []byte
@@ -43,7 +46,7 @@ func TestFileStorer_InterruptedWriteLeavesNoTrustedObject(t *testing.T) {
}
ctx := context.Background()
key := "blobs/aa/bb/aabbccddeeff"
key := testBlobKey
err = f.PutWithProgress(ctx, key, &failingReader{data: []byte("partial")}, 4096, nil)
if err == nil {
@@ -79,7 +82,7 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) {
}
ctx := context.Background()
realKey := "blobs/aa/bb/aabbccddeeff"
realKey := testBlobKey
err = f.Put(ctx, realKey, strings.NewReader("blob-bytes"))
if err != nil {
+64 -15
View File
@@ -3,6 +3,7 @@ package storage
import (
"bytes"
"context"
"crypto/rand"
"errors"
"fmt"
"io"
@@ -68,14 +69,7 @@ func (r *RcloneStorer) Put(ctx context.Context, key string, data io.Reader) erro
return fmt.Errorf("reading data: %w", err)
}
// Upload the object
_, err = operations.Rcat(ctx, r.fsys, key,
io.NopCloser(bytes.NewReader(buf)), time.Now(), nil)
if err != nil {
return fmt.Errorf("uploading object: %w", err)
}
return nil
return r.upload(ctx, key, bytes.NewReader(buf))
}
// PutWithProgress stores data with progress reporting.
@@ -89,13 +83,7 @@ func (r *RcloneStorer) PutWithProgress(
callback: progress,
}
// Upload the object
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(pr), time.Now(), nil)
if err != nil {
return fmt.Errorf("uploading object: %w", err)
}
return nil
return r.upload(ctx, key, pr)
}
// Get retrieves data from the specified key.
@@ -173,6 +161,10 @@ func (r *RcloneStorer) List(ctx context.Context, prefix string) ([]string, error
err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) {
key := obj.Remote()
if strings.HasSuffix(key, tempSuffix) {
return
}
if prefix == "" || strings.HasPrefix(key, prefix) {
keys = append(keys, key)
}
@@ -202,6 +194,10 @@ func (r *RcloneStorer) ListStream(
}
key := obj.Remote()
if strings.HasSuffix(key, tempSuffix) {
return
}
if prefix == "" || strings.HasPrefix(key, prefix) {
ch <- ObjectInfo{
Key: key,
@@ -230,6 +226,59 @@ func (r *RcloneStorer) Info() Info {
}
}
// upload writes data to key. Where the remote has a server-side move, it
// writes under a temporary name ending in tempSuffix and moves the object
// onto key once it is complete, so a killed upload cannot leave a truncated
// object at key; a remote without one is written in place. List and
// ListStream skip a temporary object left behind.
//
// rclone's own copy does this only where the remote also sets
// PartialUploads. That flag is not checked here: hdfs, for one, shows a
// file while it is written without setting it.
func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error {
if r.fsys.Features().Move == nil {
_, err := operations.Rcat(ctx, r.fsys, key, io.NopCloser(data), time.Now(), nil)
if err != nil {
return fmt.Errorf("uploading object: %w", err)
}
return nil
}
tempKey := key + "-" + rand.Text() + tempSuffix
obj, err := operations.Rcat(ctx, r.fsys, tempKey, io.NopCloser(data), time.Now(), nil)
if err != nil {
// Rcat returns the object it wrote when the written data fails its check.
if obj != nil {
_ = obj.Remove(ctx)
}
return fmt.Errorf("uploading object: %w", err)
}
// On drive, dropbox, onedrive and others the remote's own move does not
// replace an object already at key. operations.Move removes the object
// it is given first, and copies where the remote refuses the move.
existing, err := r.fsys.NewObject(ctx, key)
if errors.Is(err, fs.ErrorObjectNotFound) {
existing = nil
} else if err != nil {
_ = obj.Remove(ctx)
return fmt.Errorf("looking up existing object: %w", err)
}
_, err = operations.Move(ctx, r.fsys, existing, key, obj)
if err != nil {
_ = obj.Remove(ctx)
return fmt.Errorf("moving object into place: %w", err)
}
return nil
}
// progressReader wraps an io.Reader to track read progress.
type progressReader struct {
reader io.Reader
+307 -14
View File
@@ -1,28 +1,47 @@
package storage_test
import (
"bytes"
"context"
"errors"
"io"
"os"
"path/filepath"
"strings"
"testing"
"github.com/rclone/rclone/fs"
"github.com/rclone/rclone/fs/config/configmap"
"sneak.berlin/go/vaultik/internal/storage"
)
// The rclone backend is a thin adapter over the rclone library: it turns a
// (remote, path) pair into rclone's "remote:path" string, hands it to
// rclone, and maps rclone's own results back to the Storer interface. What
// can be tested in-process, without a configured remote or network, is that
// adapter layer — how the arguments are shaped and how construction errors
// are reported. The data-plane operations (Put/Get/List/Delete) are rclone's
// own, exercised against a real provider (drive, s3-via-rclone, ...), which
// needs a configured remote with credentials and network access and so is
// out of reach of a unit test. The shared Storer conformance suite therefore
// runs against the in-process file and s3 backends; the rclone backend
// inherits that contract once a remote is configured.
//
// These tests use rclone's ":local:" on-the-fly backend, which addresses the
// local filesystem directly without any configured remote, so construction
// runs entirely in-process.
// local filesystem directly without any configured remote, so they run
// entirely in-process. A remote that needs credentials and network access
// (drive, s3 via rclone, ...) is out of reach of a unit test.
// newRcloneStorer builds an rclone backend on rclone's local backend,
// rooted at a fresh temp directory.
//
//nolint:ireturn // conformance runs against the Storer interface by design
func newRcloneStorer(t *testing.T) storage.Storer {
t.Helper()
s, err := storage.NewRcloneStorer(context.Background(), ":local", t.TempDir())
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
return s
}
// TestRcloneStorer runs the shared Storer contract against the rclone
// backend.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorer(t *testing.T) {
runStorerConformance(t, newRcloneStorer)
}
// TestNewRcloneStorerConstruction checks that a valid remote constructs a
// backend and that Info() reports the shaped "remote:path" location.
@@ -44,6 +63,280 @@ func TestNewRcloneStorerConstruction(t *testing.T) {
}
}
// unbufferedUploadContext returns a context in which rclone neither reads
// ahead of the write nor holds a small upload in memory. Without it, a
// progress callback runs before anything is written to the remote.
func unbufferedUploadContext() context.Context {
ctx, ci := fs.AddConfig(context.Background())
ci.BufferSize = 0
ci.StreamingUploadCutoff = 0
return ctx
}
// objectAtKeyDuringUpload uploads data to testBlobKey and reports whether
// an object was at the key before the upload finished. It fails the test
// unless the key then reads back as the uploaded data.
func objectAtKeyDuringUpload(
ctx context.Context, t *testing.T, s *storage.RcloneStorer,
) bool {
t.Helper()
data := bytes.Repeat([]byte("blob-bytes"), 1000)
seen := false
err := s.PutWithProgress(ctx, testBlobKey, bytes.NewReader(data),
int64(len(data)), func(int64) error {
_, statErr := s.Stat(ctx, testBlobKey)
if statErr == nil {
seen = true
}
return nil
})
if err != nil {
t.Fatalf("PutWithProgress: %v", err)
}
got := readObject(ctx, t, s, testBlobKey)
if !bytes.Equal(got, data) {
t.Errorf("object read back as %d bytes, want the %d uploaded",
len(got), len(data))
}
return seen
}
// readObject returns the contents of the object at key.
func readObject(
ctx context.Context, t *testing.T, s *storage.RcloneStorer, key string,
) []byte {
t.Helper()
rc, err := s.Get(ctx, key)
if err != nil {
t.Fatalf("Get: %v", err)
}
defer func() { _ = rc.Close() }()
got, err := io.ReadAll(rc)
if err != nil {
t.Fatalf("reading object: %v", err)
}
return got
}
// TestRcloneStorerObjectAppearsOnlyWhenComplete checks that on a remote
// with a server-side move, such as local, nothing is at the key until the
// upload has finished, so a killed upload cannot leave a truncated object
// there.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerObjectAppearsOnlyWhenComplete(t *testing.T) {
ctx := unbufferedUploadContext()
s, err := storage.NewRcloneStorer(ctx, ":local", t.TempDir())
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
if objectAtKeyDuringUpload(ctx, t, s) {
t.Error("object was at its key before the upload finished")
}
}
// TestRcloneStorerListSkipsPartialFiles checks that a temporary file left
// by a killed upload is never listed as a key.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerListSkipsPartialFiles(t *testing.T) {
dir := t.TempDir()
ctx := context.Background()
s, err := storage.NewRcloneStorer(ctx, ":local", dir)
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
realKey := testBlobKey
err = s.Put(ctx, realKey, strings.NewReader("blob-bytes"))
if err != nil {
t.Fatalf("Put: %v", err)
}
leftover := filepath.Join(dir, realKey+"-123456.partial")
err = os.WriteFile(leftover, []byte("half"), 0o600)
if err != nil {
t.Fatalf("writing leftover temp file: %v", err)
}
keys, err := s.List(ctx, "blobs/")
if err != nil {
t.Fatalf("List: %v", err)
}
if len(keys) != 1 || keys[0] != realKey {
t.Fatalf("List should return only the real key, got %v", keys)
}
var streamed []string
for obj := range s.ListStream(ctx, "blobs/") {
if obj.Err != nil {
t.Fatalf("ListStream: %v", obj.Err)
}
streamed = append(streamed, obj.Key)
}
if len(streamed) != 1 || streamed[0] != realKey {
t.Fatalf("ListStream should return only the real key, got %v", streamed)
}
}
// newRcloneStorerOnWrappedLocal registers name as rclone's local backend
// wrapped by wrap, and builds an rclone backend on it rooted at a fresh
// temp directory. wrap changes the features the local backend reports, so
// that it behaves like a remote a unit test cannot reach.
func newRcloneStorerOnWrappedLocal(
ctx context.Context, t *testing.T, name string, wrap func(fs.Fs) fs.Fs,
) *storage.RcloneStorer {
t.Helper()
fs.Register(&fs.RegInfo{
Name: name,
NewFs: func(
ctx context.Context, _, root string, _ configmap.Mapper,
) (fs.Fs, error) {
local, err := fs.NewFs(ctx, ":local:"+root)
if err != nil {
return nil, err
}
return wrap(local), nil
},
})
s, err := storage.NewRcloneStorer(ctx, ":"+name, t.TempDir())
if err != nil {
t.Fatalf("NewRcloneStorer: %v", err)
}
return s
}
// withoutPartialUploads is rclone's local backend with the PartialUploads
// flag cleared. Like hdfs, it then has a server-side move and shows a file
// while it is written, without setting that flag.
type withoutPartialUploads struct {
fs.Fs
}
func (f *withoutPartialUploads) Features() *fs.Features {
features := *f.Fs.Features()
features.PartialUploads = false
return &features
}
// TestRcloneStorerMovesIntoPlaceWithoutPartialUploads checks that on a
// remote with a server-side move nothing is at the key until the upload has
// finished, even when rclone does not mark the remote as showing partial
// uploads.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerMovesIntoPlaceWithoutPartialUploads(t *testing.T) {
ctx := unbufferedUploadContext()
s := newRcloneStorerOnWrappedLocal(ctx, t, "withoutpartialuploads",
func(local fs.Fs) fs.Fs { return &withoutPartialUploads{Fs: local} })
if objectAtKeyDuringUpload(ctx, t, s) {
t.Error("object was at its key before the upload finished")
}
}
// withoutMove is rclone's local backend reporting no server-side move.
type withoutMove struct {
fs.Fs
}
func (f *withoutMove) Features() *fs.Features {
features := *f.Fs.Features()
features.Move = nil
return &features
}
// TestRcloneStorerWritesInPlaceWithoutMove checks that on a remote with no
// server-side move an object is written straight to its key and reads back
// from there.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerWritesInPlaceWithoutMove(t *testing.T) {
ctx := unbufferedUploadContext()
s := newRcloneStorerOnWrappedLocal(ctx, t, "withoutmove",
func(local fs.Fs) fs.Fs { return &withoutMove{Fs: local} })
if !objectAtKeyDuringUpload(ctx, t, s) {
t.Error("object was not at its key while it was uploaded")
}
}
var errNameConflict = errors.New("an object with this name already exists")
// moveRefusesExisting is rclone's local backend with a server-side move
// that, like dropbox's or onedrive's, refuses to move onto an existing
// object.
type moveRefusesExisting struct {
fs.Fs
}
func (f *moveRefusesExisting) Features() *fs.Features {
features := *f.Fs.Features()
features.Move = f.move
return &features
}
//nolint:ireturn // the signature is rclone's
func (f *moveRefusesExisting) move(
ctx context.Context, src fs.Object, remote string,
) (fs.Object, error) {
_, err := f.NewObject(ctx, remote)
if err == nil {
return nil, errNameConflict
}
return f.Fs.Features().Move(ctx, src, remote)
}
// TestRcloneStorerOverwritesWhereMoveRefusesExisting checks that writing a
// key twice replaces the object on a remote whose server-side move will not
// replace an existing object.
//
//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config
func TestRcloneStorerOverwritesWhereMoveRefusesExisting(t *testing.T) {
ctx := context.Background()
s := newRcloneStorerOnWrappedLocal(ctx, t, "moverefusesexisting",
func(local fs.Fs) fs.Fs { return &moveRefusesExisting{Fs: local} })
for _, content := range []string{"first", "second"} {
err := s.Put(ctx, testBlobKey, strings.NewReader(content))
if err != nil {
t.Fatalf("Put %q: %v", content, err)
}
}
got := readObject(ctx, t, s, testBlobKey)
if string(got) != "second" {
t.Errorf("object = %q, want %q", got, "second")
}
}
// TestNewRcloneStorerUnknownRemote checks that a remote that is not in the
// rclone config fails construction with the ErrRemoteNotFound sentinel,
// rather than silently returning a backend pointed nowhere.
+24 -6
View File
@@ -181,8 +181,11 @@ type SnapshotMetadataInfo struct {
ManifestSize int64 `json:"manifest_size"`
DatabaseSize int64 `json:"database_size"`
TotalSize int64 `json:"total_size"`
BlobCount int `json:"blob_count"`
BlobsSize int64 `json:"blobs_size"`
// Both stay nil (null in the JSON) when the snapshot's manifest was
// 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
// backup interrupted before its manifest upload leaves a directory
@@ -380,6 +383,10 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
for _, snapshotID := range snapshotIDs {
info := snapshotMetadata[snapshotID]
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
}
@@ -395,7 +402,7 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
continue
}
info.BlobCount = manifest.BlobCount
blobCount := manifest.BlobCount
var blobsSize int64
@@ -404,7 +411,8 @@ func (v *Vaultik) collectReferencedBlobsFromManifests(
blobsSize += blob.CompressedSize
}
info.BlobsSize = blobsSize
info.BlobCount = &blobCount
info.BlobsSize = &blobsSize
}
return referencedBlobs, unreadable
@@ -516,13 +524,23 @@ func (v *Vaultik) printRemoteInfoTable(result *RemoteInfoResult) {
v.stdoutf("%s", separator)
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,
truncateString(info.SnapshotID, snapshotIDColWidth),
ubytes(info.ManifestSize),
ubytes(info.DatabaseSize),
ubytes(info.TotalSize),
humanize.Comma(int64(info.BlobCount)),
ubytes(info.BlobsSize),
blobCount,
blobsSize,
)
}
+2 -2
View File
@@ -42,7 +42,7 @@ func TestTableCountForReportSurfacesReadFailure(t *testing.T) {
assert.Nil(t, missing, "a failed read is unknown, not a count")
// The rendered count for a failed read must say unknown, never 0.
assert.Equal(t, countUnknown, countText(missing))
assert.Equal(t, unknownText, countText(missing))
assert.NotEqual(t, "0", countText(missing))
}
@@ -57,7 +57,7 @@ func TestCountTextDistinguishesEmptyFromUnknown(t *testing.T) {
assert.Equal(t, "0", countText(&zero))
assert.Equal(t, "7", countText(&seven))
assert.Equal(t, countUnknown, countText(nil))
assert.Equal(t, unknownText, countText(nil))
}
// TestCountDiffUnknownWhenEitherSideUnknown checks that a delta computed
+40
View File
@@ -0,0 +1,40 @@
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,6 +4,7 @@ import (
"bytes"
"context"
"encoding/json"
"strings"
"testing"
"time"
@@ -69,6 +70,74 @@ func TestRemoteInfo_UnreadableManifestLeavesOrphansUnknown(t *testing.T) {
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
// 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
+16 -6
View File
@@ -185,6 +185,7 @@ type snapshotStats struct {
totalBlobs int
totalBytesSkipped int64
totalFilesSkipped int
totalFilesFailed int
totalFilesDeleted int
totalBytesDeleted int64
totalBytesUploaded int64
@@ -315,6 +316,7 @@ func (v *Vaultik) scanAllDirectories(
stats.totalChunks += result.ChunksCreated
stats.totalBlobs += result.BlobsCreated
stats.totalFilesSkipped += result.FilesSkipped
stats.totalFilesFailed += result.FilesFailed
stats.totalBytesSkipped += result.BytesSkipped
stats.totalFilesDeleted += result.FilesDeleted
stats.totalBytesDeleted += result.BytesDeleted
@@ -326,6 +328,7 @@ func (v *Vaultik) scanAllDirectories(
"path", dir,
"files", result.FilesScanned,
"files_skipped", result.FilesSkipped,
"files_failed", result.FilesFailed,
"bytes", result.BytesScanned,
"bytes_skipped", result.BytesSkipped,
"chunks", result.ChunksCreated,
@@ -359,9 +362,11 @@ func (v *Vaultik) finalizeSnapshotMetadata(
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{
BackupStats: snapshot.BackupStats{
FilesScanned: stats.totalFiles,
FilesScanned: stats.totalFiles - stats.totalFilesFailed,
TotalSize: stats.totalBytes + stats.totalBytesSkipped,
ChunksCreated: stats.totalChunks,
BlobsCreated: stats.totalBlobs,
@@ -409,7 +414,8 @@ func (v *Vaultik) printSnapshotSummary(
snapshotID string, startTime time.Time, stats *snapshotStats,
) {
snapshotDuration := time.Since(startTime)
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped
totalFilesChanged := stats.totalFiles - stats.totalFilesSkipped -
stats.totalFilesFailed
totalBytesAll := stats.totalBytes + stats.totalBytesSkipped
var compressionRatio float64
@@ -426,6 +432,10 @@ func (v *Vaultik) printSnapshotSummary(
v.UI.Count(stats.totalFiles),
v.UI.Count(totalFilesChanged),
v.UI.Count(stats.totalFilesSkipped))
if stats.totalFilesFailed > 0 {
filesMsg += fmt.Sprintf(", %s failed", v.UI.Count(stats.totalFilesFailed))
}
if stats.totalFilesDeleted > 0 {
filesMsg += fmt.Sprintf(", %s deleted", v.UI.Count(stats.totalFilesDeleted))
}
@@ -1756,9 +1766,9 @@ func (v *Vaultik) PruneDatabase() (*PruneResult, error) {
return result, nil
}
// countUnknown is what a count reads as when its query could not be run,
// distinct from "0", which means the table really was empty.
const countUnknown = "unknown"
// unknownText is what a count or size reads as when it could not be
// determined, distinct from "0", which is a real zero.
const unknownText = "unknown"
// tableCountForReport returns the row count of a table for the prune
// summary, or nil if the count could not be read. A read failure is
@@ -1795,7 +1805,7 @@ func countDiff(before, after *int64) *int64 {
// one that could not be queried.
func countText(count *int64) string {
if count == nil {
return countUnknown
return unknownText
}
return strconv.FormatInt(*count, 10)
@@ -0,0 +1,123 @@
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)
})
}
}