From 0bfa6dd42c820b59a6d4be7d68ee16874f29be33 Mon Sep 17 00:00:00 2001 From: sneak Date: Thu, 8 Oct 2026 07:12:38 +0000 Subject: [PATCH] Make remote nuke delete leftover .partial uploads (closes #281) The file and rclone listings skip an object whose name ends in `.partial`, the temporary name a `file://` or rclone upload writes before moving the object into place. `remote nuke` deletes only what the listings return, so it left the `.partial` objects killed uploads leave behind and still reported the destination store empty. Storer gains DeletePartialUploads. The file and rclone backends remove every `.partial` object under the prefix; S3 has none to remove, since it shows an object only once its upload completes. `remote nuke` calls it for `metadata/` and `blobs/` as its last step. Empty directories under a `file://` destination are still left behind. Model: opus-5-5 --- TODO.md | 7 ++ internal/storage/faultstore/faultstore.go | 5 ++ internal/storage/file.go | 36 +++++++++++ internal/storage/file_atomic_test.go | 42 ++++++++++++ internal/storage/rclone.go | 25 ++++++++ internal/storage/rclone_test.go | 41 ++++++++++++ internal/storage/s3.go | 6 ++ internal/storage/storer.go | 6 ++ .../vaultik/destination_validation_test.go | 4 ++ internal/vaultik/integration_test.go | 6 ++ internal/vaultik/nuke_remote_test.go | 64 +++++++++++++++++++ internal/vaultik/prune.go | 14 +++- internal/vaultik/remove_snapshot_test.go | 6 ++ 13 files changed, 260 insertions(+), 2 deletions(-) create mode 100644 internal/vaultik/nuke_remote_test.go diff --git a/TODO.md b/TODO.md index 930a1ae..46585b4 100644 --- a/TODO.md +++ b/TODO.md @@ -22,6 +22,13 @@ the tag exists and is exercised; what is left is merging `next` to # Completed Steps +- 2026-10-08: Made `remote nuke` delete the `.partial` files that + uploads cut off part-way leave on the destination store + ([issue #281](https://git.eeqj.de/sneak/vaultik/issues/281)). The + file and rclone listings skip such a file, so the command left it in + place and still reported the store empty. It now removes them under + `metadata/` and `blobs/` as its last step. + - 2026-10-08: Made a local index error while recording a directory or symlink stop a backup under `--skip-errors` ([issue #284](https://git.eeqj.de/sneak/vaultik/issues/284)). Phase 2 diff --git a/internal/storage/faultstore/faultstore.go b/internal/storage/faultstore/faultstore.go index 3aad29a..ae5d889 100644 --- a/internal/storage/faultstore/faultstore.go +++ b/internal/storage/faultstore/faultstore.go @@ -171,6 +171,11 @@ func (f *Storer) ListStream( return f.inner.ListStream(ctx, prefix) } +// DeletePartialUploads delegates unchanged. +func (f *Storer) DeletePartialUploads(ctx context.Context, prefix string) error { + return f.inner.DeletePartialUploads(ctx, prefix) +} + // Info delegates unchanged. func (f *Storer) Info() storage.Info { return f.inner.Info() diff --git a/internal/storage/file.go b/internal/storage/file.go index 7cb5653..7bb74d1 100644 --- a/internal/storage/file.go +++ b/internal/storage/file.go @@ -242,6 +242,42 @@ func (f *FileStorer) ListStream(ctx context.Context, prefix string) <-chan Objec return ch } +// DeletePartialUploads removes every file under prefix whose name ends in +// tempSuffix. A missing prefix has none to remove. +func (f *FileStorer) DeletePartialUploads(ctx context.Context, prefix string) error { + basePath := f.fullPath(prefix) + + exists, err := afero.Exists(f.fs, basePath) + if err != nil { + return fmt.Errorf("checking path: %w", err) + } + + if !exists { + return nil + } + + err = afero.Walk(f.fs, basePath, func(path string, info os.FileInfo, err error) error { + if err != nil { + return err + } + + if ctx.Err() != nil { + return ctx.Err() + } + + if info.IsDir() || !strings.HasSuffix(info.Name(), tempSuffix) { + return nil + } + + return f.fs.Remove(path) + }) + if err != nil { + return fmt.Errorf("walking directory: %w", err) + } + + return nil +} + // Info returns human-readable storage location information. func (f *FileStorer) Info() Info { return Info{ diff --git a/internal/storage/file_atomic_test.go b/internal/storage/file_atomic_test.go index b9be4e6..6030ea2 100644 --- a/internal/storage/file_atomic_test.go +++ b/internal/storage/file_atomic_test.go @@ -120,3 +120,45 @@ func TestFileStorer_ListSkipsPartialFiles(t *testing.T) { t.Fatalf("ListStream should return only the real key, got %v", streamed) } } + +// TestFileStorer_DeletePartialUploads checks that a leftover temp file is +// removed and the object at the real key is kept. +func TestFileStorer_DeletePartialUploads(t *testing.T) { + t.Parallel() + + base := t.TempDir() + + f, err := storage.NewFileStorer(base) + if err != nil { + t.Fatalf("NewFileStorer: %v", err) + } + + ctx := context.Background() + + err = f.Put(ctx, testBlobKey, strings.NewReader("blob-bytes")) + if err != nil { + t.Fatalf("Put: %v", err) + } + + leftover := filepath.Join(base, testBlobKey+"-123456.partial") + + err = os.WriteFile(leftover, []byte("half"), 0o600) + if err != nil { + t.Fatalf("writing leftover temp file: %v", err) + } + + err = f.DeletePartialUploads(ctx, "blobs/") + if err != nil { + t.Fatalf("DeletePartialUploads: %v", err) + } + + _, err = os.Stat(leftover) + if !os.IsNotExist(err) { + t.Errorf("leftover temp file was not removed: %v", err) + } + + _, err = f.Stat(ctx, testBlobKey) + if err != nil { + t.Errorf("Stat of the real key: %v", err) + } +} diff --git a/internal/storage/rclone.go b/internal/storage/rclone.go index dfc8e1f..43ac7ba 100644 --- a/internal/storage/rclone.go +++ b/internal/storage/rclone.go @@ -213,6 +213,31 @@ func (r *RcloneStorer) ListStream( return ch } +// DeletePartialUploads removes every object under prefix whose name ends +// in tempSuffix. +func (r *RcloneStorer) DeletePartialUploads(ctx context.Context, prefix string) error { + var partial []fs.Object + + err := operations.ListFn(ctx, r.fsys, func(obj fs.Object) { + key := obj.Remote() + if strings.HasPrefix(key, prefix) && strings.HasSuffix(key, tempSuffix) { + partial = append(partial, obj) + } + }) + if err != nil { + return fmt.Errorf("listing objects: %w", err) + } + + for _, obj := range partial { + err = obj.Remove(ctx) + if err != nil { + return fmt.Errorf("removing object: %w", err) + } + } + + return nil +} + // Info returns human-readable storage location information. func (r *RcloneStorer) Info() Info { location := r.remote diff --git a/internal/storage/rclone_test.go b/internal/storage/rclone_test.go index 8179e8d..c666a64 100644 --- a/internal/storage/rclone_test.go +++ b/internal/storage/rclone_test.go @@ -198,6 +198,47 @@ func TestRcloneStorerListSkipsPartialFiles(t *testing.T) { } } +// TestRcloneStorerDeletePartialUploads checks that a temporary file left +// by a killed upload is removed and the object at the real key is kept. +// +//nolint:paralleltest // NewRcloneStorer installs the process-global rclone config +func TestRcloneStorerDeletePartialUploads(t *testing.T) { + dir := t.TempDir() + ctx := context.Background() + + s, err := storage.NewRcloneStorer(ctx, ":local", dir) + if err != nil { + t.Fatalf("NewRcloneStorer: %v", err) + } + + err = s.Put(ctx, testBlobKey, strings.NewReader("blob-bytes")) + if err != nil { + t.Fatalf("Put: %v", err) + } + + leftover := filepath.Join(dir, testBlobKey+"-123456.partial") + + err = os.WriteFile(leftover, []byte("half"), 0o600) + if err != nil { + t.Fatalf("writing leftover temp file: %v", err) + } + + err = s.DeletePartialUploads(ctx, "blobs/") + if err != nil { + t.Fatalf("DeletePartialUploads: %v", err) + } + + _, err = os.Stat(leftover) + if !os.IsNotExist(err) { + t.Errorf("leftover temp file was not removed: %v", err) + } + + _, err = s.Stat(ctx, testBlobKey) + if err != nil { + t.Errorf("Stat of the real key: %v", err) + } +} + // newRcloneStorerOnWrappedLocal registers name as rclone's local backend // wrapped by wrap, and builds an rclone backend on it rooted at a fresh // temp directory. wrap changes the features the local backend reports, so diff --git a/internal/storage/s3.go b/internal/storage/s3.go index 648e697..4f2b53b 100644 --- a/internal/storage/s3.go +++ b/internal/storage/s3.go @@ -99,6 +99,12 @@ func (s *S3Storer) ListStream(ctx context.Context, prefix string) <-chan ObjectI return ch } +// DeletePartialUploads has nothing to remove: S3 shows an object only once +// its upload has completed, so an upload cut off part-way leaves none. +func (s *S3Storer) DeletePartialUploads(_ context.Context, _ string) error { + return nil +} + // Info returns human-readable storage location information. func (s *S3Storer) Info() Info { return Info{ diff --git a/internal/storage/storer.go b/internal/storage/storer.go index 2334a7d..44a4861 100644 --- a/internal/storage/storer.go +++ b/internal/storage/storer.go @@ -71,6 +71,12 @@ type Storer interface { // If an error occurs during listing, the final item will have Err set. ListStream(ctx context.Context, prefix string) <-chan ObjectInfo + // DeletePartialUploads removes every object under prefix that an + // upload cut off part-way left under a temporary name ending in + // `.partial`. The file and rclone backends' List and ListStream skip + // such an object. + DeletePartialUploads(ctx context.Context, prefix string) error + // Info returns human-readable storage location information. Info() Info } diff --git a/internal/vaultik/destination_validation_test.go b/internal/vaultik/destination_validation_test.go index d154fc9..4a083bf 100644 --- a/internal/vaultik/destination_validation_test.go +++ b/internal/vaultik/destination_validation_test.go @@ -257,6 +257,10 @@ func (s *stubLister) List(_ context.Context, _ string) ([]string, error) { return nil, errStubUnused } +func (s *stubLister) DeletePartialUploads(_ context.Context, _ string) error { + return errStubUnused +} + func (s *stubLister) Info() storage.Info { return storage.Info{} } diff --git a/internal/vaultik/integration_test.go b/internal/vaultik/integration_test.go index ded2e7e..e2891ee 100644 --- a/internal/vaultik/integration_test.go +++ b/internal/vaultik/integration_test.go @@ -155,6 +155,12 @@ func (m *MockStorer) ListStream( return ch } +// DeletePartialUploads has nothing to remove: Put stores each object +// under its key at once. +func (m *MockStorer) DeletePartialUploads(_ context.Context, _ string) error { + return nil +} + func (m *MockStorer) Info() storage.Info { return storage.Info{ Type: "mock", diff --git a/internal/vaultik/nuke_remote_test.go b/internal/vaultik/nuke_remote_test.go new file mode 100644 index 0000000..6e22699 --- /dev/null +++ b/internal/vaultik/nuke_remote_test.go @@ -0,0 +1,64 @@ +package vaultik_test + +import ( + "context" + "io/fs" + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "sneak.berlin/go/vaultik/internal/log" + "sneak.berlin/go/vaultik/internal/snapshot" +) + +// TestNukeRemoteLeavesNoFiles checks that remote nuke leaves no file +// under a file:// destination, including the `.partial` files uploads +// killed part-way leave next to a snapshot's metadata and next to where a +// blob would have been. +func TestNukeRemoteLeavesNoFiles(t *testing.T) { + log.Initialize(log.Config{}) + t.Parallel() + + ctx := context.Background() + storeDir := filepath.Join(t.TempDir(), "store") + v, repos, _ := backUpToFileDestination(ctx, t, storeDir) + + snapshots, err := repos.Snapshots.ListRecent(ctx, listRecentTestLimit) + require.NoError(t, err) + require.Len(t, snapshots, 1) + + snapshotKey := snapshot.RemoteSnapshotKey(snapshots[0].ID.String()) + hash := testBlobHashA + leftovers := []string{ + filepath.Join(storeDir, "metadata", snapshotKey, + "db.zst.age-123456.partial"), + filepath.Join(storeDir, "blobs", hash[:2], hash[2:4], + hash+"-123456.partial"), + } + + for _, leftover := range leftovers { + require.NoError(t, os.MkdirAll(filepath.Dir(leftover), 0o750)) + require.NoError(t, os.WriteFile(leftover, []byte("half an upload"), 0o600)) + } + + require.NoError(t, v.NukeRemote(true)) + + var files []string + + err = filepath.WalkDir(storeDir, + func(path string, entry fs.DirEntry, err error) error { + if err != nil { + return err + } + + if !entry.IsDir() { + files = append(files, path) + } + + return nil + }) + require.NoError(t, err) + assert.Empty(t, files) +} diff --git a/internal/vaultik/prune.go b/internal/vaultik/prune.go index 38bb135..6bd737a 100644 --- a/internal/vaultik/prune.go +++ b/internal/vaultik/prune.go @@ -24,8 +24,9 @@ var errNukeRequiresForce = errors.New( const metadataDirName = "metadata" // NukeRemote deletes every snapshot's metadata and every blob from remote -// storage. After this returns successfully the bucket prefix is empty and -// the next backup starts from scratch. +// storage, along with any object an upload cut off part-way left under a +// temporary `.partial` name. After this returns successfully the bucket +// prefix is empty and the next backup starts from scratch. // // Refuses to run unless force is true. The caller is responsible for // confirming with the user. @@ -48,6 +49,15 @@ func (v *Vaultik) NukeRemote(force bool) error { return fmt.Errorf("pruning blobs: %w", err) } + // The file and rclone listings skip `.partial` objects, so the two + // steps above never delete them. + for _, prefix := range []string{"metadata/", "blobs/"} { + err = v.Storage.DeletePartialUploads(v.ctx, prefix) + if err != nil { + return fmt.Errorf("deleting partial uploads: %w", err) + } + } + v.UI.Completef("Backup destination store is now empty.") return nil diff --git a/internal/vaultik/remove_snapshot_test.go b/internal/vaultik/remove_snapshot_test.go index bed9ace..b337c1b 100644 --- a/internal/vaultik/remove_snapshot_test.go +++ b/internal/vaultik/remove_snapshot_test.go @@ -125,6 +125,12 @@ func (s *testStorer) ListStream( return ch } +// DeletePartialUploads has nothing to remove: Put stores each object +// under its key at once. +func (s *testStorer) DeletePartialUploads(_ context.Context, _ string) error { + return nil +} + func (s *testStorer) Info() storage.Info { return storage.Info{ Type: testLabel, -- 2.54.0