mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-21 22:56:55 +00:00
* fix(s3/lifecycle): reader handles bare /buckets parent and pre-normalizes prefix extractBucketKey accepted /buckets/ but rejected /buckets (no trailing slash); some delete events emit the bare form, so bucket-root events were silently dropped. Pre-normalize BucketsPath once on Run instead of recomputing per event. * perf(s3/lifecycle): pool sha256 hashers in ShardID ShardID runs on every meta-log event before the shard filter; a fresh sha256.New per call produces measurable allocator pressure under load. sync.Pool reuses hashers across calls. * fix(s3/lifecycle): router skips hard deletes and missing-attribute events A hard delete carries no schedule-relevant state — Expiration would hit NOOP_RESOLVED at dispatch and ExpiredObjectDeleteMarker fires from a Create on the latest version. Skip rather than burn a schedule slot. Missing Attributes leaves ModTime at year 0001, which makes ExpirationDays fire immediately at dispatch. Skip the event instead. Drop the unused 'versioned' parameter from buildObjectInfo; the dispatcher's identity-CAS handles version drift in Phase 5. * fix(s3/lifecycle): EntryIdentity.MtimeNs holds true nanoseconds Both computeEntryIdentity (server) and buildIdentity (router) wrote entry.Attributes.Mtime (seconds) into a field named MtimeNs. The CAS worked because both sides agreed, but the encoding contradicted the field name and would break if either side later started using true nanoseconds. Combine Mtime*1e9 + the FuseAttributes.MtimeNs nanosecond component on both sides; the test was updated to match. * fix(s3/lifecycle): dispatcher distinguishes ctx cancel from transport errors A canceled or deadline-exceeded RPC is shutdown, not a transport failure: re-queue the Match at its original DueTime with no retry-budget burn so a quick restart can't escalate it to BLOCKED. * fix(s3/lifecycle): reader fallback prefix normalization mirrors Run The fallback path that builds prefix from r.BucketsPath when bucketsPathSlash is empty (test-only entry into extractBucketKey) was appending an unconditional '/', producing '//' if BucketsPath already ended with one. Use the same normalization Run does. * fix(s3/lifecycle): ObjectInfo.ModTime carries the nanosecond component ModTime dropped FuseAttributes.MtimeNs, leaving ExpirationDays one nanosecond off relative to EntryIdentity.MtimeNs. Pass both to time.Unix so the precision matches the CAS witness.
204 lines
7.9 KiB
Go
204 lines
7.9 KiB
Go
package s3api
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"errors"
|
|
"sort"
|
|
"strconv"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"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"
|
|
)
|
|
|
|
// LifecycleDelete executes one (rule, action) verdict: re-fetch, identity
|
|
// CAS, object-lock check, dispatch by kind. Errors surface as outcomes;
|
|
// reader cursors and pending state are the worker's concern.
|
|
func (s3a *S3ApiServer) LifecycleDelete(ctx context.Context, req *s3_lifecycle_pb.LifecycleDeleteRequest) (*s3_lifecycle_pb.LifecycleDeleteResponse, error) {
|
|
if req == nil || req.Bucket == "" || req.ObjectPath == "" {
|
|
return blocked("FATAL_EVENT_ERROR: empty bucket or object_path"), nil
|
|
}
|
|
|
|
// MPU init lives at .uploads/<id>/; not handled by getObjectEntry.
|
|
if req.ActionKind == s3_lifecycle_pb.ActionKind_ABORT_MPU {
|
|
return s3a.lifecycleAbortMPU(ctx, req)
|
|
}
|
|
|
|
entry, err := s3a.getObjectEntry(req.Bucket, req.ObjectPath, req.VersionId)
|
|
if err != nil {
|
|
if errors.Is(err, filer_pb.ErrNotFound) || errors.Is(err, ErrObjectNotFound) || errors.Is(err, ErrVersionNotFound) || errors.Is(err, ErrLatestVersionNotFound) {
|
|
return noopResolved("NOT_FOUND"), nil
|
|
}
|
|
glog.V(1).Infof("lifecycle: live fetch %s/%s@%s: %v", req.Bucket, req.ObjectPath, req.VersionId, err)
|
|
return retryLater("TRANSPORT_ERROR: " + err.Error()), nil
|
|
}
|
|
|
|
live := computeEntryIdentity(entry)
|
|
if !identityMatches(live, req.ExpectedIdentity) {
|
|
return noopResolved("STALE_IDENTITY"), nil
|
|
}
|
|
|
|
// Lifecycle never bypasses governance/compliance; the http.Request is
|
|
// only read when bypass is allowed, so nil is safe here.
|
|
if err := s3a.enforceObjectLockProtections(nil, req.Bucket, req.ObjectPath, req.VersionId, false); err != nil {
|
|
glog.V(2).Infof("lifecycle: SKIPPED_OBJECT_LOCK %s/%s@%s: %v", req.Bucket, req.ObjectPath, req.VersionId, err)
|
|
return &s3_lifecycle_pb.LifecycleDeleteResponse{
|
|
Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_SKIPPED_OBJECT_LOCK,
|
|
Reason: err.Error(),
|
|
}, nil
|
|
}
|
|
|
|
return s3a.lifecycleDispatch(ctx, req, entry)
|
|
}
|
|
|
|
func (s3a *S3ApiServer) lifecycleDispatch(ctx context.Context, req *s3_lifecycle_pb.LifecycleDeleteRequest, entry *filer_pb.Entry) (*s3_lifecycle_pb.LifecycleDeleteResponse, error) {
|
|
switch req.ActionKind {
|
|
case s3_lifecycle_pb.ActionKind_EXPIRATION_DAYS, s3_lifecycle_pb.ActionKind_EXPIRATION_DATE:
|
|
// Current-version expiration: Enabled -> delete marker; Suspended
|
|
// -> delete null + new marker; Off -> remove. Filer errors classify
|
|
// as RETRY_LATER; the worker's budget promotes to BLOCKED.
|
|
state, vErr := s3a.getVersioningState(req.Bucket)
|
|
if vErr != nil {
|
|
if errors.Is(vErr, filer_pb.ErrNotFound) {
|
|
return noopResolved("BUCKET_NOT_FOUND"), nil
|
|
}
|
|
return retryLater("TRANSPORT_ERROR: versioning lookup: " + vErr.Error()), nil
|
|
}
|
|
switch state {
|
|
case s3_constants.VersioningEnabled:
|
|
if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath); err != nil {
|
|
return retryLater("TRANSPORT_ERROR: createDeleteMarker: " + err.Error()), nil
|
|
}
|
|
return done(), nil
|
|
case s3_constants.VersioningSuspended:
|
|
// Best-effort null delete; NotFound is benign.
|
|
if err := s3a.deleteSpecificObjectVersion(req.Bucket, req.ObjectPath, "null"); err != nil {
|
|
if !errors.Is(err, filer_pb.ErrNotFound) && !errors.Is(err, ErrVersionNotFound) {
|
|
return retryLater("TRANSPORT_ERROR: deleteNullVersion: " + err.Error()), nil
|
|
}
|
|
}
|
|
if _, err := s3a.createDeleteMarker(req.Bucket, req.ObjectPath); err != nil {
|
|
return retryLater("TRANSPORT_ERROR: createDeleteMarker: " + err.Error()), nil
|
|
}
|
|
return done(), nil
|
|
default:
|
|
err := s3a.WithFilerClient(false, func(c filer_pb.SeaweedFilerClient) error {
|
|
return s3a.deleteUnversionedObjectWithClient(c, req.Bucket, req.ObjectPath)
|
|
})
|
|
if err != nil {
|
|
if errors.Is(err, filer_pb.ErrNotFound) || errors.Is(err, ErrObjectNotFound) {
|
|
return noopResolved("NOT_FOUND_AT_DELETE"), nil
|
|
}
|
|
return retryLater("TRANSPORT_ERROR: deleteUnversioned: " + err.Error()), nil
|
|
}
|
|
return done(), nil
|
|
}
|
|
|
|
case s3_lifecycle_pb.ActionKind_NONCURRENT_DAYS,
|
|
s3_lifecycle_pb.ActionKind_NEWER_NONCURRENT,
|
|
s3_lifecycle_pb.ActionKind_EXPIRED_DELETE_MARKER:
|
|
// EXPIRED_DELETE_MARKER targets the marker version itself.
|
|
if req.VersionId == "" {
|
|
return blocked("FATAL_EVENT_ERROR: version_id required for noncurrent / delete-marker delete"), nil
|
|
}
|
|
if err := s3a.deleteSpecificObjectVersion(req.Bucket, req.ObjectPath, req.VersionId); err != nil {
|
|
if errors.Is(err, filer_pb.ErrNotFound) || errors.Is(err, ErrVersionNotFound) || errors.Is(err, ErrObjectNotFound) {
|
|
return noopResolved("NOT_FOUND_AT_DELETE"), nil
|
|
}
|
|
return retryLater("TRANSPORT_ERROR: deleteSpecificVersion: " + err.Error()), nil
|
|
}
|
|
return done(), nil
|
|
|
|
case s3_lifecycle_pb.ActionKind_ABORT_MPU:
|
|
return blocked("FATAL_EVENT_ERROR: ABORT_MPU dispatched after fetch"), nil
|
|
|
|
default:
|
|
return blocked("FATAL_EVENT_ERROR: unknown action_kind " + req.ActionKind.String()), nil
|
|
}
|
|
}
|
|
|
|
func (s3a *S3ApiServer) lifecycleAbortMPU(ctx context.Context, req *s3_lifecycle_pb.LifecycleDeleteRequest) (*s3_lifecycle_pb.LifecycleDeleteResponse, error) {
|
|
// TODO(phase-5): plumb to abortMultipartUpload (currently expects an
|
|
// *s3.AbortMultipartUploadInput; lifecycle has no HTTP request).
|
|
return retryLater("ABORT_MPU not yet wired"), nil
|
|
}
|
|
|
|
// computeEntryIdentity captures (mtime, size, head fid, sorted-Extended hash):
|
|
// an overwrite changes mtime/size/fid; a metadata edit changes Extended; a
|
|
// snapshot-restore that preserves mtime+size still differs in head_fid.
|
|
func computeEntryIdentity(entry *filer_pb.Entry) *s3_lifecycle_pb.EntryIdentity {
|
|
if entry == nil {
|
|
return nil
|
|
}
|
|
id := &s3_lifecycle_pb.EntryIdentity{}
|
|
if entry.Attributes != nil {
|
|
// FuseAttributes splits the timestamp across Mtime (seconds) and
|
|
// MtimeNs (nanosecond component); EntryIdentity.MtimeNs is the
|
|
// combined nanoseconds-since-epoch value.
|
|
id.MtimeNs = entry.Attributes.Mtime*int64(1e9) + int64(entry.Attributes.MtimeNs)
|
|
id.Size = int64(entry.Attributes.FileSize)
|
|
}
|
|
if len(entry.GetChunks()) > 0 {
|
|
id.HeadFid = entry.GetChunks()[0].FileId
|
|
}
|
|
id.ExtendedHash = hashExtended(entry.Extended)
|
|
return id
|
|
}
|
|
|
|
// hashExtended is length-prefixed; a forged tag value can't collide with a
|
|
// legitimate multi-tag map.
|
|
func hashExtended(ext map[string][]byte) []byte {
|
|
if len(ext) == 0 {
|
|
return nil
|
|
}
|
|
keys := make([]string, 0, len(ext))
|
|
for k := range ext {
|
|
keys = append(keys, k)
|
|
}
|
|
sort.Strings(keys)
|
|
h := sha256.New()
|
|
for _, k := range keys {
|
|
h.Write([]byte(strconv.Itoa(len(k))))
|
|
h.Write([]byte{':'})
|
|
h.Write([]byte(k))
|
|
v := ext[k]
|
|
h.Write([]byte(strconv.Itoa(len(v))))
|
|
h.Write([]byte{':'})
|
|
h.Write(v)
|
|
}
|
|
return h.Sum(nil)
|
|
}
|
|
|
|
func identityMatches(live, want *s3_lifecycle_pb.EntryIdentity) bool {
|
|
if want == nil {
|
|
// No CAS witness (early bootstrap); skip.
|
|
return true
|
|
}
|
|
if live == nil {
|
|
return false
|
|
}
|
|
if live.MtimeNs != want.MtimeNs || live.Size != want.Size {
|
|
return false
|
|
}
|
|
if live.HeadFid != want.HeadFid {
|
|
return false
|
|
}
|
|
return bytes.Equal(live.ExtendedHash, want.ExtendedHash)
|
|
}
|
|
|
|
func done() *s3_lifecycle_pb.LifecycleDeleteResponse {
|
|
return &s3_lifecycle_pb.LifecycleDeleteResponse{Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_DONE}
|
|
}
|
|
func noopResolved(reason string) *s3_lifecycle_pb.LifecycleDeleteResponse {
|
|
return &s3_lifecycle_pb.LifecycleDeleteResponse{Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_NOOP_RESOLVED, Reason: reason}
|
|
}
|
|
func blocked(reason string) *s3_lifecycle_pb.LifecycleDeleteResponse {
|
|
return &s3_lifecycle_pb.LifecycleDeleteResponse{Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_BLOCKED, Reason: reason}
|
|
}
|
|
func retryLater(reason string) *s3_lifecycle_pb.LifecycleDeleteResponse {
|
|
return &s3_lifecycle_pb.LifecycleDeleteResponse{Outcome: s3_lifecycle_pb.LifecycleDeleteOutcome_RETRY_LATER, Reason: reason}
|
|
}
|