diff --git a/cmd/serve/s3/backend.go b/cmd/serve/s3/backend.go index 8c135e969..337cff5ea 100644 --- a/cmd/serve/s3/backend.go +++ b/cmd/serve/s3/backend.go @@ -48,13 +48,17 @@ type s3Backend struct { // warnInMemoryOnce logs a single NOTICE the first time a multipart // upload falls back to being buffered in memory. warnInMemoryOnce sync.Once + + reaperQuit chan struct{} // closed to stop the abandoned upload reaper + reaperStop sync.Once } // newBackend creates a new SimpleBucketBackend. -func newBackend(s *Server) gofakes3.Backend { +func newBackend(s *Server) *s3Backend { return &s3Backend{ - s: s, - meta: new(sync.Map), + s: s, + meta: new(sync.Map), + reaperQuit: make(chan struct{}), } } diff --git a/cmd/serve/s3/backend_test.go b/cmd/serve/s3/backend_test.go index dfbddede1..574bbc56c 100644 --- a/cmd/serve/s3/backend_test.go +++ b/cmd/serve/s3/backend_test.go @@ -37,7 +37,7 @@ func newTestBackend(t *testing.T) (*s3Backend, string) { w, err := newServer(ctx, f, &opt, &vfscommon.Opt, &proxy.Opt) require.NoError(t, err) - return newBackend(w).(*s3Backend), root + return newBackend(w), root } const rootSecret = "ROOT_LEVEL_SECRET_MARKER" diff --git a/cmd/serve/s3/multipart.go b/cmd/serve/s3/multipart.go index 304aaf69f..869127ba0 100644 --- a/cmd/serve/s3/multipart.go +++ b/cmd/serve/s3/multipart.go @@ -27,6 +27,7 @@ import ( "sort" "strings" "sync" + "time" "github.com/google/uuid" "github.com/ncw/swift/v2" @@ -61,7 +62,9 @@ type multipartUpload struct { partMD5s map[int][]byte // raw MD5 sums per part (for the final S3 multipart ETag) partSizes map[int]int64 // observed part sizes closed bool - aborted bool // closed by abort, so nothing was committed + aborted bool // closed by abort, so nothing was committed + active int // requests in flight for this upload, protecting it from the reaper + lastUsed time.Time // when the last request for this upload finished nextPart int // next part number to stream (1-based) streamBuf map[int]*pool.RW // parts received ahead of nextPart, awaiting their turn @@ -83,11 +86,27 @@ func newMultipartUpload(bucket, key, fp, streamFp string, meta map[string]string nextPart: 1, streamBuf: map[int]*pool.RW{}, bufferLimit: bufferLimit, + lastUsed: time.Now(), } up.cond = sync.NewCond(&up.mu) return up } +// startActivity marks the upload as having a request in flight. +func (up *multipartUpload) startActivity() { + up.mu.Lock() + up.active++ + up.mu.Unlock() +} + +// endActivity marks the request done and restarts the upload's idle time. +func (up *multipartUpload) endActivity() { + up.mu.Lock() + up.active-- + up.lastUsed = time.Now() + up.mu.Unlock() +} + // loadUpload looks up an in-flight upload by ID. func (b *s3Backend) loadUpload(uploadID gofakes3.UploadID) (*multipartUpload, error) { v, ok := b.multipartUploads.Load(uploadID) @@ -174,6 +193,8 @@ func (b *s3Backend) UploadPart(ctx context.Context, bucketName, objectName strin if err != nil { return "", err } + up.startActivity() + defer up.endActivity() // Wait until there is room to buffer this part, bounding the memory a // client which uploads faster than the backend drains can consume. @@ -354,6 +375,8 @@ func (b *s3Backend) CompleteMultipartUpload(ctx context.Context, bucketName, obj if err != nil { return "", "", err } + up.startActivity() + defer up.endActivity() if err := up.validate(input); err != nil { b.multipartUploads.Delete(uploadID) @@ -416,6 +439,53 @@ func (b *s3Backend) AbortMultipartUpload(ctx context.Context, bucketName, object return nil } +// startReaper starts a goroutine which aborts incomplete multipart +// uploads once they have been idle for expiry. Stopped by stopReaper. +func (b *s3Backend) startReaper(expiry time.Duration) { + interval := min(expiry/2, time.Minute) + go func() { + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-b.reaperQuit: + return + case <-ticker.C: + b.reapExpiredUploads(time.Now(), expiry) + } + } + }() +} + +// stopReaper stops the abandoned upload reaper, if running. +func (b *s3Backend) stopReaper() { + b.reaperStop.Do(func() { + close(b.reaperQuit) + }) +} + +// reapExpiredUploads aborts and cleans up multipart uploads which have had +// no request activity for longer than expiry. +func (b *s3Backend) reapExpiredUploads(now time.Time, expiry time.Duration) { + b.multipartUploads.Range(func(key, value any) bool { + uploadID := key.(gofakes3.UploadID) + up := value.(*multipartUpload) + up.mu.Lock() + expired := up.active == 0 && now.Sub(up.lastUsed) >= expiry + up.mu.Unlock() + if !expired { + return true + } + fs.Logf(up.fp, "aborting multipart upload %s idle for more than %v", uploadID, expiry) + b.multipartUploads.Delete(uploadID) + if err := up.abort(); err != nil { + fs.Errorf(up.fp, "aborting abandoned multipart upload: %v", err) + } + b.discardUpload(up) + return true + }) +} + // discardUpload cleans up after a failed or aborted upload, removing the // temporary object and any stale VFS state for it. It never removes the // object at the final path: when the parts were written straight to the diff --git a/cmd/serve/s3/multipart_test.go b/cmd/serve/s3/multipart_test.go index a0b585b56..acf36317a 100644 --- a/cmd/serve/s3/multipart_test.go +++ b/cmd/serve/s3/multipart_test.go @@ -617,6 +617,84 @@ func TestMultipartCompleteRenameFailureKeepsUpload(t *testing.T) { require.NoError(t, err, "the upload record must survive a retryable Complete failure") } +// TestMultipartReaper checks that an incomplete multipart upload abandoned +// by its client is aborted and cleaned up after --multipart-expiry, and that +// late operations on it fail with NoSuchUpload. +func TestMultipartReaper(t *testing.T) { + core, f, bucket := newMultipartTestServerOpt(t, "", false, func(opt *Options) { + opt.MultipartExpiry = fs.Duration(100 * time.Millisecond) + }) + ctx := context.Background() + const object = "abandoned.bin" + + uploadID, err := core.NewMultipartUpload(ctx, bucket, object, minio.PutObjectOptions{}) + require.NoError(t, err) + data := []byte(random.String(50 * 1024)) + _, err = core.PutObjectPart(ctx, bucket, object, uploadID, 1, bytes.NewReader(data), int64(len(data)), minio.PutObjectPartOptions{}) + require.NoError(t, err) + + // Wait for well over the expiry and the reaper interval, then the + // upload must be gone. + time.Sleep(time.Second) + _, err = core.PutObjectPart(ctx, bucket, object, uploadID, 2, bytes.NewReader(data), int64(len(data)), minio.PutObjectPartOptions{}) + require.Error(t, err) + assert.Contains(t, err.Error(), "NoSuchUpload") + err = core.AbortMultipartUpload(ctx, bucket, object, uploadID) + require.Error(t, err) + assert.Contains(t, err.Error(), "NoSuchUpload") + + // Nothing is left at the key or as a temporary object. + _, err = f.NewObject(ctx, path.Join(bucket, object)) + require.ErrorIs(t, err, fs.ErrorObjectNotFound) + requireOnly(t, f, bucket) +} + +// TestMultipartReapExpiredUploads checks the reaper's rules directly: an +// idle upload past the expiry is aborted and cleaned up, one with a request +// in flight is left alone however stale its idle time, and a fresh one is +// kept. +func TestMultipartReapExpiredUploads(t *testing.T) { + b, _, bucket := newPutTestBackend(t, "", nil) + ctx := context.Background() + _vfs, err := b.s.getVFS(ctx) + require.NoError(t, err) + + newUp := func(id string) *multipartUpload { + up := newMultipartUpload(bucket, id, bucket+"/"+id, bucket+"/"+multipartUploadPrefix+id, nil, 0) + up.fh = &stubSink{} + up.vfs = _vfs + b.multipartUploads.Store(gofakes3.UploadID(id), up) + return up + } + + const expiry = time.Hour + now := time.Now() + + newUp("fresh") + idle := newUp("idle") + busy := newUp("busy") + busy.startActivity() + for _, up := range []*multipartUpload{idle, busy} { + up.mu.Lock() + up.lastUsed = now.Add(-2 * expiry) + up.mu.Unlock() + } + + b.reapExpiredUploads(now, expiry) + + _, err = b.loadUpload("fresh") + assert.NoError(t, err, "a fresh upload must not be reaped") + _, err = b.loadUpload("busy") + assert.NoError(t, err, "an upload with a request in flight must not be reaped") + _, err = b.loadUpload("idle") + assert.ErrorIs(t, err, gofakes3.ErrNoSuchUpload, "an idle upload past the expiry must be reaped") + + // The reaped upload was aborted, not committed. + sink := idle.fh.(*stubSink) + assert.True(t, sink.closed) + assert.Equal(t, errMultipartAborted, sink.abortErr) +} + // cacheWritesVFSOpt returns VFS options with --vfs-cache-mode writes and the // given write-back delay. func cacheWritesVFSOpt(writeBack time.Duration) *vfscommon.Options { diff --git a/cmd/serve/s3/put_test.go b/cmd/serve/s3/put_test.go index 88998b1b3..242d028b5 100644 --- a/cmd/serve/s3/put_test.go +++ b/cmd/serve/s3/put_test.go @@ -69,7 +69,7 @@ func newPutTestBackend(t *testing.T, backing string, vfsOpt *vfscommon.Options) w, err := newServer(ctx, f, &opt, vfsOpt, &proxy.Opt) require.NoError(t, err) t.Cleanup(func() { _ = w.Shutdown() }) - return newBackend(w).(*s3Backend), f, bucket + return newBackend(w), f, bucket } var errBoom = errors.New("boom") diff --git a/cmd/serve/s3/s3.go b/cmd/serve/s3/s3.go index 1b0c92e61..94cb75ca5 100644 --- a/cmd/serve/s3/s3.go +++ b/cmd/serve/s3/s3.go @@ -4,6 +4,7 @@ import ( "context" _ "embed" "strings" + "time" "github.com/rclone/rclone/cmd" "github.com/rclone/rclone/cmd/serve" @@ -39,11 +40,15 @@ var OptionsInfo = fs.Options{{ }, { Name: "disable_multipart_streaming", Default: false, - Help: "Buffer multipart uploads in memory instead of streaming them to the backend (see the Multipart uploads docs section)", + Help: "Buffer multipart uploads in memory instead of streaming them to the backend", }, { Name: "multipart_streaming_buffer_limit", Default: fs.SizeSuffix(256 * 1024 * 1024), - Help: "Maximum memory buffered per streamed multipart upload for parts arriving out of order, 0 for unlimited (see the Multipart uploads docs section)", + Help: "Maximum memory buffered per streamed multipart upload for parts arriving out of order, 0 for unlimited", +}, { + Name: "multipart_expiry", + Default: fs.Duration(24 * time.Hour), + Help: "Abort incomplete multipart uploads idle for longer than this, 0 to keep forever", }}. Add(httplib.ConfigInfo). Add(httplib.AuthConfigInfo) @@ -57,6 +62,7 @@ type Options struct { NoCleanup bool `config:"no_cleanup"` DisableMultipartStreaming bool `config:"disable_multipart_streaming"` MultipartStreamingBufferLimit fs.SizeSuffix `config:"multipart_streaming_buffer_limit"` + MultipartExpiry fs.Duration `config:"multipart_expiry"` Auth httplib.AuthConfig HTTP httplib.Config } diff --git a/cmd/serve/s3/serve_s3.md b/cmd/serve/s3/serve_s3.md index 23e89ec5d..4ba573b88 100644 --- a/cmd/serve/s3/serve_s3.md +++ b/cmd/serve/s3/serve_s3.md @@ -283,6 +283,26 @@ rclone v1.75 named its temporary multipart objects `.rclone_multipart_upload_*`; leftovers from an older server are also hidden from listings and can be cleaned up the same way. +#### Abandoned uploads + +A client which starts a multipart upload and vanishes without either +completing or aborting it would otherwise hold on to its resources +forever. + +An incomplete multipart upload which has had no activity for +`--multipart-expiry` (default `24h`) is therefore aborted and cleaned +up, exactly as if the client had called `AbortMultipartUpload`, and a +`NOTICE` is logged. + +An upload with a part still being received is never expired, however +slowly the part is arriving, and each completed part restarts the +clock, so the expiry only needs to outlast the client's pauses +*between* parts, not the whole upload. + +Late operations on an expired upload fail with `NoSuchUpload`, as they +do on real S3 when a lifecycle rule has aborted the upload. Set +`--multipart-expiry 0` to keep incomplete uploads forever. + #### Disabling streaming If you pass `--disable-multipart-streaming`, multipart uploads are diff --git a/cmd/serve/s3/server.go b/cmd/serve/s3/server.go index 260508a00..9043d7716 100644 --- a/cmd/serve/s3/server.go +++ b/cmd/serve/s3/server.go @@ -11,6 +11,7 @@ import ( "net" "net/http" "strings" + "time" "github.com/go-chi/chi/v5" "github.com/rclone/gofakes3" @@ -36,6 +37,7 @@ type Server struct { f fs.Fs _vfs *vfs.VFS // don't use directly, use getVFS faker *gofakes3.GoFakeS3 + backend *s3Backend handler http.Handler proxy *proxy.Proxy ctx context.Context // for global config @@ -75,9 +77,14 @@ func newServer(ctx context.Context, f fs.Fs, opt *Options, vfsOpt *vfscommon.Opt return nil, fmt.Errorf("parsing auth list failed: %q", err) } + w.backend = newBackend(w) + if w.opt.MultipartExpiry > 0 { + w.backend.startReaper(time.Duration(w.opt.MultipartExpiry)) + } + var newLogger logger w.faker = gofakes3.New( - newBackend(w), + w.backend, gofakes3.WithHostBucket(!opt.ForcePathStyle), gofakes3.WithLogger(newLogger), gofakes3.WithRequestID(rand.Uint64()), @@ -161,6 +168,7 @@ func (w *Server) Addr() net.Addr { // Shutdown the server func (w *Server) Shutdown() error { + w.backend.stopReaper() return w.server.Shutdown() }