diff --git a/other/java/client/src/main/proto/filer.proto b/other/java/client/src/main/proto/filer.proto index 4979819e8..c33fb7770 100644 --- a/other/java/client/src/main/proto/filer.proto +++ b/other/java/client/src/main/proto/filer.proto @@ -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 /.versions map set_extended = 3; // merge into the .versions directory entry repeated string delete_extended = 4; // remove from the .versions directory entry diff --git a/weed/pb/filer.proto b/weed/pb/filer.proto index 4979819e8..c33fb7770 100644 --- a/weed/pb/filer.proto +++ b/weed/pb/filer.proto @@ -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 /.versions map set_extended = 3; // merge into the .versions directory entry repeated string delete_extended = 4; // remove from the .versions directory entry diff --git a/weed/pb/filer_pb/filer.pb.go b/weed/pb/filer_pb/filer.pb.go index 9ad4ebea0..f352b6888 100644 --- a/weed/pb/filer_pb/filer.pb.go +++ b/weed/pb/filer_pb/filer.pb.go @@ -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 /.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" + diff --git a/weed/s3api/s3api_internal_lifecycle.go b/weed/s3api/s3api_internal_lifecycle.go index 1abdfc816..a636d4aeb 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, "", 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 diff --git a/weed/s3api/s3api_object_handlers_copy.go b/weed/s3api/s3api_object_handlers_copy.go index 9e4628a9f..9c5635fdd 100644 --- a/weed/s3api/s3api_object_handlers_copy.go +++ b/weed/s3api/s3api_object_handlers_copy.go @@ -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) } diff --git a/weed/s3api/s3api_object_handlers_delete.go b/weed/s3api/s3api_object_handlers_delete.go index c4387b652..d965b4bcd 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, "", 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) diff --git a/weed/s3api/s3api_object_versioned_finalize.go b/weed/s3api/s3api_object_versioned_finalize.go index 1dd0daa1d..4ff38ff47 100644 --- a/weed/s3api/s3api_object_versioned_finalize.go +++ b/weed/s3api/s3api_object_versioned_finalize.go @@ -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 { diff --git a/weed/s3api/s3api_object_versioning.go b/weed/s3api/s3api_object_versioning.go index ff908e68f..639a5a4f2 100644 --- a/weed/s3api/s3api_object_versioning.go +++ b/weed/s3api/s3api_object_versioning.go @@ -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) } diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index 2dac0624e..3e524e86f 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -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])