diff --git a/weed/s3api/s3api_internal_lifecycle.go b/weed/s3api/s3api_internal_lifecycle.go index 164a999ba..1abdfc816 100644 --- a/weed/s3api/s3api_internal_lifecycle.go +++ b/weed/s3api/s3api_internal_lifecycle.go @@ -71,7 +71,7 @@ func (s3a *S3ApiServer) lifecycleDispatch(ctx context.Context, req *s3_lifecycle } switch state { case s3_constants.VersioningEnabled: - if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath); err != nil { + if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath, "", nil); err != nil { return retryLater("TRANSPORT_ERROR: createDeleteMarker: " + err.Error()), nil } return done(), nil @@ -82,7 +82,7 @@ func (s3a *S3ApiServer) lifecycleDispatch(ctx context.Context, req *s3_lifecycle return retryLater("TRANSPORT_ERROR: deleteNullVersion: " + err.Error()), nil } } - if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath); err != nil { + if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath, "", nil); err != nil { return retryLater("TRANSPORT_ERROR: createDeleteMarker: " + err.Error()), nil } return done(), nil diff --git a/weed/s3api/s3api_object_handlers_copy.go b/weed/s3api/s3api_object_handlers_copy.go index e6a083265..9e4628a9f 100644 --- a/weed/s3api/s3api_object_handlers_copy.go +++ b/weed/s3api/s3api_object_handlers_copy.go @@ -183,7 +183,7 @@ func (s3a *S3ApiServer) CopyObjectHandler(w http.ResponseWriter, r *http.Request } updatedEntry.Attributes.Mtime = t.Unix() - dstVersionId, etag, currentErr = s3a.finalizeCopyDestination(dstBucket, dstObject, dstVersioningState, updatedEntry) + dstVersionId, etag, currentErr = s3a.finalizeCopyDestination(dstBucket, dstObject, dstVersioningState, updatedEntry, "", nil) if currentErr != nil { return filerErrorToS3Error(currentErr) } @@ -338,20 +338,35 @@ func (s3a *S3ApiServer) CopyObjectHandler(w http.ResponseWriter, r *http.Request var dstVersionId string var etag string - // Fast path: for a non-versioned destination, the finalize is a single - // CreateEntry, so route it to the destination key's owner filer with the - // precondition and let it serialize the write locally, skipping the - // distributed lock. Versioned destinations (multi-entry: version file + - // latest pointer) keep the lock path below. + // Fast path: route the finalize to the destination's owner filer and skip the + // distributed lock. A non-versioned destination is a single CreateEntry; a + // versioning-enabled destination is a unique version file plus an atomic + // latest-pointer update on the .versions owner (FinalizeVersionedWrite). + // Suspended versioning (multi-version IsLatest rewrite) keeps the lock path. var finalizeCode s3err.ErrorCode routed := false - if dstVersioningState == "" { - if cond, condOk := buildWriteCondition(r); condOk { + if cond, condOk := buildWriteCondition(r); condOk { + switch dstVersioningState { + case "": if owner, ownerOk := s3a.routedObjectOwner(dstBucket, dstObject); ownerOk { if code, e, handled := s3a.routedCopyFinalize(owner, dstBucket, dstObject, dstEntry, cond); handled { finalizeCode, etag, routed = code, e, true } } + case s3_constants.VersioningEnabled: + if owner := s3a.objectWriteOwner(dstBucket, s3_constants.NormalizeObjectKey(dstObject)); owner != "" { + var finalizeErr error + dstVersionId, etag, finalizeErr = s3a.finalizeCopyDestination(dstBucket, dstObject, dstVersioningState, dstEntry, owner, cond) + switch { + case finalizeErr == nil: + finalizeCode = s3err.ErrNone + case errors.Is(finalizeErr, errVersionedPreconditionFailed): + finalizeCode = s3err.ErrPreconditionFailed + default: + finalizeCode = filerErrorToS3Error(finalizeErr) + } + routed = true + } } } if !routed { @@ -359,7 +374,7 @@ func (s3a *S3ApiServer) CopyObjectHandler(w http.ResponseWriter, r *http.Request return s3a.checkConditionalHeaders(r, dstBucket, dstObject) }, func() s3err.ErrorCode { var finalizeErr error - dstVersionId, etag, finalizeErr = s3a.finalizeCopyDestination(dstBucket, dstObject, dstVersioningState, dstEntry) + dstVersionId, etag, finalizeErr = s3a.finalizeCopyDestination(dstBucket, dstObject, dstVersioningState, dstEntry, "", nil) if finalizeErr != nil { return filerErrorToS3Error(finalizeErr) } @@ -478,7 +493,11 @@ func (s3a *S3ApiServer) routedCopyFinalize(owner pb.ServerAddress, dstBucket, ds } } -func (s3a *S3ApiServer) finalizeCopyDestination(dstBucket, dstObject, dstVersioningState string, dstEntry *filer_pb.Entry) (versionId string, etag string, err error) { +// finalizeCopyDestination writes the copy destination. When routedOwner is set +// (versioning-enabled destination routed off the distributed lock), the latest- +// pointer update goes through FinalizeVersionedWrite on that owner, which also +// evaluates cond; otherwise it uses the lock-based updateLatestVersionInDirectory. +func (s3a *S3ApiServer) finalizeCopyDestination(dstBucket, dstObject, dstVersioningState string, dstEntry *filer_pb.Entry, routedOwner pb.ServerAddress, cond *filer_pb.WriteCondition) (versionId string, etag string, err error) { normalizedObject := s3_constants.NormalizeObjectKey(dstObject) dstPath := util.FullPath(fmt.Sprintf("%s/%s", s3a.bucketDir(dstBucket), normalizedObject)) dstDir, dstName := dstPath.DirAndName() @@ -510,6 +529,18 @@ func (s3a *S3ApiServer) finalizeCopyDestination(dstBucket, dstObject, dstVersion return "", "", err } + if routedOwner != "" { + // Routed: precondition + demote + pointer flip run atomically on the + // .versions owner. Roll back the version file if it does not succeed. + if code := s3a.routedVersionedFinalize(routedOwner, dstBucket, normalizedObject, versionId, versionFileName, dstEntry, cond); code != s3err.ErrNone { + if rollbackErr := s3a.rollbackCopyVersion(bucketDir, versionObjectPath); rollbackErr != nil { + glog.Errorf("CopyObjectHandler: failed to rollback version %s for %s/%s after routed finalize error: %v", versionId, dstBucket, normalizedObject, rollbackErr) + } + return "", "", finalizeCodeToError(code) + } + return versionId, etag, nil + } + if err = s3a.updateLatestVersionInDirectory(dstBucket, normalizedObject, versionId, versionFileName, dstEntry); err != nil { if rollbackErr := s3a.rollbackCopyVersion(bucketDir, versionObjectPath); rollbackErr != nil { glog.Errorf("CopyObjectHandler: failed to rollback version %s for %s/%s after latest pointer update error: %v", versionId, dstBucket, normalizedObject, rollbackErr) diff --git a/weed/s3api/s3api_object_handlers_delete.go b/weed/s3api/s3api_object_handlers_delete.go index e95e63751..c4387b652 100644 --- a/weed/s3api/s3api_object_handlers_delete.go +++ b/weed/s3api/s3api_object_handlers_delete.go @@ -133,7 +133,7 @@ func (s3a *S3ApiServer) deleteVersionedObject(r *http.Request, bucket, object, v return result, s3err.ErrNone case versioningState == s3_constants.VersioningEnabled: - deleteMarkerVersionId, err := s3a.createDeleteMarker(bucket, object) + deleteMarkerVersionId, err := s3a.createDeleteMarker(bucket, object, "", nil) if err != nil { glog.Errorf("deleteVersionedObject: failed to create delete marker for %s/%s: %v", bucket, object, err) return result, s3err.ErrInternalError @@ -152,7 +152,7 @@ func (s3a *S3ApiServer) deleteVersionedObject(r *http.Request, bucket, object, v glog.Errorf("deleteVersionedObject: failed to delete null version for %s/%s: %v", bucket, object, err) return result, s3err.ErrInternalError } - deleteMarkerVersionId, err := s3a.createDeleteMarker(bucket, object) + deleteMarkerVersionId, err := s3a.createDeleteMarker(bucket, object, "", nil) if err != nil { glog.Errorf("deleteVersionedObject: failed to create delete marker for suspended versioning %s/%s: %v", bucket, object, err) return result, s3err.ErrInternalError @@ -241,6 +241,27 @@ func (s3a *S3ApiServer) DeleteObjectHandler(w http.ResponseWriter, r *http.Reque } } } + + // Fast path: a versioning-enabled delete (no specific version) creates a + // delete marker — a unique marker file plus an atomic latest-pointer flip. We + // route it to the .versions owner, off the distributed lock. Suspended + // versioning (multi-step) and version-specific deletes keep the lock path. + if !deleteHandled && versionId == "" && versioningState == s3_constants.VersioningEnabled { + if cond, condOk := buildDeleteCondition(r); condOk { + if owner := s3a.objectWriteOwner(bucket, s3_constants.NormalizeObjectKey(object)); owner != "" { + vid, err := s3a.createDeleteMarker(bucket, object, owner, cond) + switch { + case err == nil: + deleteResult = deleteMutationResult{versionId: vid, deleteMarker: true} + deleteCode, deleteHandled = s3err.ErrNone, true + case errors.Is(err, errVersionedPreconditionFailed): + deleteCode, deleteHandled = s3err.ErrPreconditionFailed, true + default: + glog.Warningf("DeleteObjectHandler: routed delete marker for %s/%s failed, falling back to lock: %v", bucket, object, err) + } + } + } + } if !deleteHandled { deleteCode = s3a.withObjectWriteLock(bucket, object, func() s3err.ErrorCode { return s3a.checkDeleteIfMatch(bucket, object, versionId, versioningState, r.Header.Get(s3_constants.IfMatch), s3err.ErrPreconditionFailed) diff --git a/weed/s3api/s3api_object_versioned_finalize.go b/weed/s3api/s3api_object_versioned_finalize.go index f4045a909..110b72991 100644 --- a/weed/s3api/s3api_object_versioned_finalize.go +++ b/weed/s3api/s3api_object_versioned_finalize.go @@ -2,6 +2,8 @@ package s3api import ( "context" + "errors" + "fmt" "net/http" "strconv" "time" @@ -13,6 +15,24 @@ import ( "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" ) +// errVersionedPreconditionFailed is returned by versioned finalize helpers whose +// contract is (…, error) when the routed FinalizeVersionedWrite reports a failed +// precondition, so callers can surface it as 412 rather than a generic 500. +var errVersionedPreconditionFailed = errors.New("versioned write precondition failed") + +// finalizeCodeToError adapts a routedVersionedFinalize result to an error: nil on +// success, the precondition sentinel on 412, else a generic error. +func finalizeCodeToError(code s3err.ErrorCode) error { + switch code { + case s3err.ErrNone: + return nil + case s3err.ErrPreconditionFailed: + return errVersionedPreconditionFailed + default: + return fmt.Errorf("versioned finalize failed: code %v", code) + } +} + // objectWriteOwner returns the filer that owns all of an object's writes, or "" // when no ring view is available. It hashes the object key with the same prefix // the non-versioned path uses (routedObjectOwner), so normal, suspended, and diff --git a/weed/s3api/s3api_object_versioning.go b/weed/s3api/s3api_object_versioning.go index 98f7da363..ff908e68f 100644 --- a/weed/s3api/s3api_object_versioning.go +++ b/weed/s3api/s3api_object_versioning.go @@ -19,6 +19,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" s3_constants "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" @@ -191,7 +192,13 @@ type ObjectVersion struct { } // createDeleteMarker creates a delete marker for versioned delete operations -func (s3a *S3ApiServer) createDeleteMarker(bucket, object string) (string, error) { +// createDeleteMarker writes a delete-marker version and makes it the latest. +// When routedOwner is set, the latest-pointer update (and the optional cond) run +// atomically via FinalizeVersionedWrite on that owner, off the distributed lock; +// otherwise it uses the lock-based updateLatestVersionInDirectory. A failed +// routed finalize rolls back the marker file and returns +// errVersionedPreconditionFailed for a precondition miss. +func (s3a *S3ApiServer) createDeleteMarker(bucket, object string, routedOwner pb.ServerAddress, cond *filer_pb.WriteCondition) (string, error) { // Clean up the object path first cleanObject := strings.TrimPrefix(object, "/") @@ -242,6 +249,19 @@ func (s3a *S3ApiServer) createDeleteMarker(bucket, object string) (string, error }, Extended: deleteMarkerExtended, } + if routedOwner != "" { + // Routed: precondition + demote + pointer flip run atomically on the + // .versions owner. Roll back the marker file if it does not succeed. + if code := s3a.routedVersionedFinalize(routedOwner, bucket, cleanObject, versionId, versionFileName, deleteMarkerEntry, cond); code != s3err.ErrNone { + if rbErr := s3a.rmObject(versionsDir, versionFileName, true, false); rbErr != nil { + glog.Errorf("createDeleteMarker: rollback %s/%s: %v", versionsDir, versionFileName, rbErr) + } + return "", finalizeCodeToError(code) + } + glog.V(2).Infof("createDeleteMarker: routed delete marker %s for %s/%s", versionId, bucket, object) + return versionId, nil + } + err = s3a.updateLatestVersionInDirectory(bucket, cleanObject, versionId, versionFileName, deleteMarkerEntry) if err != nil { glog.Errorf("createDeleteMarker: failed to update latest version in directory: %v", err)