mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-22 07:06:51 +00:00
s3: route versioned COPY and delete-marker off the DLM
Both create a unique version/marker file plus an atomic latest-pointer flip — the same shape as versioned PutObject — so both now reuse FinalizeVersionedWrite instead of the distributed lock. finalizeCopyDestination (VersioningEnabled) and createDeleteMarker take an optional owner + condition; when the .versions owner is known they run the demote + pointer flip (and the precondition) atomically on that owner via routedVersionedFinalize, rolling back the version/marker file on failure and surfacing a precondition miss as 412. CopyObjectHandler and DeleteObjectHandler take the routed path for versioning-enabled destinations and skip the lock; suspended versioning, version-specific deletes, metadata-only self-copy, and the lifecycle marker path keep the existing lock-based path. No proto or filer changes — reuses the op from the versioned-PutObject PR.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user