From 42f6ae67a563cfc25267a962037f58026c1f5d8c Mon Sep 17 00:00:00 2001 From: sneak Date: Wed, 7 Oct 2026 16:22:07 +0000 Subject: [PATCH] Write rclone uploads under a temporary name, re-upload short blobs (closes #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, skipped the upload and recorded a snapshot that could not be restored. On a remote where rclone says a file can be seen while it is still being written, an object is now written under a name ending in `.partial` and moved onto its key with the remote's server-side move, as rclone's own copy does; listings skip such names. A backup also uploads a blob again when the stored object's size differs from the blob's. Judgement call: the temporary name is used only where rclone sets its PartialUploads feature; other remotes already show an object only once complete. Model: opus-5-5 --- README.md | 11 ++ TODO.md | 12 ++ internal/snapshot/scanner.go | 17 ++- internal/storage/conformance_test.go | 5 +- internal/storage/file.go | 6 +- internal/storage/file_atomic_test.go | 7 +- internal/storage/rclone.go | 65 +++++++--- internal/storage/rclone_test.go | 146 ++++++++++++++++++++--- internal/vaultik/fault_injection_test.go | 64 ++++++++++ 9 files changed, 292 insertions(+), 41 deletions(-) diff --git a/README.md b/README.md index 529ebd1..6fb8fde 100644 --- a/README.md +++ b/README.md @@ -428,6 +428,17 @@ 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 that show a file while it is still being written (local, sftp, ftp and +smb among them) are written under a temporary name ending in `.partial` and +moved into place where the remote has a server-side move; without one they are +written in place. Other rclone remotes show an object only once its upload has +completed. A leftover `.partial` file is ignored and can be deleted. A backup +that finds a blob stored at a size other than the one it packed uploads the +blob again. + Legacy S3 configuration via `s3.*` fields (endpoint, bucket, prefix, etc.) is still supported for backward compatibility. `storage_url` takes precedence if both are set. diff --git a/TODO.md b/TODO.md index 076aae6..0958887 100644 --- a/TODO.md +++ b/TODO.md @@ -22,6 +22,18 @@ the tag exists and is exercised; what is left is merging `next` to # Completed Steps +- 2026-10-07: Stopped a backup from trusting a blob object left short by + a killed upload + ([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 where rclone says a + file can be seen while it is being written, an object is now written + under a temporary name ending in `.partial` and moved into place, and + listings skip such names. A backup also uploads a blob again when the + stored object's size differs from the blob's. + - 2026-10-07: Corrected documentation, help text and comments that were false about the code ([issue #233](https://git.eeqj.de/sneak/vaultik/issues/233)). A blob diff --git a/internal/snapshot/scanner.go b/internal/snapshot/scanner.go index b611857..2bf1627 100644 --- a/internal/snapshot/scanner.go +++ b/internal/snapshot/scanner.go @@ -1535,8 +1535,8 @@ func (s *Scanner) handleBlobReady( return nil } -// uploadBlobIfNeeded uploads the blob to storage if it doesn't already -// exist, returns whether it existed +// uploadBlobIfNeeded uploads the blob to storage unless an object of the +// blob's size is already there, and returns whether it was there func (s *Scanner) uploadBlobIfNeeded( ctx context.Context, blobPath string, @@ -1546,11 +1546,13 @@ func (s *Scanner) uploadBlobIfNeeded( ) (bool, error) { finishedBlob := blobWithReader.FinishedBlob - // Check if blob already exists (deduplication after restart) + // Check if blob already exists (deduplication after restart). An + // object of another size was left by an upload cut off part-way and + // is uploaded again. destination := s.storage.Info().Location - _, err := s.storage.Stat(ctx, blobPath) - if err == nil { + stored, err := s.storage.Stat(ctx, blobPath) + if err == nil && stored.Size == finishedBlob.Compressed { log.Info("Blob already exists in storage, skipping upload", "hash", finishedBlob.Hash, "size", humanize.Bytes(safeUint64(finishedBlob.Compressed))) @@ -1561,6 +1563,11 @@ func (s *Scanner) uploadBlobIfNeeded( return true, nil } + if err == nil { + log.Warn("Blob in storage has the wrong size, uploading it again", + "hash", finishedBlob.Hash, "stored_size", stored.Size) + } + s.ui.Beginf("Uploading blob %s (%s) to %s.", s.ui.Hex(finishedBlob.Hash), s.ui.Size(finishedBlob.Compressed), s.ui.Path(destination)) diff --git a/internal/storage/conformance_test.go b/internal/storage/conformance_test.go index 7282341..e5c35f9 100644 --- a/internal/storage/conformance_test.go +++ b/internal/storage/conformance_test.go @@ -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 diff --git a/internal/storage/file.go b/internal/storage/file.go index bcf2c3d..e203cc4 100644 --- a/internal/storage/file.go +++ b/internal/storage/file.go @@ -50,9 +50,9 @@ 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 that need it. +// 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. diff --git a/internal/storage/file_atomic_test.go b/internal/storage/file_atomic_test.go index 9006695..b9be4e6 100644 --- a/internal/storage/file_atomic_test.go +++ b/internal/storage/file_atomic_test.go @@ -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 { diff --git a/internal/storage/rclone.go b/internal/storage/rclone.go index c3a3433..2e70130 100644 --- a/internal/storage/rclone.go +++ b/internal/storage/rclone.go @@ -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,45 @@ func (r *RcloneStorer) Info() Info { } } +// upload writes data to key. On a remote where rclone says a file can be +// seen while it is still being written (its PartialUploads feature: local, +// sftp, ftp, smb and others), it writes under a temporary name ending in +// tempSuffix and moves the object onto key once it is complete, as +// rclone's own copy does, so a killed upload cannot leave a truncated +// object at key. List and ListStream skip a temporary object left behind. +func (r *RcloneStorer) upload(ctx context.Context, key string, data io.Reader) error { + features := r.fsys.Features() + if !features.PartialUploads || 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) + } + + _, err = features.Move(ctx, obj, key) + 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 diff --git a/internal/storage/rclone_test.go b/internal/storage/rclone_test.go index e941497..0f00105 100644 --- a/internal/storage/rclone_test.go +++ b/internal/storage/rclone_test.go @@ -1,28 +1,45 @@ package storage_test import ( + "bytes" "context" "errors" + "os" + "path/filepath" + "strings" "testing" + "github.com/rclone/rclone/fs" "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 +61,107 @@ func TestNewRcloneStorerConstruction(t *testing.T) { } } +// TestRcloneStorerObjectAppearsOnlyWhenComplete checks that on a remote +// where a file can be seen while it is still being written, 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) { + // Without this, rclone reads ahead of the write and holds a small + // upload in memory, so the progress callback would run before anything + // is written to the remote. + ctx, ci := fs.AddConfig(context.Background()) + ci.BufferSize = 0 + ci.StreamingUploadCutoff = 0 + + s, err := storage.NewRcloneStorer(ctx, ":local", t.TempDir()) + if err != nil { + t.Fatalf("NewRcloneStorer: %v", err) + } + + key := testBlobKey + data := bytes.Repeat([]byte("blob-bytes"), 1000) + seenEarly := false + + err = s.PutWithProgress(ctx, key, bytes.NewReader(data), int64(len(data)), + func(int64) error { + _, statErr := s.Stat(ctx, key) + if statErr == nil { + seenEarly = true + } + + return nil + }) + if err != nil { + t.Fatalf("PutWithProgress: %v", err) + } + + if seenEarly { + t.Error("object was at its key before the upload finished") + } + + info, err := s.Stat(ctx, key) + if err != nil { + t.Fatalf("Stat after upload: %v", err) + } + + if info.Size != int64(len(data)) { + t.Errorf("stored size = %d, want %d", info.Size, len(data)) + } +} + +// 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) + } +} + // 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. diff --git a/internal/vaultik/fault_injection_test.go b/internal/vaultik/fault_injection_test.go index 1595151..63184c6 100644 --- a/internal/vaultik/fault_injection_test.go +++ b/internal/vaultik/fault_injection_test.go @@ -5,6 +5,7 @@ import ( "errors" "io" "os" + "path" "path/filepath" "strings" "testing" @@ -414,6 +415,69 @@ func TestBackupRetryAfterInterruptedUploadIsRestorable(t *testing.T) { assertRestoredTree(t, fs, restoreDir, testFiles) } +// Scenario 1c: a killed upload left a short object at a blob's key, as an +// rclone remote that writes in place can. A later backup that packs the +// same blob must upload it again rather than trust the short object. See +// https://git.eeqj.de/sneak/vaultik/issues/266. +func TestBackupReplacesShortBlobObject(t *testing.T) { + log.Initialize(log.Config{}) + t.Parallel() + + fs := afero.NewOsFs() + tempDir := t.TempDir() + dataDir := filepath.Join(tempDir, "src") + restoreDir := filepath.Join(tempDir, "restored") + firstDBPath := filepath.Join(tempDir, "first.sqlite") + retryDBPath := filepath.Join(tempDir, "retry.sqlite") + + ctx := context.Background() + cfg := faultTestConfig() + testFiles := writeFaultSourceTree(t, fs, dataDir) + + store, err := storage.NewFileStorer(filepath.Join(tempDir, "remote")) + require.NoError(t, err) + + firstDB, err := database.New(ctx, firstDBPath) + require.NoError(t, err) + fullFaultBackup(ctx, t, fs, store, cfg, database.NewRepositories(firstDB), + dataDir, firstDBPath, "first") + require.NoError(t, firstDB.Close()) + + blobKeys, err := store.List(ctx, "blobs/") + require.NoError(t, err) + require.NotEmpty(t, blobKeys) + + shortKey := blobKeys[0] + require.NoError(t, store.Put(ctx, shortKey, strings.NewReader("cut off"))) + + // A backup with a new local index packs the same blobs again. + retryDB, err := database.New(ctx, retryDBPath) + require.NoError(t, err) + + retryRepos := database.NewRepositories(retryDB) + id := fullFaultBackup(ctx, t, fs, store, cfg, retryRepos, + dataDir, retryDBPath, "retry") + + blob, err := retryRepos.Blobs.GetByHash(ctx, path.Base(shortKey)) + require.NoError(t, err) + require.NotNil(t, blob, "the retry must pack the blob that was cut short") + require.NoError(t, retryDB.Close()) + + info, err := store.Stat(ctx, shortKey) + require.NoError(t, err) + assert.Equal(t, blob.CompressedSize, info.Size, + "the short object must be replaced by the whole blob") + + reader := newReaderVaultik(ctx, cfg, store, nil, fs) + require.NoError(t, reader.Restore(&vaultik.RestoreOptions{ + SnapshotID: id, + TargetDir: restoreDir, + Verify: true, + }), "the retry's snapshot must be restorable") + + assertRestoredTree(t, fs, restoreDir, testFiles) +} + // Scenario 2: the process dies during the metadata export, after the // database is uploaded but before the manifest. The destination is left // with blobs and a database but no manifest. verify and snapshot list