Author SHA1 Message Date
clawbot 34f61fb408 Make remote nuke delete leftover .partial uploads (closes #281)
check / check (push) Waiting to run
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
2026-10-08 14:17:20 +02:00
15 changed files with 270 additions and 111 deletions
+6 -7
View File
@@ -22,13 +22,12 @@ the tag exists and is exercised; what is left is merging `next` to
# Completed Steps
- 2026-10-08: Made a backup cancelled under `--skip-errors` stop at once
([issue #286](https://git.eeqj.de/sneak/vaultik/issues/286)). After
Ctrl-C or SIGTERM, phase 2 treated the cancellation error from each
remaining file like an unreadable file: it opened the file, printed an
error line for it, counted it as failed and went on to the next one.
The run now stops at the first file after the cancellation, and no
file is counted as failed because of it.
- 2026-10-08: Made `remote nuke` delete the `.partial` files that
uploads cut off part-way leave on the destination store
([issue #281](https://git.eeqj.de/sneak/vaultik/issues/281)). The
file and rclone listings skip such a file, so the command left it in
place and still reported the store empty. It now removes them under
`metadata/` and `blobs/` as its last step.
- 2026-10-08: Made a local index error while recording a directory or
symlink stop a backup under `--skip-errors`
-12
View File
@@ -1312,13 +1312,6 @@ func (s *Scanner) processPhase(
// Process each file
for _, fileToProcess := range filesToProcess {
// Check context cancellation
select {
case <-ctx.Done():
return ctx.Err()
default:
}
// Update progress
if s.progress != nil {
s.progress.GetStats().CurrentFile.Store(fileToProcess.Path)
@@ -1376,11 +1369,6 @@ func (s *Scanner) processFileWithErrorHandling(
err := s.processFileStreaming(ctx, fileToProcess, result)
if err != nil {
// A cancelled run stops here even under --skip-errors, rather than
// counting this file as failed and going on to the next one.
if ctx.Err() != nil {
return false, fmt.Errorf("processing file %s: %w", fileToProcess.Path, err)
}
// A packer/database/encryption/upload failure means the chunk's data
// may not have been stored. Skipping the file would let the snapshot
// record a file whose chunk is in no blob and cannot be restored, so
+11 -90
View File
@@ -82,32 +82,6 @@ func (f *readFailFs) Open(name string) (afero.File, error) {
return file, nil
}
// cancelOnOpenFs cancels the run when the scanner opens the target file to
// back it up, as Ctrl-C partway through processing would, and records each
// file opened after that.
type cancelOnOpenFs struct {
afero.Fs
target string
cancel context.CancelFunc
cancelled bool
openedAfterCancel []string
}
//nolint:ireturn // afero.Fs.Open is defined to return the interface.
func (f *cancelOnOpenFs) Open(name string) (afero.File, error) {
if f.cancelled {
f.openedAfterCancel = append(f.openedAfterCancel, name)
}
if name == f.target {
f.cancel()
f.cancelled = true
}
return f.Fs.Open(name)
}
// linkRemovedAfterLstatFs is the real filesystem, except that the symlink at
// target is removed right after the walk lstats it, as happens when a link is
// deleted during a backup. The scanner's readlink of it then fails.
@@ -154,16 +128,15 @@ func writeSkipErrorTestFile(t *testing.T, fs afero.Fs, path, content string) {
}
}
// runSkipErrorScan scans source on fs under ctx 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.
// 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(
ctx context.Context, t *testing.T, fs afero.Fs, source string,
skipErrors bool, uiw *ui.Writer,
t *testing.T, fs afero.Fs, source string, skipErrors bool, uiw *ui.Writer,
) (*database.Repositories, error) {
t.Helper()
db, err := database.New(ctx, ":memory:")
db, err := database.NewTestDB()
if err != nil {
t.Fatalf("create test db: %v", err)
}
@@ -188,6 +161,7 @@ func runSkipErrorScan(
SkipErrors: skipErrors,
})
ctx := context.Background()
snapshotID := "test-snapshot-skip-errors"
createTestSnapshotRecord(ctx, t, repos, snapshotID)
@@ -211,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(context.Background(), t, fs, "/source", true, nil)
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
if err == nil {
t.Fatal("expected scan to abort on the packer error, got nil")
}
@@ -238,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(context.Background(), t, fs, "/source", false, nil)
_, err := runSkipErrorScan(t, fs, "/source", false, nil)
if err == nil {
t.Fatal("expected scan to fail on the read error, got nil")
}
@@ -254,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(context.Background(), t, fs, "/source", true, nil)
repos, err := runSkipErrorScan(t, fs, "/source", true, nil)
if err != nil {
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
}
@@ -293,7 +267,7 @@ func TestScannerUnreadableSymlinkAbortsWithoutSkipErrors(t *testing.T) {
sourceDir, linkPath := writeSymlinkSource(t)
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
_, err := runSkipErrorScan(context.Background(), t, fs, sourceDir, false, nil)
_, 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)
}
@@ -309,7 +283,7 @@ func TestScannerUnreadableSymlinkSkippedWithSkipErrors(t *testing.T) {
fs := &linkRemovedAfterLstatFs{t: t, target: linkPath}
uiw := ui.NewWithColor(io.Discard, false)
repos, err := runSkipErrorScan(context.Background(), t, fs, sourceDir, true, uiw)
repos, err := runSkipErrorScan(t, fs, sourceDir, true, uiw)
if err != nil {
t.Fatalf("expected scan to complete with --skip-errors, got %v", err)
}
@@ -328,56 +302,3 @@ func TestScannerUnreadableSymlinkSkippedWithSkipErrors(t *testing.T) {
t.Fatalf("expected %s not to be recorded", linkPath)
}
}
// TestScannerCancelStopsSkipErrorsRun checks that a --skip-errors backup
// cancelled partway through processing returns the cancellation error,
// opens no further file, and reports no file as failed.
func TestScannerCancelStopsSkipErrorsRun(t *testing.T) {
t.Parallel()
tests := []struct {
name string
targetContent string
}{
// The cancellation lands while the target is being read.
{name: "while reading a file", targetContent: "first file content"},
// An empty target has no chunks, so the cancellation goes unnoticed
// until the run moves on to the next file.
{name: "between files", targetContent: ""},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
const target = "/source/a.txt"
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
fs := &cancelOnOpenFs{
Fs: afero.NewMemMapFs(), target: target, cancel: cancel,
}
writeSkipErrorTestFile(t, fs, target, tt.targetContent)
writeSkipErrorTestFile(t, fs, "/source/b.txt", "second file content")
writeSkipErrorTestFile(t, fs, "/source/c.txt", "third file content")
uiw := ui.NewWithColor(io.Discard, false)
_, err := runSkipErrorScan(ctx, t, fs, "/source", true, uiw)
if !errors.Is(err, context.Canceled) {
t.Fatalf("expected the cancellation error, got %v", err)
}
if len(fs.openedAfterCancel) != 0 {
t.Fatalf("expected no file opened after the cancellation, got %v",
fs.openedAfterCancel)
}
if uiw.ErrorCount() != 0 {
t.Fatalf("expected no file reported as failed, got %d error lines",
uiw.ErrorCount())
}
})
}
}
@@ -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()
+36
View File
@@ -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{
+42
View File
@@ -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)
}
}
+25
View File
@@ -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
+41
View File
@@ -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
+6
View File
@@ -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{
+6
View File
@@ -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
}
@@ -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{}
}
+6
View File
@@ -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",
+64
View File
@@ -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)
}
+12 -2
View File
@@ -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
+6
View File
@@ -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,