From 87ee3b63a287f7a08b9280b8037efd5f1ffb3e56 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Fri, 18 Sep 2026 01:01:52 -0700 Subject: [PATCH] s3: abort completed multipart uploads metadata-only (#11385) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * s3: abort a completed upload's leftover directory metadata-only A .uploads/ directory can outlive the object it completed into when the commit's metadata-only removal failed or the gateway died in between; the restored part entries then share chunks with the published object. AbortMultipartUpload deleted the directory recursively, chunks and all, so aborting such a leftover destroyed a committed object (#11382). Run the same check s3.clean.uploads gained in #11375 before deleting: when the object entry or a version file under .versions carries the upload id, remove .uploads/ metadata-only and answer the abort; when the lookup cannot decide, refuse with InternalError rather than risk live chunks. * s3: apply the completed-upload check to lifecycle MPU abort lifecycleAbortMPU ran the same destructive recursive delete on .uploads/. Reuse uploadCompleted so a leftover whose object entry or version file carries the upload id is removed metadata-only, and an undecidable lookup retries later instead of freeing live chunks. * s3: serialize abort's upload-dir delete with the object's commit The completed check alone leaves a race: abort can read completed=false, then an in-flight completion publishes the object over the same part chunks before the recursive delete frees them. Run the check and delete inside the object write lock, which non-routed completions hold for their whole finalize. With an owner, send the data delete as an ObjectTransaction on the object's lock key — a routed commit then either loses its upload-exists precondition after our delete or has already stamped the object, which the transaction's IF_EXTENDED_NOT_EQUAL condition detects and falls back to a metadata-only remove. lifecycleAbortMPU shares removeUploadDir so both callers get the same ordering. * s3: check for an empty object before resolving its write owner * s3: check completion at the abort's resolved object key An upload record missing ExtMultipartObjectKey skipped the completed check entirely even though the request's Key names the object. --- weed/s3api/filer_multipart.go | 134 +++++++++++- weed/s3api/filer_multipart_abort_test.go | 261 +++++++++++++++++++++++ weed/s3api/s3api_internal_lifecycle.go | 21 +- 3 files changed, 396 insertions(+), 20 deletions(-) create mode 100644 weed/s3api/filer_multipart_abort_test.go diff --git a/weed/s3api/filer_multipart.go b/weed/s3api/filer_multipart.go index 046e0aab0..67272b3cd 100644 --- a/weed/s3api/filer_multipart.go +++ b/weed/s3api/filer_multipart.go @@ -1104,22 +1104,136 @@ func (s3a *S3ApiServer) abortMultipartUpload(input *s3.AbortMultipartUploadInput glog.V(2).Infof("abortMultipartUpload input %v", input) - exists, err := s3a.exists(s3a.genUploadsFolder(*input.Bucket), *input.UploadId, true) + uploadsFolder := s3a.genUploadsFolder(*input.Bucket) + uploadEntry, err := s3a.getEntry(uploadsFolder, *input.UploadId) if err != nil { - // filer_pb.Exists reports not-found as (false, nil), so an error here is - // always a store failure; answering NoSuchUpload would leak the parts. + if isFilerNotFound(err) { + return &s3.AbortMultipartUploadOutput{}, s3err.ErrNone + } glog.Errorf("bucket %s abort upload %s: %v", *input.Bucket, *input.UploadId, err) return nil, s3err.ErrInternalError } - if exists { - err = s3a.rm(context.Background(), s3a.genUploadsFolder(*input.Bucket), *input.UploadId, true, true) - } - if err != nil { - glog.V(1).Infof("bucket %s remove upload %s: %v", *input.Bucket, *input.UploadId, err) - return nil, s3err.ErrInternalError + if uploadEntry == nil { + return &s3.AbortMultipartUploadOutput{}, s3err.ErrNone } - return &s3.AbortMultipartUploadOutput{}, s3err.ErrNone + object := string(uploadEntry.Extended[s3_constants.ExtMultipartObjectKey]) + if object == "" && input.Key != nil { + object = *input.Key + } + return &s3.AbortMultipartUploadOutput{}, s3a.withObjectWriteLock(*input.Bucket, object, nil, func() s3err.ErrorCode { + return s3a.removeUploadDir(*input.Bucket, *input.UploadId, object) + }) +} + +// removeUploadDir removes a leftover upload directory under the object write +// lock: a directory that outlived the object it completed into shares chunks +// with it and goes metadata-only; anything else frees the parts' chunks. +func (s3a *S3ApiServer) removeUploadDir(bucket, uploadId, object string) s3err.ErrorCode { + completed, err := s3a.uploadCompleted(s3a.bucketDir(bucket), uploadId, object) + if err != nil { + glog.Errorf("bucket %s remove upload %s completed check: %v", bucket, uploadId, err) + return s3err.ErrInternalError + } + if completed { + if err := s3a.rm(context.Background(), s3a.genUploadsFolder(bucket), uploadId, false, true); err != nil { + glog.V(1).Infof("bucket %s remove upload %s: %v", bucket, uploadId, err) + return s3err.ErrInternalError + } + return s3err.ErrNone + } + return s3a.routedUploadDirDelete(bucket, uploadId, object) +} + +// routedUploadDirDelete frees an open upload's chunks. The delete rides an +// ObjectTransaction on the object's lock key, so a racing routed commit either +// fails its upload-exists precondition afterwards or has already stamped the +// object — the condition then rejects the data delete and it falls back to +// metadata-only. With no owner there is no routed commit to exclude, and the +// caller's object write lock already serializes the mkFile fallback. +func (s3a *S3ApiServer) routedUploadDirDelete(bucket, uploadId, object string) s3err.ErrorCode { + uploadsFolder := s3a.genUploadsFolder(bucket) + rmUploadDir := func(isDeleteData bool) s3err.ErrorCode { + if err := s3a.rm(context.Background(), uploadsFolder, uploadId, isDeleteData, true); err != nil { + glog.V(1).Infof("bucket %s remove upload %s: %v", bucket, uploadId, err) + return s3err.ErrInternalError + } + return s3err.ErrNone + } + if object == "" { + return rmUploadDir(true) + } + owner := s3a.objectWriteOwner(bucket, object) + if owner == "" { + return rmUploadDir(true) + } + objectPath := s3a.toFilerPath(bucket, object) + resp, err := s3a.objectTxnOnFiler(owner, &filer_pb.ObjectTransactionRequest{ + LockKey: objectPath, + RouteKey: s3a.objectRouteKey(bucket, object), + ConditionKey: objectPath, + Condition: &filer_pb.WriteCondition{Clauses: []*filer_pb.WriteCondition_Clause{{ + Kind: filer_pb.WriteCondition_IF_EXTENDED_NOT_EQUAL, + ExtKey: s3_constants.SeaweedFSUploadId, + ExtValue: uploadId, + }}}, + Mutations: []*filer_pb.ObjectMutation{{ + Type: filer_pb.ObjectMutation_DELETE, + Directory: uploadsFolder, + Name: uploadId, + IsDeleteData: true, + IsRecursive: true, + }}, + }) + if err != nil { + glog.Errorf("bucket %s abort upload %s transaction: %v", bucket, uploadId, err) + return s3err.ErrInternalError + } + if resp.ErrorCode == filer_pb.FilerError_PRECONDITION_FAILED { + return rmUploadDir(false) + } + if resp.Error != "" { + glog.Errorf("bucket %s abort upload %s transaction: %s", bucket, uploadId, resp.Error) + return s3err.ErrInternalError + } + return s3err.ErrNone +} + +// uploadCompleted reports whether the upload assembled into an object: the +// object entry, or any version file under .versions, still carries the +// upload id completion stamps on it. +func (s3a *S3ApiServer) uploadCompleted(bucketDir, uploadId, objectKey string) (bool, error) { + if objectKey == "" { + return false, nil + } + name := path.Base(objectKey) + dir := path.Dir(objectKey) + if dir == "." { + dir = "" + } + objectDir := path.Join(bucketDir, dir) + + entry, err := s3a.getEntry(objectDir, name) + if err != nil && !isFilerNotFound(err) { + return false, err + } + if entry != nil && string(entry.Extended[s3_constants.SeaweedFSUploadId]) == uploadId { + return true, nil + } + + versions, _, err := s3a.list(objectDir+"/"+name+s3_constants.VersionsFolder, "", "", false, math.MaxInt32) + if err != nil { + if isFilerNotFound(err) { + return false, nil + } + return false, err + } + for _, version := range versions { + if string(version.Extended[s3_constants.SeaweedFSUploadId]) == uploadId { + return true, nil + } + } + return false, nil } type ListMultipartUploadsResult struct { diff --git a/weed/s3api/filer_multipart_abort_test.go b/weed/s3api/filer_multipart_abort_test.go new file mode 100644 index 000000000..37f3945af --- /dev/null +++ b/weed/s3api/filer_multipart_abort_test.go @@ -0,0 +1,261 @@ +package s3api + +import ( + "context" + "testing" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/pb/s3_lifecycle_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// fakeAbortFiler answers the three calls abortMultipartUpload makes: the +// .uploads/ record lookup, the completed-object lookup, the .versions +// listing, and records the DeleteEntry request it receives. +type fakeAbortFiler struct { + filer_pb.UnimplementedSeaweedFilerServer + uploadsDir string + uploadEntry *filer_pb.Entry + objectDir string + objectName string + objectEntry *filer_pb.Entry + objectErr error + versionsDir string + versions []*filer_pb.Entry + listErr error + deleteReq *filer_pb.DeleteEntryRequest +} + +func (f *fakeAbortFiler) LookupDirectoryEntry(ctx context.Context, req *filer_pb.LookupDirectoryEntryRequest) (*filer_pb.LookupDirectoryEntryResponse, error) { + switch req.Directory { + case f.uploadsDir: + if f.uploadEntry != nil && req.Name == f.uploadEntry.Name { + return &filer_pb.LookupDirectoryEntryResponse{Entry: f.uploadEntry}, nil + } + return nil, filer_pb.ErrNotFound + case f.objectDir: + if f.objectErr != nil { + return nil, f.objectErr + } + if req.Name == f.objectName && f.objectEntry != nil { + return &filer_pb.LookupDirectoryEntryResponse{Entry: f.objectEntry}, nil + } + return nil, filer_pb.ErrNotFound + } + return nil, status.Errorf(codes.Internal, "unexpected lookup in %s", req.Directory) +} + +func (f *fakeAbortFiler) ListEntries(req *filer_pb.ListEntriesRequest, stream filer_pb.SeaweedFiler_ListEntriesServer) error { + if req.Directory != f.versionsDir { + return status.Errorf(codes.Internal, "unexpected listing of %s", req.Directory) + } + if f.listErr != nil { + return f.listErr + } + for _, entry := range f.versions { + if err := stream.Send(&filer_pb.ListEntriesResponse{Entry: entry}); err != nil { + return err + } + } + return nil +} + +func (f *fakeAbortFiler) DeleteEntry(ctx context.Context, req *filer_pb.DeleteEntryRequest) (*filer_pb.DeleteEntryResponse, error) { + f.deleteReq = req + return &filer_pb.DeleteEntryResponse{}, nil +} + +func newAbortTestServer(t *testing.T, f *fakeAbortFiler) *S3ApiServer { + t.Helper() + bucketDir := (&S3ApiServer{option: &S3ApiServerOption{}}).bucketDir("b") + f.uploadsDir = bucketDir + "/" + s3_constants.MultipartUploadsFolder + f.objectDir = bucketDir + f.objectName = "a.bin" + f.versionsDir = bucketDir + "/a.bin" + s3_constants.VersionsFolder + return newFailoverTestServer(t, startFakeFiler(t, f)) +} + +func abortInput(uploadId string) *s3.AbortMultipartUploadInput { + return &s3.AbortMultipartUploadInput{ + Bucket: aws.String("b"), + Key: aws.String("a.bin"), + UploadId: aws.String(uploadId), + } +} + +func versionEntry(name, uploadId string) *filer_pb.Entry { + return &filer_pb.Entry{ + Name: name, + Extended: map[string][]byte{s3_constants.SeaweedFSUploadId: []byte(uploadId)}, + } +} + +// A leftover .uploads/ whose object entry still carries the upload id +// shares chunks with that object; abort must drop the metadata only. +func TestAbortCompletedUploadDeletesMetadataOnly(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + objectEntry: versionEntry("a.bin", "up1"), + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrNone { + t.Fatalf("code = %v, want ErrNone", code) + } + if f.deleteReq == nil || f.deleteReq.IsDeleteData { + t.Fatalf("deleteReq = %+v, want IsDeleteData=false", f.deleteReq) + } +} + +// A version file carrying the upload id is the same case: the versioned +// object's chunks are the part chunks. +func TestAbortCompletedVersionDeletesMetadataOnly(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + versions: []*filer_pb.Entry{versionEntry("v_123", "up1")}, + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrNone { + t.Fatalf("code = %v, want ErrNone", code) + } + if f.deleteReq == nil || f.deleteReq.IsDeleteData { + t.Fatalf("deleteReq = %+v, want IsDeleteData=false", f.deleteReq) + } +} + +// An upload record missing its object-key stamp can still have completed; +// the abort's Key is the fallback lookup path. +func TestAbortCompletedUploadWithoutRecordedKey(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: &filer_pb.Entry{Name: "up1", IsDirectory: true}, + objectEntry: versionEntry("a.bin", "up1"), + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrNone { + t.Fatalf("code = %v, want ErrNone", code) + } + if f.deleteReq == nil || f.deleteReq.IsDeleteData { + t.Fatalf("deleteReq = %+v, want IsDeleteData=false", f.deleteReq) + } +} + +// An upload that never completed owns its part chunks; abort frees them. +func TestAbortOpenUploadDeletesData(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + objectEntry: versionEntry("a.bin", "other-upload"), + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrNone { + t.Fatalf("code = %v, want ErrNone", code) + } + if f.deleteReq == nil || !f.deleteReq.IsDeleteData { + t.Fatalf("deleteReq = %+v, want IsDeleteData=true", f.deleteReq) + } +} + +// When the completed check cannot decide, abort must refuse rather than +// risk freeing chunks a live object references. +func TestAbortUndecidableCheckRefuses(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + objectErr: status.Error(codes.Unavailable, "store down"), + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrInternalError { + t.Fatalf("code = %v, want ErrInternalError", code) + } + if f.deleteReq != nil { + t.Fatalf("delete issued despite undecidable check: %+v", f.deleteReq) + } +} + +// No upload directory at all answers success, the same as before. +func TestAbortGoneUploadAnswersSuccess(t *testing.T) { + s3a := newAbortTestServer(t, &fakeAbortFiler{}) + + _, code := s3a.abortMultipartUpload(abortInput("gone")) + + if code != s3err.ErrNone { + t.Fatalf("code = %v, want ErrNone", code) + } +} + +// A failed .versions listing is undecidable the same way a failed object +// lookup is: refuse instead of guessing. +func TestAbortVersionsListErrorRefuses(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + listErr: status.Error(codes.Internal, "store down"), + } + s3a := newAbortTestServer(t, f) + + _, code := s3a.abortMultipartUpload(abortInput("up1")) + + if code != s3err.ErrInternalError { + t.Fatalf("code = %v, want ErrInternalError", code) + } + if f.deleteReq != nil { + t.Fatalf("delete issued despite undecidable check: %+v", f.deleteReq) + } +} + +func lifecycleAbortRequest(uploadId string) *s3_lifecycle_pb.LifecycleDeleteRequest { + return &s3_lifecycle_pb.LifecycleDeleteRequest{ + Bucket: "b", + ObjectPath: s3_constants.MultipartUploadsFolder + "/" + uploadId, + } +} + +func TestLifecycleAbortCompletedUploadDeletesMetadataOnly(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + objectEntry: versionEntry("a.bin", "up1"), + } + s3a := newAbortTestServer(t, f) + + resp, err := s3a.lifecycleAbortMPU(context.Background(), lifecycleAbortRequest("up1")) + + if err != nil || resp.Outcome != s3_lifecycle_pb.LifecycleDeleteOutcome_DONE { + t.Fatalf("resp = %v, err = %v, want DONE", resp, err) + } + if f.deleteReq == nil || f.deleteReq.IsDeleteData { + t.Fatalf("deleteReq = %+v, want IsDeleteData=false", f.deleteReq) + } +} + +func TestLifecycleAbortUndecidableCheckRetriesLater(t *testing.T) { + f := &fakeAbortFiler{ + uploadEntry: uploadRecordEntry("up1"), + objectErr: status.Error(codes.Unavailable, "store down"), + } + s3a := newAbortTestServer(t, f) + + resp, err := s3a.lifecycleAbortMPU(context.Background(), lifecycleAbortRequest("up1")) + + if err != nil || resp.Outcome != s3_lifecycle_pb.LifecycleDeleteOutcome_RETRY_LATER { + t.Fatalf("resp = %v, err = %v, want RETRY_LATER", resp, err) + } + if f.deleteReq != nil { + t.Fatalf("delete issued despite undecidable check: %+v", f.deleteReq) + } +} diff --git a/weed/s3api/s3api_internal_lifecycle.go b/weed/s3api/s3api_internal_lifecycle.go index 79950dced..20e76b62c 100644 --- a/weed/s3api/s3api_internal_lifecycle.go +++ b/weed/s3api/s3api_internal_lifecycle.go @@ -12,6 +12,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/pb/s3_lifecycle_pb" "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" "github.com/seaweedfs/seaweedfs/weed/s3api/s3lifecycle" stats_collect "github.com/seaweedfs/seaweedfs/weed/stats" ) @@ -172,24 +173,24 @@ func (s3a *S3ApiServer) lifecycleAbortMPU(ctx context.Context, req *s3_lifecycle // Pre-check existence: filer.DeleteEntry suppresses ErrNotFound and // returns success, so without this check an already-aborted upload // would report DONE instead of the correct NOOP_RESOLVED. - exists, err := s3a.exists(uploadsFolder, uploadID, true) + uploadEntry, err := s3a.getEntry(uploadsFolder, uploadID) if err != nil { if errors.Is(err, filer_pb.ErrNotFound) { return noopResolved("NOT_FOUND"), nil } - return retryLater("TRANSPORT_ERROR: exists: " + err.Error()), nil + return retryLater("TRANSPORT_ERROR: getEntry: " + err.Error()), nil } - if !exists { + if uploadEntry == nil { return noopResolved("NOT_FOUND"), nil } - if err := s3a.rm(ctx, uploadsFolder, uploadID, true, true); err != nil { - if errors.Is(err, filer_pb.ErrNotFound) { - return noopResolved("NOT_FOUND_AT_DELETE"), nil - } - glog.V(1).Infof("lifecycle abort_mpu %s/%s: %v", req.Bucket, req.ObjectPath, err) - return retryLater("TRANSPORT_ERROR: rm: " + err.Error()), nil + object := string(uploadEntry.Extended[s3_constants.ExtMultipartObjectKey]) + code := s3a.withObjectWriteLock(req.Bucket, object, nil, func() s3err.ErrorCode { + return s3a.removeUploadDir(req.Bucket, uploadID, object) + }) + if code == s3err.ErrNone { + return done(), nil } - return done(), nil + return retryLater("TRANSPORT_ERROR: removeUploadDir: " + s3err.GetAPIError(code).Code), nil } // checkSoleSurvivorMarker returns nil to proceed with the delete, or a