mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-27 19:37:00 +00:00
s3: route non-conditional suspended DELETE off the DLM
A suspended DELETE removes the null version (the main object entry) and adds a delete marker. Now that suspended PUT routes on the object key, a DELETE left on the distributed lock wouldn't serialize against it (the per-path entry lock and the distributed lock are independent). Route it on the object key too. FinalizeVersionedWrite gains an optional delete_path it removes under the same lock before flipping the pointer, so "delete null + add marker" is one atomic step on the owner. createDeleteMarker threads it through (empty for the versioning-enabled marker path). The DeleteObjectHandler routes a no-versionId suspended delete when there is no If-Match — that condition targets the main object, not the .versions pointer the op evaluates, so a conditional suspended delete stays on the lock path. Object-lock can't apply (it requires versioning enabled, never suspended).
This commit is contained in:
@@ -292,6 +292,7 @@ message UpdateEntryResponse {
|
||||
// from set_extended[prior_latest_key].
|
||||
message FinalizeVersionedWriteRequest {
|
||||
string lock_key = 13; // object path; the per-path lock all of this object's writes share
|
||||
string delete_path = 14; // optional entry to delete under the lock first (suspended delete: the null version)
|
||||
string versions_dir = 1; // full path of <object>/.versions
|
||||
map<string, bytes> set_extended = 3; // merge into the .versions directory entry
|
||||
repeated string delete_extended = 4; // remove from the .versions directory entry
|
||||
|
||||
@@ -292,6 +292,7 @@ message UpdateEntryResponse {
|
||||
// from set_extended[prior_latest_key].
|
||||
message FinalizeVersionedWriteRequest {
|
||||
string lock_key = 13; // object path; the per-path lock all of this object's writes share
|
||||
string delete_path = 14; // optional entry to delete under the lock first (suspended delete: the null version)
|
||||
string versions_dir = 1; // full path of <object>/.versions
|
||||
map<string, bytes> set_extended = 3; // merge into the .versions directory entry
|
||||
repeated string delete_extended = 4; // remove from the .versions directory entry
|
||||
|
||||
@@ -1558,6 +1558,7 @@ func (x *UpdateEntryResponse) GetMetadataEvent() *SubscribeMetadataResponse {
|
||||
type FinalizeVersionedWriteRequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
LockKey string `protobuf:"bytes,13,opt,name=lock_key,json=lockKey,proto3" json:"lock_key,omitempty"` // object path; the per-path lock all of this object's writes share
|
||||
DeletePath string `protobuf:"bytes,14,opt,name=delete_path,json=deletePath,proto3" json:"delete_path,omitempty"` // optional entry to delete under the lock first (suspended delete: the null version)
|
||||
VersionsDir string `protobuf:"bytes,1,opt,name=versions_dir,json=versionsDir,proto3" json:"versions_dir,omitempty"` // full path of <object>/.versions
|
||||
SetExtended map[string][]byte `protobuf:"bytes,3,rep,name=set_extended,json=setExtended,proto3" json:"set_extended,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` // merge into the .versions directory entry
|
||||
DeleteExtended []string `protobuf:"bytes,4,rep,name=delete_extended,json=deleteExtended,proto3" json:"delete_extended,omitempty"` // remove from the .versions directory entry
|
||||
@@ -1613,6 +1614,13 @@ func (x *FinalizeVersionedWriteRequest) GetLockKey() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *FinalizeVersionedWriteRequest) GetDeletePath() string {
|
||||
if x != nil {
|
||||
return x.DeletePath
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *FinalizeVersionedWriteRequest) GetVersionsDir() string {
|
||||
if x != nil {
|
||||
return x.VersionsDir
|
||||
@@ -5983,9 +5991,11 @@ const file_filer_proto_rawDesc = "" +
|
||||
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
|
||||
"\x05value\x18\x02 \x01(\fR\x05value:\x028\x01\"a\n" +
|
||||
"\x13UpdateEntryResponse\x12J\n" +
|
||||
"\x0emetadata_event\x18\x01 \x01(\v2#.filer_pb.SubscribeMetadataResponseR\rmetadataEvent\"\x9b\x05\n" +
|
||||
"\x0emetadata_event\x18\x01 \x01(\v2#.filer_pb.SubscribeMetadataResponseR\rmetadataEvent\"\xbc\x05\n" +
|
||||
"\x1dFinalizeVersionedWriteRequest\x12\x19\n" +
|
||||
"\block_key\x18\r \x01(\tR\alockKey\x12!\n" +
|
||||
"\block_key\x18\r \x01(\tR\alockKey\x12\x1f\n" +
|
||||
"\vdelete_path\x18\x0e \x01(\tR\n" +
|
||||
"deletePath\x12!\n" +
|
||||
"\fversions_dir\x18\x01 \x01(\tR\vversionsDir\x12[\n" +
|
||||
"\fset_extended\x18\x03 \x03(\v28.filer_pb.FinalizeVersionedWriteRequest.SetExtendedEntryR\vsetExtended\x12'\n" +
|
||||
"\x0fdelete_extended\x18\x04 \x03(\tR\x0edeleteExtended\x12(\n" +
|
||||
|
||||
@@ -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, "", nil); 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, "", nil); err != nil {
|
||||
if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath, "", nil, ""); err != nil {
|
||||
return retryLater("TRANSPORT_ERROR: createDeleteMarker: " + err.Error()), nil
|
||||
}
|
||||
return done(), nil
|
||||
|
||||
@@ -532,7 +532,7 @@ func (s3a *S3ApiServer) finalizeCopyDestination(dstBucket, dstObject, dstVersion
|
||||
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 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)
|
||||
}
|
||||
|
||||
@@ -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, "", nil)
|
||||
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, "", nil)
|
||||
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
|
||||
@@ -249,7 +249,7 @@ func (s3a *S3ApiServer) DeleteObjectHandler(w http.ResponseWriter, r *http.Reque
|
||||
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)
|
||||
vid, err := s3a.createDeleteMarker(bucket, object, owner, cond, "")
|
||||
switch {
|
||||
case err == nil:
|
||||
deleteResult = deleteMutationResult{versionId: vid, deleteMarker: true}
|
||||
@@ -262,6 +262,27 @@ func (s3a *S3ApiServer) DeleteObjectHandler(w http.ResponseWriter, r *http.Reque
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Fast path: a suspended-versioning delete (no specific version) removes the
|
||||
// null version (the main object entry) and adds a delete marker — the op does
|
||||
// both atomically under the object-key lock, so it serializes with the (also
|
||||
// routed) suspended PUT. Routed only without If-Match: that condition targets
|
||||
// the main object, not the .versions pointer the op evaluates, so a
|
||||
// conditional suspended delete stays on the lock path.
|
||||
if !deleteHandled && versionId == "" && versioningState == s3_constants.VersioningSuspended {
|
||||
if cond, condOk := buildDeleteCondition(r); condOk && cond.Kind == filer_pb.WriteCondition_NONE {
|
||||
if owner := s3a.objectWriteOwner(bucket, s3_constants.NormalizeObjectKey(object)); owner != "" {
|
||||
nullPath := s3a.toFilerPath(bucket, s3_constants.NormalizeObjectKey(object))
|
||||
vid, err := s3a.createDeleteMarker(bucket, object, owner, cond, nullPath)
|
||||
if err == nil {
|
||||
deleteResult = deleteMutationResult{versionId: vid, deleteMarker: true}
|
||||
deleteCode, deleteHandled = s3err.ErrNone, true
|
||||
} else {
|
||||
glog.Warningf("DeleteObjectHandler: routed suspended delete 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)
|
||||
|
||||
@@ -98,10 +98,11 @@ func versionedPointerExtended(versionId, versionFileName string, versionEntry *f
|
||||
// version file and skips the object write lock and gateway precondition. A
|
||||
// transport/in-band error maps to InternalError (the client retries); there is
|
||||
// no silent non-atomic fallback.
|
||||
func (s3a *S3ApiServer) routedVersionedFinalize(owner pb.ServerAddress, bucket, object, versionId, versionFileName string, versionEntry *filer_pb.Entry, cond *filer_pb.WriteCondition) s3err.ErrorCode {
|
||||
func (s3a *S3ApiServer) routedVersionedFinalize(owner pb.ServerAddress, bucket, object, versionId, versionFileName string, versionEntry *filer_pb.Entry, cond *filer_pb.WriteCondition, deletePath string) s3err.ErrorCode {
|
||||
set, del := versionedPointerExtended(versionId, versionFileName, versionEntry)
|
||||
req := &filer_pb.FinalizeVersionedWriteRequest{
|
||||
LockKey: s3a.toFilerPath(bucket, object),
|
||||
DeletePath: deletePath,
|
||||
VersionsDir: s3a.toFilerPath(bucket, object+s3_constants.VersionsFolder),
|
||||
SetExtended: set,
|
||||
DeleteExtended: del,
|
||||
@@ -148,7 +149,7 @@ func (s3a *S3ApiServer) versionedAfterCreate(r *http.Request, bucket, object, ve
|
||||
cond, condOk := buildWriteCondition(r)
|
||||
if owner != "" && condOk {
|
||||
return func(versionEntry *filer_pb.Entry) s3err.ErrorCode {
|
||||
return s3a.routedVersionedFinalize(owner, bucket, object, versionId, versionFileName, versionEntry, cond)
|
||||
return s3a.routedVersionedFinalize(owner, bucket, object, versionId, versionFileName, versionEntry, cond, "")
|
||||
}, true
|
||||
}
|
||||
return func(versionEntry *filer_pb.Entry) s3err.ErrorCode {
|
||||
|
||||
@@ -198,7 +198,7 @@ type ObjectVersion struct {
|
||||
// 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) {
|
||||
func (s3a *S3ApiServer) createDeleteMarker(bucket, object string, routedOwner pb.ServerAddress, cond *filer_pb.WriteCondition, deleteNullPath string) (string, error) {
|
||||
// Clean up the object path first
|
||||
cleanObject := strings.TrimPrefix(object, "/")
|
||||
|
||||
@@ -252,7 +252,7 @@ func (s3a *S3ApiServer) createDeleteMarker(bucket, object string, routedOwner pb
|
||||
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 code := s3a.routedVersionedFinalize(routedOwner, bucket, cleanObject, versionId, versionFileName, deleteMarkerEntry, cond, deleteNullPath); code != s3err.ErrNone {
|
||||
if rbErr := s3a.rmObject(versionsDir, versionFileName, true, false); rbErr != nil {
|
||||
glog.Errorf("createDeleteMarker: rollback %s/%s: %v", versionsDir, versionFileName, rbErr)
|
||||
}
|
||||
|
||||
@@ -352,6 +352,16 @@ func (fs *FilerServer) FinalizeVersionedWrite(ctx context.Context, req *filer_pb
|
||||
}
|
||||
}
|
||||
|
||||
// Optional: delete an entry under the same lock first (suspended delete
|
||||
// removes the "null" version at the main object path before adding the
|
||||
// marker). Not-found is fine (idempotent).
|
||||
if req.DeletePath != "" {
|
||||
if derr := fs.filer.DeleteEntryMetaAndData(ctx, util.FullPath(req.DeletePath), false, false, true, req.IsFromOtherCluster, req.Signatures, 0); derr != nil && derr != filer_pb.ErrNotFound {
|
||||
resp.Error = derr.Error()
|
||||
return resp, nil
|
||||
}
|
||||
}
|
||||
|
||||
// 2. Stamp the previously-latest version as noncurrent BEFORE the pointer
|
||||
// flip (so the lifecycle router observes it). Best-effort.
|
||||
newFileName := string(req.SetExtended[req.PriorLatestKey])
|
||||
|
||||
Reference in New Issue
Block a user