diff --git a/cmd/erasure-healing.go b/cmd/erasure-healing.go index 8afc126ca..8fb371270 100644 --- a/cmd/erasure-healing.go +++ b/cmd/erasure-healing.go @@ -1068,6 +1068,11 @@ func (er erasureObjects) HealObject(ctx context.Context, bucket, object, version newReqInfo = logger.NewReqInfo("", "", globalDeploymentID(), "", "Heal", bucket, object) } healCtx := logger.SetReqInfo(GlobalContext, newReqInfo) + if opts.NoLock { + // The caller owns the namespace lock. Stop if that lock's context is + // canceled instead of continuing metadata writes after losing it. + healCtx = logger.SetReqInfo(ctx, newReqInfo) + } // Healing directories handle it separately. if HasSuffix(object, SlashSeparator) { diff --git a/cmd/erasure-multipart.go b/cmd/erasure-multipart.go index 9de330e28..67eab87e9 100644 --- a/cmd/erasure-multipart.go +++ b/cmd/erasure-multipart.go @@ -1155,17 +1155,10 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str return oi, toObjectErr(err, bucket, object, uploadID) } - // A trusted SSE-C replica completion re-orders the Object Lock carried in the - // upload metadata against the version it is about to replace, read on this - // erasure set under the write lock held above, so a hold or retention that - // reached the version after this upload was initiated is not rolled back at - // completion (issue #120). Scoped to SSE-C uploads, the only ones this issue - // routes through completion. - // - // Scope: correct for a single erasure set. A multi-pool deployment (duplicate - // versions across pools, ModTime ties, cross-pool lock authority) is out of - // scope and tracked in pgsty/silo#133. - if opts.ReplicaLockReconcile && crypto.SSEC.IsEncrypted(fi.Metadata) { + // Reconcile against the upload's persisted version, under the object lock. + // Multi-pool callers supply a resolver spanning all pools, including those + // draining their contents; the upload itself stays in its original pool. + if opts.ReplicaLockReconcile { // A persisted upload records the null version as an empty VersionID; look // it up as the null version so the reconcile reads the addressed version's // stored lock, not the latest version's. @@ -1173,7 +1166,11 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str if lookupVersionID == "" { lookupVersionID = nullVersionID } - curr, gerr := er.getObjectInfo(ctx, bucket, object, ObjectOptions{ + getObjectInfo := er.getObjectInfo + if opts.replicaObjectInfo != nil { + getObjectInfo = opts.replicaObjectInfo + } + curr, gerr := getObjectInfo(ctx, bucket, object, ObjectOptions{ VersionID: lookupVersionID, Versioned: opts.Versioned, VersionSuspended: opts.VersionSuspended, @@ -1182,6 +1179,7 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str switch { case gerr == nil: reconcileStoredObjectLock(fi.Metadata, storedObjectLockState(curr.UserDefined)) + reconcileStoredObjectTags(fi.Metadata, curr.UserDefined) case isErrVersionNotFound(gerr) || isErrObjectNotFound(gerr): // No existing version to order against: keep the upload's own accepted // lock, including a pre-upgrade upload that persisted values without diff --git a/cmd/erasure-object.go b/cmd/erasure-object.go index bade406e6..c8b270519 100644 --- a/cmd/erasure-object.go +++ b/cmd/erasure-object.go @@ -133,6 +133,11 @@ func (er erasureObjects) CopyObject(ctx context.Context, srcBucket, srcObject, d return fi.ToObjectInfo(srcBucket, srcObject, srcOpts.Versioned || srcOpts.VersionSuspended), toObjectErr(errMethodNotAllowed, srcBucket, srcObject) } + if dstOpts.ReplicaLockReconcile { + reconcileStoredObjectLock(srcInfo.UserDefined, storedObjectLockState(fi.Metadata)) + reconcileStoredObjectTags(srcInfo.UserDefined, fi.Metadata) + } + filterOnlineDisksInplace(fi, metaArr, onlineDisks) versionID := srcInfo.VersionID @@ -1280,7 +1285,11 @@ func (er erasureObjects) putObject(ctx context.Context, bucket string, object st opts.NoLock = true } - obj, err := er.getObjectInfo(ctx, bucket, object, opts) + getObjectInfo := er.getObjectInfo + if opts.replicaObjectInfo != nil { + getObjectInfo = opts.replicaObjectInfo + } + obj, err := getObjectInfo(ctx, bucket, object, opts) // A destination read that fails for a reason other than not-found must not // be taken as a passed precondition or as absent lock state. if err != nil && !isErrVersionNotFound(err) && !isErrObjectNotFound(err) { @@ -1297,17 +1306,12 @@ func (er erasureObjects) putObject(ctx context.Context, bucket string, object st } } - // Order this trusted SSE-C replica's Object Lock against the addressed - // version's stored state, read on this erasure set under the write lock, - // so a value that lost the ordering cannot overwrite a newer one committed - // after the handler decided (issue #120). Only reconcile against an - // existing version; on not-found the write's own accepted lock is kept. - // - // Scope: correct for a single erasure set. A multi-pool deployment - // (duplicate versions across pools, ModTime ties, cross-pool lock - // authority) is out of scope and tracked in pgsty/silo#133. + // Re-read under the write lock, using the pools-layer resolver when + // present. A missing version keeps the incoming accepted state; an + // existing version contributes independently ordered lock and tags. if opts.ReplicaLockReconcile && err == nil { reconcileStoredObjectLock(opts.UserDefined, storedObjectLockState(obj.UserDefined)) + reconcileStoredObjectTags(opts.UserDefined, obj.UserDefined) } } diff --git a/cmd/erasure-server-pool-consistency.go b/cmd/erasure-server-pool-consistency.go new file mode 100644 index 000000000..aa374f775 --- /dev/null +++ b/cmd/erasure-server-pool-consistency.go @@ -0,0 +1,317 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "context" + "maps" + "sort" + "strings" + "time" + + madmin "github.com/minio/madmin-go/v3" + xhttp "github.com/minio/minio/internal/http" +) + +// objectPoolInfos reads the addressed version in every pool, including pools +// draining their contents. The caller must hold the pools-layer object lock. +// An unreadable pool may contain a newer version or metadata; it is not absence. +func (z *erasureServerPools) objectPoolInfos(ctx context.Context, bucket, object string, opts ObjectOptions) ([]PoolObjInfo, error) { + opts.NoLock = true + opts.CheckPrecondFn = nil + var copies []PoolObjInfo + for i, pool := range z.serverPools { + oi, err := pool.GetObjectInfo(ctx, bucket, object, opts) + if err == nil || (oi.DeleteMarker && (isErrObjectNotFound(err) || isErrMethodNotAllowed(err))) { + copies = append(copies, PoolObjInfo{Index: i, ObjInfo: oi}) + continue + } + if !isErrObjectNotFound(err) && !isErrVersionNotFound(err) { + return nil, err + } + } + sort.Slice(copies, func(i, j int) bool { + a, b := copies[i], copies[j] + if a.ObjInfo.ModTime.Equal(b.ObjInfo.ModTime) { + return a.Index < b.Index + } + return a.ObjInfo.ModTime.After(b.ObjInfo.ModTime) + }) + if len(copies) == 0 { + if opts.VersionID != "" { + return nil, VersionNotFound{Bucket: bucket, Object: decodeDirObject(object), VersionID: opts.VersionID} + } + return nil, ObjectNotFound{Bucket: bucket, Object: decodeDirObject(object)} + } + return copies, nil +} + +// Healing also commits object metadata. Both queued healing and drive healing +// must use the same namespace as pooled writes, rather than racing them under +// a destination set's independent namespace. +func (z *erasureServerPools) healObjectInPool(ctx context.Context, er *erasureObjects, bucket, object, versionID string, opts madmin.HealOpts) (madmin.HealResultItem, error) { + if !z.SinglePool() && !opts.NoLock { + lk := z.NewNSLock(bucket, object) + lkctx, err := lk.GetLock(ctx, globalOperationTimeout) + if err != nil { + return madmin.HealResultItem{}, err + } + ctx = lkctx.Context() + defer lk.Unlock(lkctx) + opts.NoLock = true + } + return er.HealObject(ctx, bucket, object, versionID, opts) +} + +// mergePoolLockState resolves each independently ordered field. An ordered +// removal is a value in its own right; an absent unordered field is not one. +func mergePoolLockState(copies []PoolObjInfo) objectLockState { + state := storedObjectLockState(copies[0].ObjInfo.UserDefined) + for _, copy := range copies[1:] { + next := storedObjectLockState(copy.ObjInfo.UserDefined) + retentionTime, _ := time.Parse(time.RFC3339Nano, next.retentionTimestamp) + if state.retentionIsOlderThan(retentionTime) || (state.retentionTimestamp == "" && state.mode == "" && next.mode != "") { + state.mode, state.retainUntil, state.retentionTimestamp = next.mode, next.retainUntil, next.retentionTimestamp + } + holdTime, _ := time.Parse(time.RFC3339Nano, next.legalHoldTimestamp) + if state.legalHoldIsOlderThan(holdTime) || (state.legalHoldTimestamp == "" && state.legalHold == "" && next.legalHold != "") { + state.legalHold, state.legalHoldTimestamp = next.legalHold, next.legalHoldTimestamp + } + } + return state +} + +func replaceObjectLockMetadata(metadata map[string]string, state objectLockState) { + for _, key := range []string{ + strings.ToLower(xhttp.AmzObjectLockMode), strings.ToLower(xhttp.AmzObjectLockRetainUntilDate), + strings.ToLower(xhttp.AmzObjectLockLegalHold), + ReservedMetadataPrefixLower + ObjectLockRetentionTimestamp, + ReservedMetadataPrefixLower + ObjectLockLegalHoldTimestamp, + } { + // Metadata updates use maps.Copy at the set layer. Empty values also + // clear a previous value there, including timestamp-only removals. + if _, exists := metadata[key]; exists { + metadata[key] = "" + } + } + state.restoreRetention(metadata) + state.restoreLegalHold(metadata) +} + +func mergedPoolObjectInfo(copies []PoolObjInfo) ObjectInfo { + oi := copies[0].ObjInfo + oi.UserDefined = maps.Clone(oi.UserDefined) + if oi.UserDefined == nil { + oi.UserDefined = make(map[string]string) + } + replaceObjectLockMetadata(oi.UserDefined, mergePoolLockState(copies)) + for _, copy := range copies[1:] { + ts := copy.ObjInfo.UserDefined[ReservedMetadataPrefixLower+TaggingTimestamp] + stamp, _ := time.Parse(time.RFC3339Nano, ts) + if olderThan(oi.UserDefined[ReservedMetadataPrefixLower+TaggingTimestamp], stamp) { + oi.UserDefined[ReservedMetadataPrefixLower+TaggingTimestamp] = ts + oi.UserDefined[xhttp.AmzObjectTagging] = copy.ObjInfo.UserDefined[xhttp.AmzObjectTagging] + } + } + oi.UserTags = oi.UserDefined[xhttp.AmzObjectTagging] + return oi +} + +func (z *erasureServerPools) replicaObjectInfo(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { + if opts.VersionID == "" { + opts.VersionID = nullVersionID + } + copies, err := z.objectPoolInfos(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + if copies[0].ObjInfo.DeleteMarker { + return copies[0].ObjInfo, MethodNotAllowed{Bucket: bucket, Object: decodeDirObject(object), VersionID: opts.VersionID} + } + return mergedPoolObjectInfo(copies), nil +} + +// metadataPoolInfos resolves an unqualified metadata request to one logical +// version before collecting its copies, rather than mixing different versions. +func (z *erasureServerPools) metadataPoolInfos(ctx context.Context, bucket, object string, opts ObjectOptions) ([]PoolObjInfo, error) { + copies, err := z.objectPoolInfos(ctx, bucket, object, opts) + if err != nil { + return nil, err + } + if copies[0].ObjInfo.DeleteMarker { + return nil, MethodNotAllowed{Bucket: bucket, Object: decodeDirObject(object), VersionID: opts.VersionID} + } + if opts.VersionID == "" { + opts.VersionID = copies[0].ObjInfo.VersionID + if opts.VersionID == "" { + opts.VersionID = nullVersionID + } + return z.objectPoolInfos(ctx, bucket, object, opts) + } + return copies, nil +} + +func (z *erasureServerPools) updatePoolMetadata(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { + copies, err := z.metadataPoolInfos(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + updated := mergedPoolObjectInfo(copies) + before := maps.Clone(updated.UserDefined) + if opts.EvalMetadataFn != nil { + if _, err := opts.EvalMetadataFn(&updated, nil); err != nil { + return ObjectInfo{}, err + } + } + changes := make(map[string]string) + for key, value := range updated.UserDefined { + if prior, ok := before[key]; !ok || prior != value { + changes[key] = value + } + } + for key := range before { + if _, ok := updated.UserDefined[key]; !ok { + changes[key] = "" + } + } + state := storedObjectLockState(updated.UserDefined) + opts.VersionID = updated.VersionID + if opts.VersionID == "" { + opts.VersionID = nullVersionID + } + opts.NoLock = true + opts.EvalMetadataFn = func(oi *ObjectInfo, _ error) (ReplicateDecision, error) { + maps.Copy(oi.UserDefined, changes) + replaceObjectLockMetadata(oi.UserDefined, state) + for _, key := range []string{xhttp.AmzObjectTagging, ReservedMetadataPrefixLower + TaggingTimestamp} { + value, exists := updated.UserDefined[key] + if exists || oi.UserDefined[key] != "" { + oi.UserDefined[key] = value + } + } + return ReplicateDecision{}, nil + } + var primary ObjectInfo + for _, copy := range copies { + oi, err := z.serverPools[copy.Index].PutObjectMetadata(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + if copy.Index == copies[0].Index { + primary = oi + } + } + return primary, nil +} + +func reconcileStoredObjectTags(metadata, stored map[string]string) { + key := ReservedMetadataPrefixLower + TaggingTimestamp + stamp, err := time.Parse(time.RFC3339Nano, stored[key]) + if err != nil { + return + } + incoming, err := time.Parse(time.RFC3339Nano, metadata[key]) + if err != nil || !stamp.Before(incoming) { + metadata[key] = stored[key] + metadata[xhttp.AmzObjectTagging] = stored[xhttp.AmzObjectTagging] + } +} + +// retireReplicaCopies runs only after committing a replacement. Failures are +// returned to the caller, so a stale copy cannot be hidden behind a successful +// response. Data movement owns its source cleanup and does not use this helper. +func (z *erasureServerPools) retireReplicaCopies(ctx context.Context, bucket, object string, keep int, oi ObjectInfo) error { + versionID := oi.VersionID + if versionID == "" { + versionID = nullVersionID + } + copies, err := z.objectPoolInfos(ctx, bucket, object, ObjectOptions{VersionID: versionID}) + if err != nil { + return err + } + for _, copy := range copies { + if copy.Index == keep { + continue + } + _, err := z.serverPools[copy.Index].DeleteObject(ctx, bucket, object, + ObjectOptions{VersionID: versionID, NoLock: true, NoAuditLog: true}) + if err != nil && !isErrObjectNotFound(err) && !isErrVersionNotFound(err) { + return err + } + } + return nil +} + +// deleteObjectConditional evaluates the condition once against the logical +// version, then removes all its copies under the same lock as pooled writers. +func (z *erasureServerPools) deleteObjectConditional(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { + copies, err := z.objectPoolInfos(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + primary := copies[0] + if opts.CheckPrecondFn(primary.ObjInfo) { + return ObjectInfo{}, PreConditionFailed{} + } + opts.CheckPrecondFn = nil + opts.NoLock = true + if opts.EvalRetentionBypassFn != nil || opts.EvalMetadataFn != nil { + versions, err := z.metadataPoolInfos(ctx, bucket, object, opts) + if err != nil { + return ObjectInfo{}, err + } + logical := mergedPoolObjectInfo(versions) + if opts.EvalRetentionBypassFn != nil { + if err := opts.EvalRetentionBypassFn(logical, nil); err != nil { + return ObjectInfo{}, err + } + opts.EvalRetentionBypassFn = nil + } + if opts.EvalMetadataFn != nil { + decision, err := opts.EvalMetadataFn(&logical, nil) + if err != nil { + return ObjectInfo{}, err + } + if decision.ReplicateAny() { + opts.SetDeleteReplicationState(decision, opts.VersionID) + } + opts.EvalMetadataFn = nil + } + } + if opts.VersionID == "" && (opts.Versioned || opts.VersionSuspended) { + // A single new delete marker hides the current version. Older versions + // remain history and must not be removed by an unqualified DELETE. + oi, err := z.serverPools[primary.Index].DeleteObject(ctx, bucket, object, opts) + oi.Name = decodeDirObject(object) + oi.replicationDecision = opts.DeleteReplication.ReplicateDecisionStr + return oi, err + } + // Retire non-authoritative copies first. If cleanup fails, retain the + // authoritative version and report the error instead of acknowledging a + // deletion that would expose an older copy. + for _, copy := range copies[1:] { + _, err := z.serverPools[copy.Index].DeleteObject(ctx, bucket, object, opts) + if err != nil && !isErrObjectNotFound(err) && !isErrVersionNotFound(err) { + return ObjectInfo{}, err + } + } + oi, err := z.serverPools[primary.Index].DeleteObject(ctx, bucket, object, opts) + oi.Name = decodeDirObject(object) + oi.replicationDecision = opts.DeleteReplication.ReplicateDecisionStr + return oi, err +} diff --git a/cmd/erasure-server-pool-consistency_test.go b/cmd/erasure-server-pool-consistency_test.go new file mode 100644 index 000000000..99c117628 --- /dev/null +++ b/cmd/erasure-server-pool-consistency_test.go @@ -0,0 +1,575 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "bytes" + "context" + "fmt" + "io" + "maps" + "strings" + "sync" + "testing" + "time" + + madmin "github.com/minio/madmin-go/v3" + xhttp "github.com/minio/minio/internal/http" +) + +func consistencyPools(t *testing.T) (*erasureServerPools, string) { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + obj, dirs, err := prepareErasurePoolsWithContext(ctx) + if err != nil { + cancel() + t.Fatal(err) + } + z := obj.(*erasureServerPools) + previous := newObjectLayerFn() + setObjectLayer(z) + t.Cleanup(func() { cancel(); z.Shutdown(context.Background()); removeRoots(dirs); setObjectLayer(previous) }) + bucket := "pool-consistency" + if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil { + t.Fatal(err) + } + return z, bucket +} + +func putConsistencyObject(t *testing.T, z *erasureServerPools, bucket, object string, pool int, body string, opts ObjectOptions) ObjectInfo { + t.Helper() + oi, err := z.serverPools[pool].PutObject(t.Context(), bucket, object, + mustGetPutObjReader(t, bytes.NewBufferString(body), int64(len(body)), "", ""), opts) + if err != nil { + t.Fatal(err) + } + return oi +} + +func TestPoolsConditionalDeleteVersionSelection(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "split-versions" + old := putConsistencyObject(t, z, bucket, object, 1, "old", ObjectOptions{Versioned: true, MTime: time.Now().Add(-time.Hour)}) + latest := putConsistencyObject(t, z, bucket, object, 0, "latest", ObjectOptions{Versioned: true}) + _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + Versioned: true, VersionID: old.VersionID, HasIfMatch: true, + CheckPrecondFn: func(oi ObjectInfo) bool { return oi.ETag != old.ETag }, + }) + if err != nil { + t.Fatalf("delete addressed version in pool 1: %v", err) + } + if _, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: old.VersionID}); !isErrVersionNotFound(err) { + t.Errorf("addressed version survived: %v", err) + } + if oi, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); err != nil || oi.VersionID != latest.VersionID { + t.Errorf("delete changed the latest version: %+v, %v", oi, err) + } +} + +func TestPoolsConditionalDeleteDuplicateVersion(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "duplicate-version" + oi := putConsistencyObject(t, z, bucket, object, 0, "payload", ObjectOptions{Versioned: true}) + putConsistencyObject(t, z, bucket, object, 1, "payload", ObjectOptions{Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime}) + _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, HasIfMatch: true, + CheckPrecondFn: func(current ObjectInfo) bool { return current.ETag != oi.ETag }, + }) + if err != nil { + t.Fatal(err) + } + if _, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) { + t.Fatalf("success left a readable duplicate version: %v", err) + } +} + +func TestPoolsConditionalDeleteReportsOtherPoolFailure(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "delete-failure" + putConsistencyObject(t, z, bucket, object, 1, "older", ObjectOptions{MTime: time.Now().Add(-time.Hour)}) + oi := putConsistencyObject(t, z, bucket, object, 0, "latest", ObjectOptions{}) + set := z.serverPools[1].getHashedSet(object) + getDisks := set.getDisks + faulty := append([]StorageAPI(nil), getDisks()...) + for i := range faulty { + faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object} + } + set.getDisks = func() []StorageAPI { return faulty } + defer func() { set.getDisks = getDisks }() + _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + HasIfMatch: true, CheckPrecondFn: func(current ObjectInfo) bool { return current.ETag != oi.ETag }, + }) + if err == nil { + t.Fatal("delete reported success despite the other pool's write failure") + } +} + +func TestPoolsConditionalDeleteSerializesPut(t *testing.T) { + testPoolsConditionalDeleteWriter(t, false) +} + +func TestPoolsConditionalDeleteSerializesCompletion(t *testing.T) { + testPoolsConditionalDeleteWriter(t, true) +} + +func testPoolsConditionalDeleteWriter(t *testing.T, multipart bool) { + z, bucket := consistencyPools(t) + const object = "concurrent-put" + initial := putConsistencyObject(t, z, bucket, object, 1, "before", ObjectOptions{}) + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + reader := mustGetPutObjReader(t, bytes.NewBufferString("after"), 5, "", "") + destination := 1 + write := func() error { + _, err := z.PutObject(ctx, bucket, object, reader, ObjectOptions{DataMovement: true, SrcPoolIdx: 0, DstPoolIdx: &destination}) + return err + } + if multipart { + mp, err := z.serverPools[1].NewMultipartUpload(ctx, bucket, object, ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + part, err := z.serverPools[1].PutObjectPart(ctx, bucket, object, mp.UploadID, 1, reader, ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + write = func() error { + _, err := z.CompleteMultipartUpload(ctx, bucket, object, mp.UploadID, + []CompletePart{{PartNumber: 1, ETag: part.ETag}}, ObjectOptions{}) + return err + } + } + checked, resume := make(chan struct{}), make(chan struct{}) + var resumeOnce sync.Once + release := func() { resumeOnce.Do(func() { close(resume) }) } + defer release() + deleted := make(chan error, 1) + go func() { + _, err := z.DeleteObject(ctx, bucket, object, ObjectOptions{ + HasIfMatch: true, + CheckPrecondFn: func(oi ObjectInfo) bool { + close(checked) + <-resume + return oi.ETag != initial.ETag + }, + }) + deleted <- err + }() + select { + case <-checked: + case err := <-deleted: + t.Fatalf("delete did not evaluate the precondition: %v", err) + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + written := make(chan error, 1) + go func() { written <- write() }() + var early bool + select { + case err := <-written: + early = true + if err != nil { + t.Errorf("concurrent PUT: %v", err) + } + case <-time.After(250 * time.Millisecond): + } + release() + if err := <-deleted; err != nil { + t.Fatal(err) + } + if !early { + if err := <-written; err != nil { + t.Fatal(err) + } + } + if early { + t.Error("PUT committed while DELETE was between its comparison and removal") + } + if _, err := z.GetObjectInfo(ctx, bucket, object, ObjectOptions{}); err != nil { + t.Errorf("matching the old ETag removed the concurrent replacement: %v", err) + } +} + +type consistencyGateReader struct { + io.Reader + entered, resume chan struct{} + once sync.Once + release sync.Once +} + +func (r *consistencyGateReader) Read(p []byte) (int, error) { + r.once.Do(func() { close(r.entered); <-r.resume }) + return r.Reader.Read(p) +} + +func TestPoolsReplicaSerializesMetadataAndHealing(t *testing.T) { + for _, heal := range []bool{false, true} { + t.Run(fmt.Sprintf("heal=%t", heal), func(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "metadata-race" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + meta := poolLockMetadata("GOVERNANCE", "OFF", old, old) + oi := putConsistencyObject(t, z, bucket, object, 1, "data", ObjectOptions{Versioned: true, UserDefined: meta}) + z.poolMetaMutex.Lock() + z.poolMeta.Pools[0].Decommission = &PoolDecommissionInfo{} + z.poolMetaMutex.Unlock() + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + gate := &consistencyGateReader{Reader: strings.NewReader("data"), entered: make(chan struct{}), resume: make(chan struct{})} + release := func() { gate.release.Do(func() { close(gate.resume) }) } + defer release() + reader := mustGetPutObjReader(t, gate, 4, "", "") + written := make(chan error, 1) + go func() { + _, err := z.PutObject(ctx, bucket, object, reader, ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, + ReplicaLockReconcile: true, UserDefined: maps.Clone(meta), + }) + written <- err + }() + select { + case <-gate.entered: + case err := <-written: + t.Fatalf("write failed before consuming the body: %v", err) + case <-ctx.Done(): + t.Fatal(ctx.Err()) + } + mutated := make(chan error, 1) + go func() { + var err error + if heal { + _, err = z.healObjectInPool(ctx, z.serverPools[1].getHashedSet(object), bucket, object, oi.VersionID, madmin.HealOpts{}) + } else { + _, err = z.PutObjectMetadata(ctx, bucket, object, ObjectOptions{ + VersionID: oi.VersionID, MTime: oi.ModTime, + EvalMetadataFn: func(current *ObjectInfo, _ error) (ReplicateDecision, error) { + current.UserDefined[strings.ToLower(xhttp.AmzObjectLockLegalHold)] = "ON" + current.UserDefined[ReservedMetadataPrefixLower+ObjectLockLegalHoldTimestamp] = recent + return ReplicateDecision{}, nil + }, + }) + } + mutated <- err + }() + var early bool + select { + case err := <-mutated: + early = true + t.Errorf("mutation completed during the paused replica write: %v", err) + case <-time.After(250 * time.Millisecond): + } + release() + if err := <-written; err != nil { + t.Fatal(err) + } + if !early { + if err := <-mutated; err != nil { + t.Fatal(err) + } + } + got, err := z.GetObjectInfo(ctx, bucket, object, ObjectOptions{VersionID: oi.VersionID}) + if err != nil { + t.Fatal(err) + } + if !heal && storedObjectLockState(got.UserDefined).legalHold != "ON" { + t.Error("replica write rolled back the newer metadata update") + } + }) + } +} + +func TestPoolsConditionalDeletePreservesVersionHistory(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "conditional-marker" + older := putConsistencyObject(t, z, bucket, object, 1, "older", ObjectOptions{Versioned: true, MTime: time.Now().Add(-time.Hour)}) + latest := putConsistencyObject(t, z, bucket, object, 0, "latest", ObjectOptions{Versioned: true}) + _, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + Versioned: true, + CheckPrecondFn: func(oi ObjectInfo) bool { return oi.ETag != older.ETag }, + }) + if !isErrPreconditionFailed(err) { + t.Fatalf("wrong latest ETag: %v", err) + } + marker, err := z.DeleteObject(t.Context(), bucket, object, ObjectOptions{ + Versioned: true, + CheckPrecondFn: func(oi ObjectInfo) bool { return oi.ETag != latest.ETag }, + }) + if err != nil || !marker.DeleteMarker { + t.Fatalf("logical delete did not create a marker: %+v, %v", marker, err) + } + if got, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); !isErrObjectNotFound(err) || !got.DeleteMarker { + t.Errorf("delete marker did not hide the split object: %+v, %v", got, err) + } + for _, version := range []string{older.VersionID, latest.VersionID} { + if _, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: version}); err != nil { + t.Errorf("conditional marker removed history %s: %v", version, err) + } + } +} + +func poolLockMetadata(retention, hold, retentionTime, holdTime string) map[string]string { + return map[string]string{ + strings.ToLower(xhttp.AmzObjectLockMode): retention, + strings.ToLower(xhttp.AmzObjectLockRetainUntilDate): "2030-01-01T00:00:00Z", + strings.ToLower(xhttp.AmzObjectLockLegalHold): hold, + ReservedMetadataPrefixLower + ObjectLockRetentionTimestamp: retentionTime, + ReservedMetadataPrefixLower + ObjectLockLegalHoldTimestamp: holdTime, + } +} + +func TestPoolsReplicaIndependentLockWinners(t *testing.T) { + for _, tc := range []struct { + name string + multipart, rebalance bool + }{ + {"put-decommission", false, false}, + {"put-rebalance", false, true}, + {"multipart-decommission", true, false}, + {"multipart-rebalance", true, true}, + } { + t.Run(tc.name, func(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "split-lock-state" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + first := poolLockMetadata("GOVERNANCE", "OFF", recent, old) + second := poolLockMetadata("", "ON", old, recent) + oi := putConsistencyObject(t, z, bucket, object, 0, "data", ObjectOptions{Versioned: true, UserDefined: first}) + putConsistencyObject(t, z, bucket, object, 1, "data", ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, UserDefined: second, + }) + opts := ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, + ReplicaLockReconcile: true, UserDefined: poolLockMetadata("", "OFF", old, old), + } + var uploadID string + var parts []CompletePart + if tc.multipart { + mp, err := z.serverPools[1].NewMultipartUpload(t.Context(), bucket, object, opts) + if err != nil { + t.Fatal(err) + } + uploadID = mp.UploadID + part, err := z.serverPools[1].PutObjectPart(t.Context(), bucket, object, uploadID, 1, + mustGetPutObjReader(t, bytes.NewBufferString("data"), 4, "", ""), ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + parts = []CompletePart{{PartNumber: 1, ETag: part.ETag}} + } + // The upload and both versions predate the routing change. Retention's + // winner remains in the draining pool; legal hold's winner is in pool 1. + if tc.rebalance { + z.rebalMu.Lock() + z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{ + {Participating: true, Info: rebalanceInfo{Status: rebalStarted}}, {}, + }} + z.rebalMu.Unlock() + } else { + z.poolMetaMutex.Lock() + z.poolMeta.Pools[0].Decommission = &PoolDecommissionInfo{} + z.poolMetaMutex.Unlock() + } + var err error + if tc.multipart { + _, err = z.CompleteMultipartUpload(t.Context(), bucket, object, uploadID, parts, + ObjectOptions{Versioned: true, MTime: oi.ModTime, ReplicaLockReconcile: true}) + } else { + _, err = z.PutObject(t.Context(), bucket, object, + mustGetPutObjReader(t, bytes.NewBufferString("data"), 4, "", ""), opts) + } + if err != nil { + t.Fatal(err) + } + got, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}) + if err != nil { + t.Fatal(err) + } + state := storedObjectLockState(got.UserDefined) + if state.mode != "GOVERNANCE" || state.retentionTimestamp != recent || state.legalHold != "ON" || state.legalHoldTimestamp != recent { + t.Errorf("pooled read lost independently ordered lock state: %+v", state) + } + if _, err := z.serverPools[0].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) { + t.Errorf("successful replacement left its stale competing copy: %v", err) + } + }) + } +} + +func TestPoolsReplicaSoleDrainingOwner(t *testing.T) { + for _, rebalance := range []bool{false, true} { + for _, null := range []bool{false, true} { + t.Run(fmt.Sprintf("rebalance=%t/null=%t", rebalance, null), func(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "sole-owner" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + meta := poolLockMetadata("", "ON", recent, recent) + delete(meta, strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)) + oi := putConsistencyObject(t, z, bucket, object, 0, "data", ObjectOptions{Versioned: !null, UserDefined: meta}) + versionID := oi.VersionID + if null { + versionID = nullVersionID + // A different latest version must not donate its lock to the null version. + putConsistencyObject(t, z, bucket, object, 1, "other-version", ObjectOptions{Versioned: true}) + } + if rebalance { + z.rebalMu.Lock() + z.rebalMeta = &rebalanceMeta{PoolStats: []*rebalanceStats{ + {Participating: true, Info: rebalanceInfo{Status: rebalStarted}}, {}, + }} + z.rebalMu.Unlock() + } else { + z.poolMetaMutex.Lock() + z.poolMeta.Pools[0].Decommission = &PoolDecommissionInfo{} + z.poolMetaMutex.Unlock() + } + _, err := z.PutObject(t.Context(), bucket, object, + mustGetPutObjReader(t, bytes.NewBufferString("data"), 4, "", ""), ObjectOptions{ + Versioned: !null, VersionSuspended: null, VersionID: versionID, MTime: oi.ModTime, + ReplicaLockReconcile: true, UserDefined: poolLockMetadata("GOVERNANCE", "OFF", old, old), + }) + if err != nil { + t.Fatal(err) + } + got, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: versionID}) + if err != nil { + t.Fatal(err) + } + state := storedObjectLockState(got.UserDefined) + if state.mode != "" || state.retainUntil != "" || state.retentionTimestamp != recent || state.legalHold != "ON" { + t.Errorf("lost a removal or legal hold from the sole draining owner: %+v", state) + } + }) + } + } +} + +func TestPoolsMetadataUpdateUsesMergedVersion(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "metadata-update" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + oi := putConsistencyObject(t, z, bucket, object, 0, "data", ObjectOptions{ + Versioned: true, UserDefined: poolLockMetadata("GOVERNANCE", "OFF", recent, old), + }) + putConsistencyObject(t, z, bucket, object, 1, "data", ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, UserDefined: poolLockMetadata("", "ON", old, recent), + }) + latest := putConsistencyObject(t, z, bucket, object, 1, "latest", ObjectOptions{Versioned: true}) + called := 0 + _, err := z.PutObjectMetadata(t.Context(), bucket, object, ObjectOptions{ + VersionID: oi.VersionID, MTime: oi.ModTime, + EvalMetadataFn: func(current *ObjectInfo, _ error) (ReplicateDecision, error) { + called++ + state := storedObjectLockState(current.UserDefined) + if state.mode != "GOVERNANCE" || state.legalHold != "ON" { + return ReplicateDecision{}, fmt.Errorf("metadata policy evaluated stale state: %+v", state) + } + current.UserDefined["custom-update"] = "preserved" + return ReplicateDecision{}, nil + }, + }) + if err != nil || called != 1 { + t.Fatalf("metadata update: %v; callback count %d", err, called) + } + for i, pool := range z.serverPools { + got, err := pool.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}) + if err != nil { + t.Fatal(err) + } + state := storedObjectLockState(got.UserDefined) + if state.mode != "GOVERNANCE" || state.legalHold != "ON" || got.UserDefined["custom-update"] != "preserved" { + t.Errorf("pool %d did not receive the merged update: %+v", i, got.UserDefined) + } + } + if current, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}); err != nil || current.VersionID != latest.VersionID { + t.Errorf("updating an older version changed the latest: %+v, %v", current, err) + } +} + +func TestPoolsReplicaMetadataCopyReconcilesLockAndTags(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "replica-metadata-copy" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + meta := poolLockMetadata("GOVERNANCE", "OFF", recent, old) + meta[xhttp.AmzObjectTagging] = "key=new" + meta[ReservedMetadataPrefixLower+TaggingTimestamp] = recent + oi := putConsistencyObject(t, z, bucket, object, 0, "data", ObjectOptions{Versioned: true, UserDefined: meta}) + putConsistencyObject(t, z, bucket, object, 1, "data", ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, UserDefined: poolLockMetadata("", "ON", old, recent), + }) + stale := oi + stale.metadataOnly = true + stale.UserDefined = maps.Clone(oi.UserDefined) + stale.UserDefined[xhttp.AmzObjectTagging] = "key=old" + stale.UserDefined[ReservedMetadataPrefixLower+TaggingTimestamp] = old + _, err := z.CopyObject(t.Context(), bucket, object, bucket, object, stale, + ObjectOptions{VersionID: oi.VersionID}, ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, ReplicaLockReconcile: true, + }) + if err != nil { + t.Fatal(err) + } + got, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}) + if err != nil { + t.Fatal(err) + } + state := storedObjectLockState(got.UserDefined) + if state.mode != "GOVERNANCE" || state.legalHold != "ON" || got.UserTags != "key=new" { + t.Errorf("metadata replication rolled back a newer field: %+v", got.UserDefined) + } +} + +func TestPoolsReplicaCleanupFailureCanRetry(t *testing.T) { + z, bucket := consistencyPools(t) + const object = "replica-cleanup-failure" + old, recent := "2026-09-09T09:00:00Z", "2026-09-09T10:00:00Z" + oi := putConsistencyObject(t, z, bucket, object, 0, "data", ObjectOptions{ + Versioned: true, UserDefined: poolLockMetadata("GOVERNANCE", "ON", recent, recent), + }) + z.poolMetaMutex.Lock() + z.poolMeta.Pools[0].Decommission = &PoolDecommissionInfo{} + z.poolMetaMutex.Unlock() + set := z.serverPools[0].getHashedSet(object) + getDisks := set.getDisks + faulty := append([]StorageAPI(nil), getDisks()...) + for i := range faulty { + faulty[i] = accessMoveDeleteFaultDisk{StorageAPI: faulty[i], bucket: bucket, object: object, version: oi.VersionID} + } + set.getDisks = func() []StorageAPI { return faulty } + defer func() { set.getDisks = getDisks }() + write := func() error { + _, err := z.PutObject(t.Context(), bucket, object, + mustGetPutObjReader(t, bytes.NewBufferString("data"), 4, "", ""), ObjectOptions{ + Versioned: true, VersionID: oi.VersionID, MTime: oi.ModTime, ReplicaLockReconcile: true, + UserDefined: poolLockMetadata("", "OFF", old, old), + }) + return err + } + if err := write(); err == nil { + t.Fatal("replica replacement hid a competing-copy cleanup failure") + } + if got, err := z.serverPools[1].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}); err != nil || storedObjectLockState(got.UserDefined).legalHold != "ON" { + t.Fatalf("cleanup failure lost the committed reconciled version: %+v, %v", got, err) + } + set.getDisks = getDisks + if err := write(); err != nil { + t.Fatalf("retry could not finish cleanup: %v", err) + } + if _, err := z.serverPools[0].GetObjectInfo(t.Context(), bucket, object, ObjectOptions{VersionID: oi.VersionID}); !isErrVersionNotFound(err) { + t.Errorf("retry left the competing version: %v", err) + } +} diff --git a/cmd/erasure-server-pool-decom_test.go b/cmd/erasure-server-pool-decom_test.go index 7f4c8ef17..b1312bba0 100644 --- a/cmd/erasure-server-pool-decom_test.go +++ b/cmd/erasure-server-pool-decom_test.go @@ -23,6 +23,10 @@ import ( ) func prepareErasurePools() (ObjectLayer, []string, error) { + return prepareErasurePoolsWithContext(context.Background()) +} + +func prepareErasurePoolsWithContext(ctx context.Context) (ObjectLayer, []string, error) { nDisks := 32 fsDirs, err := getRandomDisks(nDisks) if err != nil { @@ -32,7 +36,7 @@ func prepareErasurePools() (ObjectLayer, []string, error) { pools := mustGetPoolEndpoints(0, fsDirs[:16]...) pools = append(pools, mustGetPoolEndpoints(1, fsDirs[16:]...)...) - objLayer, _, err := initObjectLayer(context.Background(), pools) + objLayer, _, err := initObjectLayer(ctx, pools) if err != nil { removeRoots(fsDirs) return nil, nil, err diff --git a/cmd/erasure-server-pool.go b/cmd/erasure-server-pool.go index b5eda01d5..3755f324c 100644 --- a/cmd/erasure-server-pool.go +++ b/cmd/erasure-server-pool.go @@ -645,6 +645,12 @@ func (z *erasureServerPools) getPoolIdxNoLock(ctx context.Context, bucket, objec // if none are found falls back to most available space pool, this function is // designed to be only used by PutObject, CopyObject (newObject creation) and NewMultipartUpload. func (z *erasureServerPools) getPoolIdx(ctx context.Context, bucket, object string, size int64, dstPoolIdx *int) (idx int, err error) { + return z.getWritePoolIdx(ctx, bucket, object, size, dstPoolIdx, false) +} + +// getWritePoolIdx keeps the write-allocation policy when the caller already +// holds the pools-layer object lock and must not reacquire a set read lock. +func (z *erasureServerPools) getWritePoolIdx(ctx context.Context, bucket, object string, size int64, dstPoolIdx *int, noLock bool) (idx int, err error) { if dstPoolIdx != nil { if *dstPoolIdx < 0 || *dstPoolIdx >= len(z.serverPools) { return -1, errInvalidArgument @@ -656,6 +662,7 @@ func (z *erasureServerPools) getPoolIdx(ctx context.Context, bucket, object stri } pinfo, _, err := z.getPoolInfoExistingWithOpts(ctx, bucket, object, ObjectOptions{ + NoLock: noLock, SkipDecommissioned: true, SkipRebalancing: true, }) @@ -1164,7 +1171,18 @@ func (z *erasureServerPools) PutObject(ctx context.Context, bucket string, objec return z.serverPools[0].PutObject(ctx, bucket, object, data, opts) } - idx, err := z.getPoolIdx(ctx, bucket, object, data.Size(), dataMovementDstPool(opts)) + if !opts.NoLock { + lk := z.NewNSLock(bucket, object) + lkctx, err := lk.GetLock(ctx, globalOperationTimeout) + if err != nil { + return ObjectInfo{}, err + } + ctx = lkctx.Context() + defer lk.Unlock(lkctx) + } + opts.NoLock = true + + idx, err := z.getWritePoolIdx(ctx, bucket, object, data.Size(), dataMovementDstPool(opts), true) if err != nil { return ObjectInfo{}, err } @@ -1178,7 +1196,15 @@ func (z *erasureServerPools) PutObject(ctx context.Context, bucket string, objec } } - return z.serverPools[idx].PutObject(ctx, bucket, object, data, opts) + if opts.ReplicaLockReconcile || opts.DataMovement { + opts.ReplicaLockReconcile = true + opts.replicaObjectInfo = z.replicaObjectInfo + } + oi, err := z.serverPools[idx].PutObject(ctx, bucket, object, data, opts) + if err == nil && opts.ReplicaLockReconcile && !opts.DataMovement { + err = z.retireReplicaCopies(ctx, bucket, object, idx, oi) + } + return oi, err } func (z *erasureServerPools) deletePrefix(ctx context.Context, bucket string, prefix string) error { @@ -1199,15 +1225,7 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob object = encodeDirObject(object) } - // Acquire a write lock before deleting the object. - // - // NOTE: this lock is taken at the server-pool level. The conditional - // (If-Match) precondition below relies on this lock making the read-check- - // delete sequence atomic. That holds for a single erasure set: the write - // path (PutObject) locks at the destination set, which shares this lock's - // namespace only within one set. Multi-pool conditional-delete atomicity - // (concurrent writers across pools, cross-pool version selection) is a - // separate concern tracked as a follow-up. + // Serialize logical deletion with pooled writes and metadata updates. lk := z.NewNSLock(bucket, object) lkctx, err := lk.GetLock(ctx, globalDeleteOperationTimeout) if err != nil { @@ -1242,6 +1260,10 @@ func (z *erasureServerPools) DeleteObject(ctx context.Context, bucket string, ob return objInfo, err } + if !z.SinglePool() && opts.CheckPrecondFn != nil { + return z.deleteObjectConditional(ctx, bucket, object, opts) + } + gopts := opts gopts.NoLock = true @@ -1377,8 +1399,12 @@ func (z *erasureServerPools) deleteObjectFromAllPools(ctx context.Context, bucke // the delete call tries to clean up other pools during DeleteObject call. objInfo = dobjects[0] objInfo.Name = decodeDirObject(object) - err = derrs[0] - return objInfo, err + for _, derr := range derrs { + if derr != nil && !isErrObjectNotFound(derr) && !isErrVersionNotFound(derr) { + return objInfo, derr + } + } + return objInfo, derrs[0] } func (z *erasureServerPools) DeleteObjects(ctx context.Context, bucket string, objects []ObjectToDelete, opts ObjectOptions) ([]DeletedObject, []error) { @@ -1466,6 +1492,22 @@ func (z *erasureServerPools) CopyObject(ctx context.Context, srcBucket, srcObjec dstOpts.NoLock = true } + if !z.SinglePool() && cpSrcDstSame && srcInfo.metadataOnly && dstOpts.ReplicaLockReconcile && srcOpts.VersionID == dstOpts.VersionID { + copies, err := z.metadataPoolInfos(ctx, dstBucket, dstObject, dstOpts) + if err != nil { + return ObjectInfo{}, err + } + stored := mergedPoolObjectInfo(copies) + reconcileStoredObjectLock(srcInfo.UserDefined, storedObjectLockState(stored.UserDefined)) + reconcileStoredObjectTags(srcInfo.UserDefined, stored.UserDefined) + idx := copies[0].Index + oi, err := z.serverPools[idx].CopyObject(ctx, srcBucket, srcObject, dstBucket, dstObject, srcInfo, srcOpts, dstOpts) + if err == nil { + err = z.retireReplicaCopies(ctx, dstBucket, dstObject, idx, oi) + } + return oi, err + } + poolIdx, err := z.getPoolIdxNoLock(ctx, dstBucket, dstObject, srcInfo.Size, dataMovementDstPool(dstOpts)) if err != nil { return objInfo, err @@ -1477,6 +1519,17 @@ func (z *erasureServerPools) CopyObject(ctx context.Context, srcBucket, srcObjec } } + if !z.SinglePool() && dstOpts.ReplicaLockReconcile { + dstOpts.replicaObjectInfo = z.replicaObjectInfo + if !dstOpts.DataMovement { + defer func() { + if err == nil { + err = z.retireReplicaCopies(ctx, dstBucket, dstObject, poolIdx, objInfo) + } + }() + } + } + // CopyObjectHandler predicts the outcome of this decision in // copyRewritesObjectData(); keep the two in sync. if cpSrcDstSame && srcInfo.metadataOnly { @@ -1511,6 +1564,8 @@ func (z *erasureServerPools) CopyObject(ctx context.Context, srcBucket, srcObjec EncryptFn: dstOpts.EncryptFn, WantChecksum: dstOpts.WantChecksum, WantServerSideChecksumType: dstOpts.WantServerSideChecksumType, + ReplicaLockReconcile: dstOpts.ReplicaLockReconcile, + replicaObjectInfo: dstOpts.replicaObjectInfo, } return z.serverPools[poolIdx].PutObject(ctx, dstBucket, dstObject, srcInfo.PutObjReader, putOpts) @@ -2158,6 +2213,19 @@ func (z *erasureServerPools) CompleteMultipartUpload(ctx context.Context, bucket } }() + if !z.SinglePool() { + if !opts.NoLock { + lk := z.NewNSLock(bucket, encodeDirObject(object)) + lkctx, err := lk.GetLock(ctx, globalOperationTimeout) + if err != nil { + return objInfo, err + } + ctx = lkctx.Context() + defer lk.Unlock(lkctx) + } + opts.NoLock = true + } + // Hold write locks to verify uploaded parts, also disallows any // parallel PutObjectPart() requests. uploadIDLock := z.NewNSLock(bucket, pathJoin(object, uploadID)) @@ -2172,13 +2240,21 @@ func (z *erasureServerPools) CompleteMultipartUpload(ctx context.Context, bucket return z.serverPools[0].CompleteMultipartUpload(ctx, bucket, object, uploadID, uploadedParts, opts) } + if opts.ReplicaLockReconcile || opts.DataMovement { + opts.ReplicaLockReconcile = true + opts.replicaObjectInfo = z.replicaObjectInfo + } + for idx, pool := range z.serverPools { if z.IsSuspended(idx) { continue } objInfo, err = pool.CompleteMultipartUpload(ctx, bucket, object, uploadID, uploadedParts, opts) if err == nil { - return objInfo, nil + if opts.ReplicaLockReconcile && !opts.DataMovement { + err = z.retireReplicaCopies(ctx, bucket, encodeDirObject(object), idx, objInfo) + } + return objInfo, err } if _, ok := err.(InvalidUploadID); ok { // upload id not found move to next pool @@ -2705,7 +2781,7 @@ func (z *erasureServerPools) HealObject(ctx context.Context, bucket, object, ver wg.Add(1) go func(idx int, pool *erasureSets) { defer wg.Done() - result, err := pool.HealObject(ctx, bucket, object, versionID, opts) + result, err := z.healObjectInPool(ctx, pool.getHashedSet(object), bucket, object, versionID, opts) result.Object = decodeDirObject(result.Object) errs[idx] = err results[idx] = result @@ -2974,15 +3050,7 @@ func (z *erasureServerPools) PutObjectMetadata(ctx context.Context, bucket, obje defer lk.Unlock(lkctx) } - opts.MetadataChg = true - opts.NoLock = true - // We don't know the size here set 1GiB at least. - idx, err := z.getPoolIdxExistingWithOpts(ctx, bucket, object, opts) - if err != nil { - return ObjectInfo{}, err - } - - return z.serverPools[idx].PutObjectMetadata(ctx, bucket, object, opts) + return z.updatePoolMetadata(ctx, bucket, object, opts) } // PutObjectTags - replace or add tags to an existing object @@ -3003,44 +3071,31 @@ func (z *erasureServerPools) PutObjectTags(ctx context.Context, bucket, object s defer lk.Unlock(lkctx) } - opts.MetadataChg = true - opts.NoLock = true - - // We don't know the size here set 1GiB at least. - idx, err := z.getPoolIdxExistingWithOpts(ctx, bucket, object, opts) + copies, err := z.metadataPoolInfos(ctx, bucket, object, opts) if err != nil { return ObjectInfo{}, err } - - return z.serverPools[idx].PutObjectTags(ctx, bucket, object, tags, opts) + opts.NoLock = true + opts.VersionID = copies[0].ObjInfo.VersionID + if opts.VersionID == "" { + opts.VersionID = nullVersionID + } + var primary ObjectInfo + for _, copy := range copies { + oi, err := z.serverPools[copy.Index].PutObjectTags(ctx, bucket, object, tags, opts) + if err != nil { + return ObjectInfo{}, err + } + if copy.Index == copies[0].Index { + primary = oi + } + } + return primary, nil } // DeleteObjectTags - delete object tags from an existing object func (z *erasureServerPools) DeleteObjectTags(ctx context.Context, bucket, object string, opts ObjectOptions) (ObjectInfo, error) { - object = encodeDirObject(object) - if z.SinglePool() { - return z.serverPools[0].DeleteObjectTags(ctx, bucket, object, opts) - } - - if !opts.NoLock { - // Lock the object before deleting tags. - lk := z.NewNSLock(bucket, object) - lkctx, err := lk.GetLock(ctx, globalOperationTimeout) - if err != nil { - return ObjectInfo{}, err - } - ctx = lkctx.Context() - defer lk.Unlock(lkctx) - } - - opts.MetadataChg = true - opts.NoLock = true - idx, err := z.getPoolIdxExistingWithOpts(ctx, bucket, object, opts) - if err != nil { - return ObjectInfo{}, err - } - - return z.serverPools[idx].DeleteObjectTags(ctx, bucket, object, opts) + return z.PutObjectTags(ctx, bucket, object, "", opts) } // GetObjectTags - get object tags from an existing object @@ -3050,7 +3105,7 @@ func (z *erasureServerPools) GetObjectTags(ctx context.Context, bucket, object s return z.serverPools[0].GetObjectTags(ctx, bucket, object, opts) } - oi, _, err := z.getLatestObjectInfoWithIdx(ctx, bucket, object, opts) + oi, err := z.GetObjectInfo(ctx, bucket, object, opts) if err != nil { return nil, err } diff --git a/cmd/global-heal.go b/cmd/global-heal.go index 22bf54ca6..ca78df818 100644 --- a/cmd/global-heal.go +++ b/cmd/global-heal.go @@ -166,6 +166,12 @@ func (er *erasureObjects) healErasureSet(ctx context.Context, buckets []string, if objAPI == nil { return errServerNotInitialized } + healInPool := er.HealObject + if z, ok := objAPI.(*erasureServerPools); ok && !z.SinglePool() { + healInPool = func(ctx context.Context, bucket, object, versionID string, opts madmin.HealOpts) (madmin.HealResultItem, error) { + return z.healObjectInPool(ctx, er, bucket, object, versionID, opts) + } + } started := tracker.Started if started.IsZero() || started.Equal(timeSentinel) { @@ -419,7 +425,7 @@ func (er *erasureObjects) healErasureSet(ctx context.Context, buckets []string, var result healEntryResult fivs, err := entry.fileInfoVersions(bucket) if err != nil { - res, err := er.HealObject(ctx, bucket, encodedEntryName, "", + res, err := healInPool(ctx, bucket, encodedEntryName, "", madmin.HealOpts{ ScanMode: scanMode, Remove: healDeleteDangling, @@ -455,7 +461,7 @@ func (er *erasureObjects) healErasureSet(ctx context.Context, buckets []string, continue } - res, err := er.HealObject(ctx, bucket, encodedEntryName, + res, err := healInPool(ctx, bucket, encodedEntryName, version.VersionID, madmin.HealOpts{ ScanMode: scanMode, Remove: healDeleteDangling, diff --git a/cmd/object-api-interface.go b/cmd/object-api-interface.go index 68c9e230f..681024c00 100644 --- a/cmd/object-api-interface.go +++ b/cmd/object-api-interface.go @@ -99,10 +99,13 @@ type ObjectOptions struct { ReplicationSourceTaggingTimestamp time.Time // set if MinIOSourceTaggingTimestamp received ReplicationSourceLegalholdTimestamp time.Time // set if MinIOSourceObjectLegalholdTimestamp received ReplicationSourceRetentionTimestamp time.Time // set if MinIOSourceObjectRetentionTimestamp received - ReplicaLockReconcile bool // set for a trusted SSE-C replica full write/completion: re-order Object Lock against the destination version read under the write lock (single erasure set; see pgsty/silo#133) + ReplicaLockReconcile bool // re-order a trusted replica write against the destination version under the object write lock DeletePrefix bool // set true to enforce a prefix deletion, only application for DeleteObject API, DeletePrefixObject bool // set true when object's erasure set is resolvable by object name (using getHashedSetIndex) + // Pools-layer version resolver; the caller holds the shared object lock. + replicaObjectInfo GetObjectInfoFn + Speedtest bool // object call specifically meant for SpeedTest code, set to 'true' when invoked by SpeedtestHandler. // Use the maximum parity (N/2), used when saving server configuration files diff --git a/cmd/object-handlers.go b/cmd/object-handlers.go index a4dfab00b..86321943b 100644 --- a/cmd/object-handlers.go +++ b/cmd/object-handlers.go @@ -1832,6 +1832,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re applyReplicatedObjectLock(srcInfo.UserDefined, storedLock, replicaTrusted, retentionMode, retentionDate, legalHold, dstOpts.ReplicationSourceRetentionTimestamp, dstOpts.ReplicationSourceLegalholdTimestamp) + dstOpts.ReplicaLockReconcile = replicaTrusted && dstOpts.VersionID != "" if replicaTrusted { // An SSE-C key rotation snapshots every stored reserved key into @@ -2414,10 +2415,8 @@ func (api objectAPIHandlers) PutObjectHandler(w http.ResponseWriter, r *http.Req // The decision above orders against the version as read here; let the // object layer re-run it against the version read under the write lock // that guards the replacement, so a newer hold or retention committed in - // between is not rolled back (issue #120). Scoped to the SSE-C replica - // retransmit this issue enables, keyed on the incoming write's restored - // SSE-C seal, the same predicate as the duplicate-version exemption. - opts.ReplicaLockReconcile = isReplicaTrusted(ctx) && opts.VersionID != "" && crypto.SSEC.IsEncrypted(metadata) + // between is not rolled back, across all pools and encryption modes. + opts.ReplicaLockReconcile = isReplicaTrusted(ctx) && opts.VersionID != "" } if s3Err != ErrNone { writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL)