fix(storage): serialize and reconcile multi-pool object updates

Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-11 20:01:04 +08:00
parent b32f2d9dd0
commit e59a3d938e
10 changed files with 1052 additions and 86 deletions
+5
View File
@@ -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) {
+10 -12
View File
@@ -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
+14 -10
View File
@@ -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)
}
}
+317
View File
@@ -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 <http://www.gnu.org/licenses/>.
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
}
+575
View File
@@ -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 <http://www.gnu.org/licenses/>.
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)
}
}
+5 -1
View File
@@ -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
+111 -56
View File
@@ -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
}
+8 -2
View File
@@ -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,
+4 -1
View File
@@ -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
+3 -4
View File
@@ -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)