s3: abort completed multipart uploads metadata-only (#11385)

* s3: abort a completed upload's leftover directory metadata-only

A .uploads/<id> 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 <key>.versions carries the
upload id, remove .uploads/<id> 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/<id>. 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.
This commit is contained in:
Chris Lu
2026-09-18 01:01:52 -07:00
committed by GitHub
parent 0ca1c19821
commit 87ee3b63a2
3 changed files with 396 additions and 20 deletions
+124 -10
View File
@@ -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 <key>.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 {
+261
View File
@@ -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/<id> 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/<id> 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)
}
}
+11 -10
View File
@@ -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