Compare commits

Author SHA1 Message Date
sneak b4284b2170 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, and the
processing start time the rate is measured from is the first 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 01:48:40 +00: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
14 changed files with 738 additions and 107 deletions
+10
View File
@@ -428,6 +428,16 @@ a local or mounted filesystem. Useful for testing or backing up to a NAS.
**Rclone** (`rclone://remote/path`): Uses rclone's 70+ supported cloud
providers. Requires rclone to be configured separately (`rclone config`).
An upload cut off part-way leaves nothing under the object's name on S3, which
shows an object only once its upload has completed, and on the local filesystem
backend, which writes a temporary file and renames it into place. Rclone
remotes with a server-side move (local and sftp among them) are written under a
temporary name ending in `.partial` and moved into place. Rclone remotes without
one are written in place: where such a remote shows a file while it is still
being written, a killed upload can leave a truncated object under its name,
which a later backup takes for complete. A leftover `.partial` file is ignored
and can be deleted.
Legacy S3 configuration via `s3.*` fields (endpoint, bucket, prefix, etc.) is
still supported for backward compatibility. `storage_url` takes precedence if
both are set.
+35
View File
@@ -22,6 +22,41 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps
- 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, and the rate is
measured from when the first path's processing started.
- 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: 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
+25 -28
View File
@@ -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,16 +296,17 @@ func RunOperation(
return err
}
// RunWithApp returns only after OnStop has waited for the goroutine
// to return, whether the shutdown came from op finishing or from an
// interrupt, so the goroutine's write to failed or interrupted is in
// place by the time we read it. If OnStop timed out, the goroutine
// may still be running, 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
+14 -7
View File
@@ -67,9 +67,9 @@ type ProgressStats struct {
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
TotalSize atomic.Int64 // Size to process in the paths scanned so far
TotalFiles atomic.Int64 // Files to process in the paths scanned so far
ProcessStartTime atomic.Value // stores time.Time; set by the first path
StartTime time.Time
mu sync.RWMutex
lastDetailTime time.Time
@@ -148,10 +148,17 @@ 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. The processed counts
// run across every path, so the total does too, and the processing start
// time, which the rate is measured from, is the first path's.
func (pr *ProgressReporter) AddTotalSize(size int64) {
pr.stats.TotalSize.Add(size)
_, started := pr.stats.ProcessStartTime.Load().(time.Time)
if !started {
pr.stats.ProcessStartTime.Store(time.Now().UTC())
}
}
// Helper functions
+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())
}
}
+24 -21
View File
@@ -407,8 +407,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 +889,10 @@ func (s *Scanner) scanPhase(
}
// Handle symlinks and directories
if handled := s.recordSpecialEntry(
filePath, info, existingFiles, collector, result); handled {
return nil
handled, err := s.recordSpecialEntry(
filePath, info, existingFiles, collector, result)
if handled {
return err
}
// Skip other non-regular files (devices, sockets, etc.)
@@ -933,22 +934,25 @@ func (s *Scanner) scanPhase(
}
// recordSpecialEntry records symlinks and directories (which have no
// data to chunk) and reports whether it handled the entry.
// data to chunk) and reports whether it handled the entry. For a symlink
// whose target cannot be read it returns handleWalkError's result.
func (s *Scanner) recordSpecialEntry(
filePath string, info os.FileInfo,
existingFiles map[string]struct{},
collector *scanCollector, result *ScanResult,
) bool {
) (bool, error) {
// Handle symlinks
if info.Mode()&os.ModeSymlink != 0 {
file := s.buildSymlinkEntry(filePath, info)
if file != nil {
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
file, err := s.buildSymlinkEntry(filePath, info)
if err != nil {
return true, s.handleWalkError(filePath, err)
}
return true
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
s.updateScanEntryStats(result, true, info)
return true, nil
}
// Handle directories (record for permission/ownership preservation
@@ -958,10 +962,10 @@ func (s *Scanner) recordSpecialEntry(
existingFiles[filePath] = struct{}{}
collector.addToProcess(filePath, info, file)
return true
return true, nil
}
return false
return false, nil
}
// handleWalkError deals with a filesystem error surfaced by the walk:
@@ -1111,13 +1115,12 @@ func (s *Scanner) printScanProgressLine(
}
// buildSymlinkEntry creates a File record for a symlink.
// Returns nil if the link target cannot be read.
func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.File {
func (s *Scanner) buildSymlinkEntry(
path string, info os.FileInfo,
) (*database.File, error) {
target, err := os.Readlink(path)
if err != nil {
log.Debug("Cannot read symlink target", "path", path, "error", err)
return nil
return nil, err
}
var uid, gid uint32
@@ -1136,7 +1139,7 @@ func (s *Scanner) buildSymlinkEntry(path string, info os.FileInfo) *database.Fil
UID: uid,
GID: gid,
LinkTarget: types.FilePath(target),
}
}, nil
}
// buildDirectoryEntry creates a File record for a directory.
+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)
}
}
+25 -8
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{
+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.