From 89e609f0635dd69c602a37bacf024130fccd807c 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 and move them into place (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 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 --- README.md | 10 ++ TODO.md | 11 ++ internal/storage/conformance_test.go | 5 +- internal/storage/file.go | 7 +- internal/storage/file_atomic_test.go | 7 +- internal/storage/rclone.go | 79 ++++++++-- internal/storage/rclone_test.go | 227 +++++++++++++++++++++++++-- 7 files changed, 310 insertions(+), 36 deletions(-) diff --git a/README.md b/README.md index 529ebd1..78eaea7 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/TODO.md b/TODO.md index 076aae6..ae7d39d 100644 --- a/TODO.md +++ b/TODO.md @@ -22,6 +22,17 @@ the tag exists and is exercised; what is left is merging `next` to # Completed Steps +- 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: 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/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..7cb5653 100644 --- a/internal/storage/file.go +++ b/internal/storage/file.go @@ -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. 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..dfc8e1f 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,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 diff --git a/internal/storage/rclone_test.go b/internal/storage/rclone_test.go index e941497..d73f114 100644 --- a/internal/storage/rclone_test.go +++ b/internal/storage/rclone_test.go @@ -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,186 @@ func TestNewRcloneStorerConstruction(t *testing.T) { } } +// 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) { + // 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) + } +} + +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) { + fs.Register(&fs.RegInfo{ + Name: "moverefusesexisting", + 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 &moveRefusesExisting{Fs: local}, nil + }, + }) + + ctx := context.Background() + + s, err := storage.NewRcloneStorer(ctx, ":moverefusesexisting", t.TempDir()) + if err != nil { + t.Fatalf("NewRcloneStorer: %v", err) + } + + for _, content := range []string{"first", "second"} { + err = s.Put(ctx, testBlobKey, strings.NewReader(content)) + if err != nil { + t.Fatalf("Put %q: %v", content, err) + } + } + + rc, err := s.Get(ctx, testBlobKey) + 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) + } + + 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.