diff --git a/cmd/serve/s3/backend.go b/cmd/serve/s3/backend.go index 292fdb114..913f0dc39 100644 --- a/cmd/serve/s3/backend.go +++ b/cmd/serve/s3/backend.go @@ -12,16 +12,23 @@ import ( "sync" "time" + "github.com/google/uuid" "github.com/ncw/swift/v2" "github.com/rclone/gofakes3" "github.com/rclone/rclone/fs" + "github.com/rclone/rclone/fs/operations" "github.com/rclone/rclone/vfs" + "github.com/rclone/rclone/vfs/vfscommon" ) var ( emptyPrefix = &gofakes3.Prefix{} ) +// putObjectPrefix is prepended to the leaf name of the temporary object a +// PutObject upload is written to before it is renamed into place. +const putObjectPrefix = ".rclone_put_object_" + // s3Backend implements the gofacess3.Backend interface to make an S3 // backend for gofakes3. It also implements gofakes3.MultipartBackend so that // multipart uploads stream straight through to the underlying Fs via @@ -311,6 +318,9 @@ func (b *s3Backend) TouchObject(ctx context.Context, fp string, meta map[string] } // PutObject creates or overwrites the object with the given name. +// +// A failed or interrupted upload should never disturb what is stored +// at the key. func (b *s3Backend) PutObject( ctx context.Context, bucketName, objectName string, @@ -343,24 +353,63 @@ func (b *s3Backend) PutObject( } } - f, err := _vfs.Create(fp) + // Upload via a temporary object renamed into place if needed. The handle + // is opened read-write, which the VFS routes through the cache - where a + // failed upload commits what it has on Close and can't be abandoned - + // from --vfs-cache-mode minimal up, not just writes. + fsys := _vfs.Fs() + tmpFp := fp + if (fsys.Features().PartialUploads || _vfs.Opt.CacheMode >= vfscommon.CacheModeMinimal) && operations.CanServerSideMove(fsys) { + tmpFp = path.Join(objectDir, putObjectPrefix+uuid.New().String()) + } + + // cleanup discards a failed upload, removing the temporary object + // (never the object at fp) and any stale VFS state for fp. + cleanup := func() { + if tmpFp != fp { + b.forgetPath(ctx, tmpFp) + _ = _vfs.Remove(tmpFp) + } else { + b.forgetPath(ctx, fp) + } + } + + f, err := _vfs.Create(tmpFp) if err != nil { return result, err } - if _, err := io.Copy(f, input); err != nil { - // remove file when i/o error occurred (FsPutErr) - _ = f.Close() - _ = _vfs.Remove(fp) + n, err := io.Copy(f, input) + if err == nil && size >= 0 && n != size { + // The body ended cleanly but short of its declared size + err = gofakes3.ErrIncompleteBody + } + if err != nil { + // The upload from the client failed part way through - abort the + // write so a streaming upload fails rather than committing a + // truncated object. + if aborter, ok := f.(interface{ CloseWithError(error) error }); ok { + _ = aborter.CloseWithError(err) + } else { + _ = f.Close() + } + cleanup() return result, err } if err := f.Close(); err != nil { - // remove file when close error occurred (FsPutErr) - _ = _vfs.Remove(fp) + cleanup() return result, err } + // Rename the temporary object into place + if tmpFp != fp { + if err := _vfs.Rename(tmpFp, fp); err != nil { + cleanup() + return result, err + } + } + _, err = _vfs.Stat(fp) if err != nil { return result, err @@ -541,6 +590,10 @@ func (b *s3Backend) CopyObject(ctx context.Context, srcBucket, srcKey, dstBucket meta["mtime"] = swift.TimeToFloatString(cStat.ModTime()) } + // PutObject rejects a body shorter than its declared size, so a copy + // whose source is being overwritten as it is read - its stated size no + // longer matching its readable bytes - fails with IncompleteBody rather + // than storing a mixture. _, err = b.PutObject(ctx, dstBucket, dstKey, meta, c.Contents, c.Size) if err != nil { return diff --git a/cmd/serve/s3/list.go b/cmd/serve/s3/list.go index 76908e8c0..a95bb8dee 100644 --- a/cmd/serve/s3/list.go +++ b/cmd/serve/s3/list.go @@ -23,8 +23,8 @@ func (b *s3Backend) entryListR(_vfs *vfs.VFS, bucketName, fdPath, name string, a for _, entry := range dirEntries { object := entry.Name() - // Hide the in-progress multipart uploads - if strings.HasPrefix(object, multipartUploadPrefix) { + // Hide the temporary objects of in-progress uploads + if strings.HasPrefix(object, multipartUploadPrefix) || strings.HasPrefix(object, putObjectPrefix) { continue } diff --git a/cmd/serve/s3/put_test.go b/cmd/serve/s3/put_test.go new file mode 100644 index 000000000..5cb266b63 --- /dev/null +++ b/cmd/serve/s3/put_test.go @@ -0,0 +1,201 @@ +// PutObject error handling tests for serve s3. +// +// A PUT which fails part way through its body must never disturb the object +// already stored at the key. + +package s3 + +import ( + "bytes" + "context" + "errors" + "fmt" + "io" + "path" + "testing" + "time" + + "github.com/rclone/gofakes3" + "github.com/rclone/rclone/cmd/serve/proxy" + "github.com/rclone/rclone/fs" + "github.com/rclone/rclone/fstest" + "github.com/rclone/rclone/lib/random" + "github.com/rclone/rclone/vfs" + "github.com/rclone/rclone/vfs/vfscommon" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// errorReader yields data then fails with err. +type errorReader struct { + data []byte + err error +} + +func (r *errorReader) Read(p []byte) (int, error) { + if len(r.data) == 0 { + return 0, r.err + } + n := copy(p, r.data) + r.data = r.data[n:] + return n, nil +} + +// newPutTestBackend returns an s3Backend for direct PutObject tests, backed by +// the named remote, plus the backing Fs and a bucket created on it. vfsOpt +// overrides the VFS options (nil for the defaults). +func newPutTestBackend(t *testing.T, backing string, vfsOpt *vfscommon.Options) (*s3Backend, fs.Fs, string) { + fstest.Initialise() + ctx := context.Background() + if backing == "" { + backing = t.TempDir() + } + f, err := fs.NewFs(ctx, backing) + require.NoError(t, err) + bucket := fmt.Sprintf("test-%d", testBackingCounter.Add(1)) + require.NoError(t, f.Mkdir(ctx, bucket)) + if vfsOpt == nil { + vfsOpt = &vfscommon.Opt + } + // The VFS is cached per remote, so a shared ":memory:" backing reuses a + // VFS whose cached root listing predates the bucket just created; forget + // it so the new bucket is visible. + if root, err := vfs.New(ctx, f, vfsOpt).Root(); err == nil { + root.ForgetAll() + } + opt := Opt + opt.HTTP.ListenAddr = []string{endpoint} + w, err := newServer(ctx, f, &opt, vfsOpt, &proxy.Opt) + require.NoError(t, err) + t.Cleanup(func() { _ = w.Shutdown() }) + return newBackend(w).(*s3Backend), f, bucket +} + +var errBoom = errors.New("boom") + +// failureModes enumerates the ways a PUT body can fail to arrive in full: +// an arbitrary read error, the unexpected EOF a dropped connection gives +// (which must not be mistaken for a clean end of stream), and a body which +// ends cleanly short of its declared size. +var failureModes = []struct { + name string + readerErr error + wantErr error +}{ + {"ReadError", errBoom, errBoom}, + {"UnexpectedEOF", io.ErrUnexpectedEOF, io.ErrUnexpectedEOF}, + {"ShortBody", io.EOF, gofakes3.ErrIncompleteBody}, +} + +// failPut makes a PutObject call to bucket/object whose body fails with +// readerErr part way through and asserts wantErr is passed back. +func failPut(t *testing.T, b *s3Backend, bucket, object string, readerErr, wantErr error) { + _, err := b.PutObject(context.Background(), bucket, object, map[string]string{}, + &errorReader{data: []byte(random.String(50)), err: readerErr}, 1000) + require.ErrorIs(t, err, wantErr) +} + +// TestPutObjectFailurePreservesExisting checks that a failed PUT leaves the +// object stored at the key untouched - neither removed nor overwritten with +// truncated data - with no temporary object left behind. +func TestPutObjectFailurePreservesExisting(t *testing.T) { + for _, tc := range testRemotes { + for _, fm := range failureModes { + t.Run(tc.name+"/"+fm.name, func(t *testing.T) { + b, f, bucket := newPutTestBackend(t, tc.backing, nil) + ctx := context.Background() + const object = "existing.txt" + + existing := []byte(random.String(100)) + _, err := b.PutObject(ctx, bucket, object, map[string]string{}, bytes.NewReader(existing), int64(len(existing))) + require.NoError(t, err) + assert.Equal(t, existing, readObject(t, f, bucket, object)) + requireOnly(t, f, bucket, object) + + failPut(t, b, bucket, object, fm.readerErr, fm.wantErr) + + assert.Equal(t, existing, readObject(t, f, bucket, object)) + requireOnly(t, f, bucket, object) + + // The object must still be served correctly - in particular + // HeadObject must not report the failed upload's size. + head, err := b.HeadObject(ctx, bucket, object) + require.NoError(t, err) + assert.Equal(t, int64(len(existing)), head.Size) + }) + } + } +} + +// TestPutObjectFailureNewKey checks that a failed PUT to a key with no +// existing object leaves nothing behind - no partial object at the key and +// no temporary object. +func TestPutObjectFailureNewKey(t *testing.T) { + for _, tc := range testRemotes { + for _, fm := range failureModes { + t.Run(tc.name+"/"+fm.name, func(t *testing.T) { + b, f, bucket := newPutTestBackend(t, tc.backing, nil) + ctx := context.Background() + const object = "new.txt" + + failPut(t, b, bucket, object, fm.readerErr, fm.wantErr) + + _, err := f.NewObject(ctx, path.Join(bucket, object)) + require.ErrorIs(t, err, fs.ErrorObjectNotFound) + requireOnly(t, f, bucket) + + // The failed upload must not be served either + _, err = b.HeadObject(ctx, bucket, object) + require.Error(t, err) + }) + } + } +} + +// waitForObject waits for bucket/object to appear on the backing Fs (e.g. +// after the VFS write-back delay). +func waitForObject(t *testing.T, f fs.Fs, bucket, object string) { + ctx := context.Background() + deadline := time.Now().Add(10 * time.Second) + for time.Now().Before(deadline) { + if _, err := f.NewObject(ctx, path.Join(bucket, object)); err == nil { + return + } + time.Sleep(50 * time.Millisecond) + } + t.Fatalf("object %s/%s never appeared on the backing remote", bucket, object) +} + +// TestPutObjectFailureCached checks the failed PUT semantics with the cache +// modes that write through the VFS cache - minimal and above, since the PUT +// opens its handle read-write: the truncated file the failed upload leaves +// in the cache must be discarded, not written back over the object at the +// key. +func TestPutObjectFailureCached(t *testing.T) { + for _, cm := range []vfscommon.CacheMode{vfscommon.CacheModeMinimal, vfscommon.CacheModeWrites} { + for _, tc := range testRemotes { + t.Run(cm.String()+"/"+tc.name, func(t *testing.T) { + vfsOpt := vfscommon.Opt + vfsOpt.CacheMode = cm + vfsOpt.WriteBack = fs.Duration(100 * time.Millisecond) + b, f, bucket := newPutTestBackend(t, tc.backing, &vfsOpt) + const object = "cached.txt" + + existing := []byte(random.String(100)) + _, err := b.PutObject(context.Background(), bucket, object, map[string]string{}, bytes.NewReader(existing), int64(len(existing))) + require.NoError(t, err) + waitForObject(t, f, bucket, object) + assert.Equal(t, existing, readObject(t, f, bucket, object)) + + failPut(t, b, bucket, object, errBoom, errBoom) + + // Wait out several write-back intervals to catch the truncated + // data being written back before checking the object is + // untouched. + time.Sleep(time.Second) + assert.Equal(t, existing, readObject(t, f, bucket, object)) + requireOnly(t, f, bucket, object) + }) + } + } +} diff --git a/cmd/serve/s3/serve_s3.md b/cmd/serve/s3/serve_s3.md index 3c8f9df5a..4d214a0f0 100644 --- a/cmd/serve/s3/serve_s3.md +++ b/cmd/serve/s3/serve_s3.md @@ -88,6 +88,24 @@ access_key_id = ACCESS_KEY_ID secret_access_key = SECRET_ACCESS_KEY ``` +### Object uploads (PUT) + +A `PutObject` upload only ever changes the object at its key atomically, on +success, a failed or interrupted PUT neither removes nor overwrites the +object already stored at the key, and never leaves a partial object visible +at it. + +Remotes that upload atomically (e.g. object stores such as `s3`) are streamed +straight to the destination. On remotes where a partial upload would +otherwise be visible (e.g. `local`), and whenever `--vfs-cache-mode` is +`writes` or above, the upload is written to a temporary object that is +renamed into place on success; these remotes need to support a server-side +move or copy for this (nearly all do - without move or copy the upload is +written directly and a failed PUT may leave a partial object at the key). If +`serve s3` is killed part-way through an upload the temporary object (named +with a leading `.rclone_put_object_`) may be left behind; it is hidden from +S3 listings but must be removed manually. + ### Multipart uploads By default `serve s3` **streams** each multipart upload, in part-number