diff --git a/cmd/admin-bucket-handlers.go b/cmd/admin-bucket-handlers.go index 01984e4de..a33e18f2f 100644 --- a/cmd/admin-bucket-handlers.go +++ b/cmd/admin-bucket-handlers.go @@ -76,7 +76,7 @@ func (a adminAPIHandlers) PutBucketQuotaConfigHandler(w http.ResponseWriter, r * return } - quotaConfig, err := parseBucketQuota(bucket, data) + _, err = parseBucketQuota(bucket, data) if err != nil { writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) return @@ -94,9 +94,6 @@ func (a adminAPIHandlers) PutBucketQuotaConfigHandler(w http.ResponseWriter, r * Quota: data, UpdatedAt: updatedAt, } - if quotaConfig.Size == 0 && quotaConfig.Quota == 0 { - bucketMeta.Quota = nil - } // Call site replication hook. replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, bucketMeta)) @@ -922,7 +919,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r * continue } - configData, err := json.Marshal(bucketPolicy) + configData, err := canonicalBucketPolicy(bucketPolicy) if err != nil { rpt.SetStatus(bucket, fileName, err) continue @@ -1081,18 +1078,39 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r * continue } var merged BucketMetadata + var commitAt time.Time err := func() error { lockCtx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket) if err != nil { return err } defer unlock() - merged, err = loadBucketMetadataParse(lockCtx, objectAPI, bucket, true) + merged, err = loadBucketMetadataParse(lockCtx, objectAPI, bucket, false) if err != nil { return err } + if err := ensureBucketMetadataCreated(lockCtx, objectAPI, &merged); err != nil { + return err + } + commitAt = UTCNow() + for _, file := range replicatedBucketConfigs { + if _, ok := fields[file]; ok { + commitAt = localBucketConfigUpdatedAt(merged, file, commitAt) + } + } applyImportedBucketMetadata(&merged, *meta, fields) - return globalBucketMetadataSys.saveMetadata(lockCtx, objectAPI, merged) + for _, file := range replicatedBucketConfigs { + if _, ok := fields[file]; !ok { + continue + } + data, at := replicatedBucketConfig(&merged, file) + payload, _, err := bucketConfigPayload(bucket, file, *data, len(merged.ObjectLockConfigXML) != 0) + if err != nil { + return err + } + *data, *at = payload, commitAt + } + return globalBucketMetadataSys.saveMetadata(lockCtx, objectAPI, &merged) }() if err != nil { rpt.SetStatus(bucket, "", err) @@ -1100,7 +1118,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r * } *meta = merged globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket) - hook := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: updatedAt} + hook := madmin.SRBucketMeta{Bucket: bucket, UpdatedAt: commitAt} var hookNeeded bool if _, ok := fields[bucketQuotaConfigFile]; ok { hook.Quota = meta.QuotaConfigJSON @@ -1108,7 +1126,7 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r * } if _, ok := fields[bucketPolicyConfig]; ok { hook.Policy = meta.PolicyConfigJSON - hookNeeded = true + hookNeeded = hookNeeded || len(hook.Policy) != 0 } if _, ok := fields[bucketVersioningConfig]; ok { hook.Versioning = enc(meta.VersioningConfigXML) @@ -1129,6 +1147,12 @@ func (a adminAPIHandlers) ImportBucketMetadataHandler(w http.ResponseWriter, r * if hookNeeded { err = globalSiteReplicationSys.BucketMetaHook(ctx, hook) } + if _, ok := fields[bucketPolicyConfig]; ok && len(meta.PolicyConfigJSON) == 0 { + // An omitted bulk Policy cannot express deletion. + err = errors.Join(err, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{ + Type: madmin.SRBucketMetaTypePolicy, Bucket: bucket, UpdatedAt: commitAt, + })) + } if _, ok := fields[bucketCorsConfig]; ok { // CORS carries its own timestamp, so it replicates through the // dedicated event rather than the shared bucket metadata hook. It diff --git a/cmd/bucket-metadata-replication.go b/cmd/bucket-metadata-replication.go new file mode 100644 index 000000000..087230e66 --- /dev/null +++ b/cmd/bucket-metadata-replication.go @@ -0,0 +1,318 @@ +// Copyright (c) 2015-2026 MinIO, Inc. +// +// This file is part of MinIO 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" + "encoding/json" + "errors" + "sort" + "time" + + "github.com/minio/minio-go/v7/pkg/tags" + bucketsse "github.com/minio/minio/internal/bucket/encryption" + objectlock "github.com/minio/minio/internal/bucket/object/lock" + "github.com/minio/minio/internal/bucket/versioning" + "github.com/pgsty/silo-pkg/v3/policy" +) + +// Only these fields share the site-replication source-time ordering contract. +// Object Lock is applied before Versioning, whose effective document depends on it. +var replicatedBucketConfigs = [...]string{ + objectLockConfig, bucketVersioningConfig, bucketPolicyConfig, + bucketTaggingConfig, bucketSSEConfig, bucketQuotaConfigFile, +} + +func replicatedBucketConfig(meta *BucketMetadata, file string) (*[]byte, *time.Time) { + switch file { + case bucketPolicyConfig: + return &meta.PolicyConfigJSON, &meta.PolicyConfigUpdatedAt + case bucketTaggingConfig: + return &meta.TaggingConfigXML, &meta.TaggingConfigUpdatedAt + case bucketSSEConfig: + return &meta.EncryptionConfigXML, &meta.EncryptionConfigUpdatedAt + case bucketQuotaConfigFile: + return &meta.QuotaConfigJSON, &meta.QuotaConfigUpdatedAt + case bucketVersioningConfig: + return &meta.VersioningConfigXML, &meta.VersioningConfigUpdatedAt + case objectLockConfig: + return &meta.ObjectLockConfigXML, &meta.ObjectLockConfigUpdatedAt + } + return nil, nil +} + +func bucketConfigUpdateOnly(file string) bool { + return file == bucketVersioningConfig || file == objectLockConfig +} + +// Reuse the persistence rule before comparison, so an accepted Versioning +// event and the document Save actually writes have the same comparison key. +func effectiveBucketVersioning(data []byte, lockEnabled bool) []byte { + if lockEnabled { + config, err := versioning.ParseConfig(bytes.NewReader(data)) + if err != nil || !config.Enabled() || config.PrefixesExcluded() { + return enabledBucketVersioningConfig + } + } + return data +} + +// A parsed policy still contains map-backed sets with nondeterministic Marshal +// order. Sort every set array recursively, including statements and conditions. +// RawMessage keeps integer values intact; decoding through float64 would not. +func canonicalBucketPolicyJSON(data json.RawMessage) (json.RawMessage, error) { + data = bytes.TrimSpace(data) + if len(data) == 0 { + return nil, nil + } + switch data[0] { + case '{': + var obj map[string]json.RawMessage + if err := json.Unmarshal(data, &obj); err != nil { + return nil, err + } + for key, value := range obj { + var err error + obj[key], err = canonicalBucketPolicyJSON(value) + if err != nil { + return nil, err + } + } + return json.Marshal(obj) + case '[': + var arr []json.RawMessage + if err := json.Unmarshal(data, &arr); err != nil { + return nil, err + } + for i := range arr { + var err error + arr[i], err = canonicalBucketPolicyJSON(arr[i]) + if err != nil { + return nil, err + } + } + sort.Slice(arr, func(i, j int) bool { return bytes.Compare(arr[i], arr[j]) < 0 }) + return json.Marshal(arr) + default: + var compact bytes.Buffer + if err := json.Compact(&compact, data); err != nil { + return nil, err + } + return compact.Bytes(), nil + } +} + +// Encode the validated policy fields explicitly: BPStatement's required +// Action/Resource tags otherwise try to marshal empty sets for the supported +// NotAction/NotResource alternatives. This stays within the existing schema. +func canonicalBucketPolicy(cfg *policy.BucketPolicy) ([]byte, error) { + if cfg.IsEmpty() { + return nil, nil + } + doc := map[string]any{"Version": cfg.Version} + if cfg.ID != "" { + doc["ID"] = cfg.ID + } + statements := make([]map[string]any, 0, len(cfg.Statements)) + for _, st := range cfg.Statements { + statement := map[string]any{"Effect": st.Effect, "Principal": st.Principal} + if st.SID != "" { + statement["Sid"] = st.SID + } + if len(st.Actions) != 0 { + statement["Action"] = st.Actions + } + if len(st.NotActions) != 0 { + statement["NotAction"] = st.NotActions + } + if len(st.Resources) != 0 { + statement["Resource"] = st.Resources + } + if len(st.NotResources) != 0 { + statement["NotResource"] = st.NotResources + } + if len(st.Conditions) != 0 { + statement["Condition"] = st.Conditions + } + statements = append(statements, statement) + } + doc["Statement"] = statements + data, err := json.Marshal(doc) + if err != nil { + return nil, err + } + return canonicalBucketPolicyJSON(data) +} + +// Validate with the same parsers as Save. Policy's established empty-policy +// semantics are deletion; a parsed zero quota is still a live document. +func bucketConfigPayload(bucket, file string, data []byte, lockEnabled bool) ([]byte, []byte, error) { + if len(data) == 0 { + return nil, nil, nil + } + var err error + key := data + switch file { + case bucketPolicyConfig: + var cfg *policy.BucketPolicy + cfg, err = policy.ParseBucketPolicyConfig(bytes.NewReader(data), bucket) + if err == nil { + if cfg.IsEmpty() { + return nil, nil, nil + } + key, err = canonicalBucketPolicy(cfg) + } + case bucketQuotaConfigFile: + cfg, parseErr := parseBucketQuota(bucket, data) + err = parseErr + if err == nil { + key, err = json.Marshal(cfg) + } + case bucketTaggingConfig: + _, err = tags.ParseBucketXML(bytes.NewReader(data)) + case bucketSSEConfig: + _, err = bucketsse.ParseBucketSSEConfig(bytes.NewReader(data)) + case objectLockConfig: + _, err = objectlock.ParseObjectLockConfig(bytes.NewReader(data)) + case bucketVersioningConfig: + data = effectiveBucketVersioning(data, lockEnabled) + key = data + _, err = versioning.ParseConfig(bytes.NewReader(data)) + } + return data, key, err +} + +type bucketConfigState struct { + data, key []byte + at time.Time + real, valid bool +} + +func newBucketConfigState(bucket, file string, data []byte, at, created time.Time, lockEnabled bool) (bucketConfigState, error) { + data, key, err := bucketConfigPayload(bucket, file, data, lockEnabled) + if err != nil { + return bucketConfigState{}, err + } + if at.IsZero() { + at = created + } + valid := !created.IsZero() && !at.Before(created) + real := valid && at.After(created) + if bucketConfigUpdateOnly(file) && len(data) == 0 { + real = false + } + return bucketConfigState{data: data, key: key, at: at, real: real, valid: valid}, nil +} + +func (s bucketConfigState) candidate() bool { + return s.valid && (s.real || len(s.data) != 0) +} + +func compareBucketConfigStates(a, b bucketConfigState) int { + if a.valid != b.valid { + if a.valid { + return 1 + } + return -1 + } + if a.real != b.real { + if a.real { + return 1 + } + return -1 + } + if a.real { + if n := a.at.Compare(b.at); n != 0 { + return n + } + // At equal source time a real deletion wins, preventing resurrection. + if (len(a.data) == 0) != (len(b.data) == 0) { + if len(a.data) == 0 { + return 1 + } + return -1 + } + } + return bytes.Compare(a.key, b.key) +} + +func localBucketConfigUpdatedAt(meta BucketMetadata, file string, now time.Time) time.Time { + _, at := replicatedBucketConfig(&meta, file) + for _, lower := range []time.Time{meta.Created, *at} { + if !now.After(lower) { + now = lower.Add(time.Nanosecond) + } + } + return now.UTC() +} + +func ensureBucketMetadataCreated(ctx context.Context, obj ObjectLayer, meta *BucketMetadata) error { + if !meta.Created.IsZero() { + return nil + } + info, err := obj.GetBucketInfo(ctx, meta.Name, BucketOptions{NoMetadata: true}) + if err != nil { + return err + } + if info.Created.IsZero() { + return errors.New("bucket metadata creation time is unknown") + } + meta.Created = info.Created.UTC() + return nil +} + +// applyBucketConfig runs under metadata.lock, on freshly loaded metadata. It +// changes only the selected field; the caller persists once after all checks. +func applyBucketConfig(meta *BucketMetadata, file string, data []byte, at time.Time) (bool, error) { + if bucketConfigUpdateOnly(file) && len(data) == 0 { + return false, nil + } + current, currentAt := replicatedBucketConfig(meta, file) + lockEnabled := len(meta.ObjectLockConfigXML) != 0 + incoming, err := newBucketConfigState(meta.Name, file, data, at, meta.Created, lockEnabled) + if err != nil { + return false, err + } + if !incoming.candidate() { + return false, nil + } + local, err := newBucketConfigState(meta.Name, file, *current, *currentAt, meta.Created, lockEnabled) + if err != nil { + return false, err + } + if compareBucketConfigStates(incoming, local) <= 0 { + return false, nil + } + *current, *currentAt = bytes.Clone(incoming.data), incoming.at.UTC() + return true, nil +} + +func rebaseBucketConfigDefaults(meta *BucketMetadata, oldCreated time.Time) { + if meta.Created.Equal(oldCreated) { + return + } + // These six fields alone use Created to distinguish a baseline from a + // tombstone. Preserve actual source times when adopting an existing bucket. + for _, file := range replicatedBucketConfigs { + _, at := replicatedBucketConfig(meta, file) + if at.IsZero() || at.Equal(oldCreated) { + *at = meta.Created + } + } +} diff --git a/cmd/bucket-metadata-sys.go b/cmd/bucket-metadata-sys.go index 8096e717d..81065df34 100644 --- a/cmd/bucket-metadata-sys.go +++ b/cmd/bucket-metadata-sys.go @@ -126,19 +126,38 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) { } } -func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool) (updatedAt time.Time, err error) { +type bucketMetadataUpdate struct { + meta BucketMetadata + updatedAt time.Time + changed bool +} + +func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket, configFile string, configData []byte, parse, lifecycleDelete bool) (time.Time, error) { + result, err := sys.updateAndParseMetadata(ctx, bucket, configFile, configData, parse, lifecycleDelete, nil) + return result.updatedAt, err +} + +func (sys *BucketMetadataSys) updateAndParseMetadata(ctx context.Context, bucket string, configFile string, configData []byte, parse, lifecycleDelete bool, sourceTime *time.Time) (result bucketMetadataUpdate, err error) { objAPI := newObjectLayerFn() if objAPI == nil { - return updatedAt, errServerNotInitialized + return result, errServerNotInitialized } if isMinioMetaBucketName(bucket) { - return updatedAt, errInvalidArgument + return result, errInvalidArgument + } + // Load deletions without parsed caches (notably quota), and compare the + // six replicated fields against the raw document under the same lock. + if data, _ := replicatedBucketConfig(&result.meta, configFile); data != nil { + parse = false + if bucketConfigUpdateOnly(configFile) && len(configData) == 0 { + return result, nil + } } notifyCtx := ctx ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket) if err != nil { - return updatedAt, err + return result, err } err = func() error { @@ -158,55 +177,61 @@ func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, return err } } - updatedAt = UTCNow() - switch configFile { - case bucketPolicyConfig: - meta.PolicyConfigJSON = configData - meta.PolicyConfigUpdatedAt = updatedAt - case bucketNotificationConfig: - meta.NotificationConfigXML = configData - meta.NotificationConfigUpdatedAt = updatedAt - case bucketLifecycleConfig: - meta.LifecycleConfigXML = configData - meta.LifecycleConfigUpdatedAt = updatedAt - case bucketSSEConfig: - meta.EncryptionConfigXML = configData - meta.EncryptionConfigUpdatedAt = updatedAt - case bucketTaggingConfig: - meta.TaggingConfigXML = configData - meta.TaggingConfigUpdatedAt = updatedAt - case bucketQuotaConfigFile: - meta.QuotaConfigJSON = configData - meta.QuotaConfigUpdatedAt = updatedAt - case objectLockConfig: - meta.ObjectLockConfigXML = configData - meta.ObjectLockConfigUpdatedAt = updatedAt - case bucketVersioningConfig: - meta.VersioningConfigXML = configData - meta.VersioningConfigUpdatedAt = updatedAt - case bucketReplicationConfig: - meta.ReplicationConfigXML = configData - meta.ReplicationConfigUpdatedAt = updatedAt - case bucketTargetsFile: - meta.BucketTargetsConfigJSON, meta.BucketTargetsConfigMetaJSON, err = encryptBucketMetadata(ctx, meta.Name, configData, kms.Context{ - bucket: meta.Name, - bucketTargetsFile: bucketTargetsFile, - }) - if err != nil { - return fmt.Errorf("Error encrypting bucket target metadata %w", err) + updatedAt := UTCNow() + if data, _ := replicatedBucketConfig(&meta, configFile); data != nil { + if err := ensureBucketMetadataCreated(ctx, objAPI, &meta); err != nil { + return err + } + if sourceTime == nil || sourceTime.IsZero() { + updatedAt = localBucketConfigUpdatedAt(meta, configFile, updatedAt) + } else { + updatedAt = sourceTime.UTC() + } + changed, err := applyBucketConfig(&meta, configFile, configData, updatedAt) + if err != nil { + return err + } + if !changed { + return nil + } + } else { + switch configFile { + case bucketNotificationConfig: + meta.NotificationConfigXML = configData + meta.NotificationConfigUpdatedAt = updatedAt + case bucketLifecycleConfig: + meta.LifecycleConfigXML = configData + meta.LifecycleConfigUpdatedAt = updatedAt + case bucketReplicationConfig: + meta.ReplicationConfigXML = configData + meta.ReplicationConfigUpdatedAt = updatedAt + case bucketTargetsFile: + meta.BucketTargetsConfigJSON, meta.BucketTargetsConfigMetaJSON, err = encryptBucketMetadata(ctx, meta.Name, configData, kms.Context{ + bucket: meta.Name, + bucketTargetsFile: bucketTargetsFile, + }) + if err != nil { + return fmt.Errorf("Error encrypting bucket target metadata %w", err) + } + meta.BucketTargetsConfigUpdatedAt = updatedAt + meta.BucketTargetsConfigMetaUpdatedAt = updatedAt + default: + return fmt.Errorf("Unknown bucket %s metadata update requested %s", bucket, configFile) } - meta.BucketTargetsConfigUpdatedAt = updatedAt - meta.BucketTargetsConfigMetaUpdatedAt = updatedAt - default: - return fmt.Errorf("Unknown bucket %s metadata update requested %s", bucket, configFile) } - return sys.saveMetadata(ctx, objAPI, meta) + if err := sys.saveMetadata(ctx, objAPI, &meta); err != nil { + return err + } + result = bucketMetadataUpdate{meta: meta, updatedAt: updatedAt, changed: true} + return nil }() if err != nil { - return updatedAt, err + return result, err } - globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) // Do not use caller context here - return updatedAt, nil + if result.changed { + globalNotificationSys.LoadBucketMetadata(bgContext(notifyCtx), bucket) + } + return result, nil } func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) error { @@ -219,7 +244,7 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err return errInvalidArgument } - if err := sys.saveMetadata(ctx, objAPI, meta); err != nil { + if err := sys.saveMetadata(ctx, objAPI, &meta); err != nil { return err } @@ -229,7 +254,7 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err // saveMetadata persists and publishes metadata locally. Callers performing a // read-modify-write must hold metadata.lock and release it before peer fan-out. -func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLayer, meta BucketMetadata) error { +func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLayer, meta *BucketMetadata) error { // A writer may have queued for metadata.lock before DeleteBucket completed. // Recheck the physical bucket under that lock, before recreating metadata. if _, err := objAPI.GetBucketInfo(ctx, meta.Name, BucketOptions{NoMetadata: true}); err != nil { @@ -238,7 +263,7 @@ func (sys *BucketMetadataSys) saveMetadata(ctx context.Context, objAPI ObjectLay if err := meta.Save(ctx, objAPI); err != nil { return err } - sys.Set(meta.Name, meta) + sys.Set(meta.Name, *meta) return nil } @@ -351,7 +376,7 @@ func (sys *BucketMetadataSys) UpdateExpiryLCConfig(ctx context.Context, bucket s } meta.LifecycleConfigXML = configData meta.LifecycleConfigUpdatedAt = UTCNow() - return sys.saveMetadata(ctx, objAPI, meta) + return sys.saveMetadata(ctx, objAPI, &meta) }() if err != nil { return err diff --git a/cmd/bucket-metadata.go b/cmd/bucket-metadata.go index 0bb59bb9e..8bf5cbe82 100644 --- a/cmd/bucket-metadata.go +++ b/cmd/bucket-metadata.go @@ -388,15 +388,7 @@ func (b *BucketMetadata) parseAllConfigs(ctx context.Context, objectAPI ObjectLa } else { b.objectLockConfig = nil } - if b.objectLockConfig != nil { - // Object Lock requires every object to be versioned. Whatever the lock - // document contains, a suspended or prefix-excluded versioning document - // is replaced by plain Enabled versioning; Save persists the result. - config, versioningErr := versioning.ParseConfig(bytes.NewReader(b.VersioningConfigXML)) - if versioningErr != nil || !config.Enabled() || config.PrefixesExcluded() { - b.VersioningConfigXML = enabledBucketVersioningConfig - } - } + b.VersioningConfigXML = effectiveBucketVersioning(b.VersioningConfigXML, b.objectLockConfig != nil) if len(b.VersioningConfigXML) != 0 { b.versioningConfig, err = versioning.ParseConfig(bytes.NewReader(b.VersioningConfigXML)) diff --git a/cmd/bucket-policy-handlers.go b/cmd/bucket-policy-handlers.go index 20d5afe61..585a094db 100644 --- a/cmd/bucket-policy-handlers.go +++ b/cmd/bucket-policy-handlers.go @@ -100,13 +100,13 @@ func (api objectAPIHandlers) PutBucketPolicyHandler(w http.ResponseWriter, r *ht return } - configData, err := json.Marshal(bucketPolicy) + configData, err := canonicalBucketPolicy(bucketPolicy) if err != nil { writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) return } - updatedAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, configData) + result, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketPolicyConfig, configData, false, false, nil) if err != nil { writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) return @@ -116,8 +116,8 @@ func (api objectAPIHandlers) PutBucketPolicyHandler(w http.ResponseWriter, r *ht replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{ Type: madmin.SRBucketMetaTypePolicy, Bucket: bucket, - Policy: bucketPolicyBytes, - UpdatedAt: updatedAt, + Policy: result.meta.PolicyConfigJSON, + UpdatedAt: result.updatedAt, })) // Success. diff --git a/cmd/bucket-versioning-handler.go b/cmd/bucket-versioning-handler.go index bd09c2367..cd4399d50 100644 --- a/cmd/bucket-versioning-handler.go +++ b/cmd/bucket-versioning-handler.go @@ -97,7 +97,7 @@ func (api objectAPIHandlers) PutBucketVersioningHandler(w http.ResponseWriter, r return } - updatedAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, configData) + result, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketVersioningConfig, configData, false, false, nil) if err != nil { writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) return @@ -107,12 +107,12 @@ func (api objectAPIHandlers) PutBucketVersioningHandler(w http.ResponseWriter, r // // We encode the xml bytes as base64 to ensure there are no encoding // errors. - cfgStr := base64.StdEncoding.EncodeToString(configData) + cfgStr := base64.StdEncoding.EncodeToString(result.meta.VersioningConfigXML) replLogIf(ctx, globalSiteReplicationSys.BucketMetaHook(ctx, madmin.SRBucketMeta{ Type: madmin.SRBucketMetaTypeVersionConfig, Bucket: bucket, Versioning: &cfgStr, - UpdatedAt: updatedAt, + UpdatedAt: result.updatedAt, })) writeSuccessResponseHeadersOnly(w) diff --git a/cmd/site-replication-metadata_test.go b/cmd/site-replication-metadata_test.go new file mode 100644 index 000000000..45e538f00 --- /dev/null +++ b/cmd/site-replication-metadata_test.go @@ -0,0 +1,624 @@ +// Copyright (c) 2015-2026 MinIO, Inc. +// +// This file is part of MinIO 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" + "encoding/base64" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "path" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/minio/madmin-go/v3" + "github.com/minio/minio/internal/auth" +) + +func TestPeerBucketMetadataSourceTimeAndDeletion(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: testPeerBucketMetadataSourceTimeAndDeletion}) +} + +func testPeerBucketMetadataSourceTimeAndDeletion(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) { + ctx := t.Context() + created := UTCNow().Add(-time.Hour) + putAt := created.Add(10 * time.Minute) + delAt := putAt.Add(time.Minute) + enc := func(s string) *string { v := base64.StdEncoding.EncodeToString([]byte(s)); return &v } + base := newBucketMetadata(bucket) + base.SetCreatedAt(created) + base.defaultTimestamps() + quotaJSON, err := json.Marshal(madmin.BucketQuota{Quota: 1024, Type: madmin.HardQuota}) + if err != nil { + t.Fatal(err) + } + cases := []struct { + name, file string + put madmin.SRBucketMeta + value func(BucketMetadata) []byte + stamp func(BucketMetadata) time.Time + exported func(madmin.SRBucketInfo) time.Time + deletable bool + }{ + {"policy", bucketPolicyConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypePolicy, Policy: []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket))}, func(m BucketMetadata) []byte { return m.PolicyConfigJSON }, func(m BucketMetadata) time.Time { return m.PolicyConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.PolicyUpdatedAt }, true}, + {"tags", bucketTaggingConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeTags, Tags: enc(`keyold`)}, func(m BucketMetadata) []byte { return m.TaggingConfigXML }, func(m BucketMetadata) time.Time { return m.TaggingConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.TagConfigUpdatedAt }, true}, + {"sse", bucketSSEConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeSSEConfig, SSEConfig: enc(`AES256`)}, func(m BucketMetadata) []byte { return m.EncryptionConfigXML }, func(m BucketMetadata) time.Time { return m.EncryptionConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.SSEConfigUpdatedAt }, true}, + {"quota", bucketQuotaConfigFile, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeQuotaConfig, Quota: quotaJSON}, func(m BucketMetadata) []byte { return m.QuotaConfigJSON }, func(m BucketMetadata) time.Time { return m.QuotaConfigUpdatedAt }, func(m madmin.SRBucketInfo) time.Time { return m.QuotaConfigUpdatedAt }, true}, + {"versioning", bucketVersioningConfig, madmin.SRBucketMeta{Type: madmin.SRBucketMetaTypeVersionConfig, Versioning: enc(`Enabled`)}, func(m BucketMetadata) []byte { return m.VersioningConfigXML }, func(m BucketMetadata) time.Time { return m.VersioningConfigUpdatedAt }, nil, false}, + {"objectlock", objectLockConfig, newSRBucketObjectLockMeta(bucket, enc(`EnabledGOVERNANCE30`), putAt), func(m BucketMetadata) []byte { return m.ObjectLockConfigXML }, func(m BucketMetadata) time.Time { return m.ObjectLockConfigUpdatedAt }, nil, false}, + } + for _, tc := range cases { + t.Run(backend+"/"+tc.name, func(t *testing.T) { + if err := globalBucketMetadataSys.save(ctx, base); err != nil { + t.Fatal(err) + } + item := tc.put + item.Bucket, item.UpdatedAt = bucket, putAt + apply := func(item madmin.SRBucketMeta) { + t.Helper() + rec := applySRBucketMetaViaAdmin(t, cred, item) + if rec.Code != http.StatusOK { + t.Fatalf("admin apply returned %d: %s", rec.Code, rec.Body.String()) + } + } + read := func() BucketMetadata { + t.Helper() + m, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + return m + } + apply(item) + put := read() + if len(tc.value(put)) == 0 { + t.Fatal("PUT did not establish live config") + } + if !tc.stamp(put).Equal(putAt) { + t.Errorf("SOURCE_TIME: persisted %s, want source %s", tc.stamp(put), putAt) + } + apply(madmin.SRBucketMeta{Type: item.Type, Bucket: bucket, UpdatedAt: delAt}) + after := read() + if !tc.deletable { + if !bytes.Equal(tc.value(put), tc.value(after)) || !tc.stamp(put).Equal(tc.stamp(after)) { + t.Error("NIL_NOOP: update-only config changed") + } + t.Log("nil payload is correctly a no-op") + return + } + if len(tc.value(after)) != 0 { + t.Error("NEWER_DELETE: HTTP 200, but live config remained after newer source DELETE") + } + if _, err := globalBucketMetadataSys.Delete(ctx, bucket, tc.file); err != nil { + t.Fatal(err) + } + tombstone := read() + if len(tc.value(tombstone)) != 0 { + t.Fatal("local DELETE failed to establish tombstone") + } + apply(item) + after = read() + if len(tc.value(after)) != 0 { + t.Error("STALE_RESURRECTION: older source PUT resurrected a locally deleted config") + } else { + t.Log("older source PUT did not resurrect config") + } + }) + } +} + +func TestPeerBucketMetadataBulkOrdering(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) { + t.Run(backend, func(t *testing.T) { + ctx := t.Context() + meta, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + newAt := meta.Created.Add(time.Hour) + newXML := `keynew` + meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = []byte(newXML), newAt + if err := globalBucketMetadataSys.save(ctx, meta); err != nil { + t.Fatal(err) + } + oldXML := base64.StdEncoding.EncodeToString([]byte(`keyold`)) + item := madmin.SRBucketMeta{Bucket: bucket, Tags: &oldXML, UpdatedAt: newAt.Add(-time.Minute)} + rec := applySRBucketMetaViaAdmin(t, cred, item) + if rec.Code != http.StatusOK { + t.Fatalf("bulk apply returned %d: %s", rec.Code, rec.Body.String()) + } + after, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(after.TaggingConfigXML, []byte(newXML)) || !after.TaggingConfigUpdatedAt.Equal(newAt) { + t.Errorf("BULK_STALE_OVERWRITE: older bulk event overwrote newer config, value=%q time=%s", after.TaggingConfigXML, after.TaggingConfigUpdatedAt) + } + }) + }}) +} + +func TestPeerBucketMetadataOrderingUnderLock(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { + t.Run(backend, func(t *testing.T) { + ctx := t.Context() + meta, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + oldAt, newAt := meta.Created.Add(time.Hour), meta.Created.Add(2*time.Hour) + oldXML := base64.StdEncoding.EncodeToString([]byte(`keyolder`)) + newXML := []byte(`keynewest`) + lockCtx, unlock, err := lockBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + locked := true + defer func() { + if locked { + unlock() + } + }() + ready := make(chan struct{}, 1) + hook := func(name string) { + if name == bucket { + select { + case ready <- struct{}{}: + default: + } + } + } + lockBucketMetadataAcquireHook.Store(&hook) + defer lockBucketMetadataAcquireHook.Store(nil) + done := make(chan error, 1) + go func() { done <- globalSiteReplicationSys.PeerBucketTaggingHandler(ctx, bucket, &oldXML, oldAt) }() + select { + case <-ready: + case <-time.After(5 * time.Second): + t.Fatal("older event never reached metadata lock") + } + // A newer writer commits while holding the existing metadata lock. + // The old peer event has already checked the pre-commit cache. + meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = newXML, newAt + if err := globalBucketMetadataSys.saveMetadata(lockCtx, obj, &meta); err != nil { + t.Fatal(err) + } + unlock() + locked = false + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(5 * time.Second): + t.Fatal("older event did not finish") + } + after, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + if string(after.TaggingConfigXML) != string(newXML) { + t.Errorf("CHECK_OUTSIDE_LOCK: queued older event replaced newer committed tags with %q", after.TaggingConfigXML) + } + }) + }}) +} + +// Count actual metadata persistence, including writes that would be invisible +// in a value-only assertion after a duplicate or rejected event. +type bucketConfigWriteCounter struct { + ObjectLayer + writes atomic.Int64 +} + +func (o *bucketConfigWriteCounter) PutObject(ctx context.Context, bucket, object string, data *PutObjReader, opts ObjectOptions) (ObjectInfo, error) { + if bucket == minioMetaBucket && path.Base(object) == bucketMetadataFile { + o.writes.Add(1) + } + return o.ObjectLayer.PutObject(ctx, bucket, object, data, opts) +} + +func TestBucketConfigStateOrdering(t *testing.T) { + created := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + live := []byte(`ka`) + other := bytes.ReplaceAll(live, []byte("a"), []byte("b")) + state := func(data []byte, at, birth time.Time) bucketConfigState { + t.Helper() + s, err := newBucketConfigState("bucket", bucketTaggingConfig, data, at, birth, false) + if err != nil { + t.Fatal(err) + } + return s + } + baseline := state(nil, created, created) + baselineLive := state(live, created, created) + put := state(live, created.Add(time.Second), created) + deleted := state(nil, put.at, created) + lateBaseline := state(other, created.Add(time.Hour), created.Add(time.Hour)) + tests := []struct { + name string + lower, higher bucketConfigState + }{ + {"baseline initialization", baseline, baselineLive}, + {"real before later baseline", lateBaseline, put}, + {"delete wins tie", put, deleted}, + {"live key tie", put, state(other, put.at, created)}, + {"new PUT after deletion", deleted, state(live, put.at.Add(time.Second), created)}, + {"baseline key ignores creation time", state(live, created.Add(time.Hour), created.Add(time.Hour)), state(other, created, created)}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if compareBucketConfigStates(tc.lower, tc.higher) >= 0 || compareBucketConfigStates(tc.higher, tc.lower) <= 0 { + t.Fatal("ordering is not antisymmetric or chose the wrong winner") + } + }) + } + for _, s := range []bucketConfigState{baseline, state(live, created.Add(-time.Second), created), state(live, created, time.Time{})} { + if s.candidate() { + t.Fatal("default, pre-creation, or unknown-generation state became a source") + } + } + for _, file := range []string{bucketVersioningConfig, objectLockConfig} { + s, err := newBucketConfigState("bucket", file, nil, put.at, created, false) + if err != nil || s.candidate() { + t.Fatalf("empty update-only candidate: %+v %v", s, err) + } + } +} + +func TestBucketPolicyReplicationKey(t *testing.T) { + // Permute independent set arrays and parse separately, as two sites would. + a := []byte(`{"Version":"2012-10-17","Statement":[{"Sid":"one","Effect":"Allow","Principal":{"AWS":["b","a"]},"Action":["s3:GetObject","s3:PutObject"],"Resource":["arn:aws:s3:::bucket/b*","arn:aws:s3:::bucket/a*"],"Condition":{"StringLike":{"s3:prefix":["b*","a*"]}}},{"Sid":"two","Effect":"Deny","Principal":"*","NotAction":["s3:GetObject","s3:PutObject"],"NotResource":["arn:aws:s3:::bucket/d*","arn:aws:s3:::bucket/c*"]}]}`) + // Use a condition valid for object actions. + a = bytes.ReplaceAll(a, []byte("s3:prefix"), []byte("aws:UserAgent")) + b := []byte(`{"Statement":[{"NotResource":["arn:aws:s3:::bucket/c*","arn:aws:s3:::bucket/d*"],"NotAction":["s3:PutObject","s3:GetObject"],"Principal":"*","Effect":"Deny","Sid":"two"},{"Condition":{"StringLike":{"aws:UserAgent":["a*","b*"]}},"Resource":["arn:aws:s3:::bucket/a*","arn:aws:s3:::bucket/b*"],"Action":["s3:PutObject","s3:GetObject"],"Principal":{"AWS":["a","b"]},"Effect":"Allow","Sid":"one"}],"Version":"2012-10-17"}`) + _, key, err := bucketConfigPayload("bucket", bucketPolicyConfig, a, false) + if err != nil { + t.Fatal(err) + } + for i := range 30 { + _, next, err := bucketConfigPayload("bucket", bucketPolicyConfig, b, false) + if err != nil || !bytes.Equal(key, next) { + t.Fatalf("independent parse %d: %s != %s (%v)", i, key, next, err) + } + } + _, sid, err := bucketConfigPayload("bucket", bucketPolicyConfig, bytes.ReplaceAll(a, []byte(`"one"`), []byte(`"other"`)), false) + if err != nil || bytes.Equal(key, sid) { + t.Fatalf("Sid difference lost: %v", err) + } + n, err := canonicalBucketPolicyJSON([]byte(`{"n":[9007199254740993,9007199254740992]}`)) + if err != nil || string(n) != `{"n":[9007199254740992,9007199254740993]}` { + t.Fatalf("integer precision lost: %s %v", n, err) + } +} + +func TestPeerBucketMetadataWireAtomicity(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) { + t.Run(backend, func(t *testing.T) { + ctx := t.Context() + meta, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + stamp := meta.Created.Add(time.Hour) + policyData := []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket)) + meta.PolicyConfigJSON, meta.PolicyConfigUpdatedAt = policyData, stamp + meta.QuotaConfigJSON, meta.QuotaConfigUpdatedAt = []byte(`{"quota":1024,"quotatype":"hard"}`), stamp + meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = []byte(`ka`), stamp + if err = globalBucketMetadataSys.save(ctx, meta); err != nil { + t.Fatal(err) + } + counter := &bucketConfigWriteCounter{ObjectLayer: obj} + setObjectLayer(counter) + defer setObjectLayer(obj) + read := func() BucketMetadata { + t.Helper() + m, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + return m + } + apply := func(fragment string, at time.Time, ok bool) { + t.Helper() + raw := fmt.Sprintf(`{"bucket":%q,"updatedAt":%q,%s}`, bucket, at.Format(time.RFC3339Nano), fragment) + var item madmin.SRBucketMeta + if err := json.Unmarshal([]byte(raw), &item); err != nil { + t.Fatal(err) + } + rec := applySRBucketMetaViaAdmin(t, cred, item) + if (rec.Code == http.StatusOK) != ok { + t.Fatalf("%s: %d %s", raw, rec.Code, rec.Body.String()) + } + } + // Omitted/null *string and empty update-only payloads are all no-ops. + apply(`"policy":null,"quota":null,"tags":null,"versioningConfig":"","objectLockConfig":""`, stamp.Add(time.Minute), true) + current := read() + if len(current.PolicyConfigJSON) != 0 || len(current.QuotaConfigJSON) == 0 { + t.Fatal("explicit null semantics lost") + } + if q, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err != nil || q.Quota != 0 { + t.Fatalf("zero quota cache: %+v %v", q, err) + } + before := counter.writes.Load() + apply(`"tags":null,"versioningConfig":"","objectLockConfig":""`, stamp.Add(2*time.Minute), true) + if counter.writes.Load() != before { + t.Fatal("omitted fields or empty update-only wrote metadata") + } + // An invalid second field cannot persist the valid deletion in this bulk. + apply(`"tags":"","quota":{"quota":1,"quotatype":"invalid"}`, stamp.Add(2*time.Minute), false) + if counter.writes.Load() != before || len(read().TaggingConfigXML) == 0 { + t.Fatal("invalid bulk partially committed") + } + apply(`"tags":""`, stamp.Add(2*time.Minute), true) + before = counter.writes.Load() + apply(`"tags":""`, stamp.Add(2*time.Minute), true) + apply(`"tags":""`, stamp.Add(time.Minute), true) + if counter.writes.Load() != before { + t.Fatal("duplicate or old event wrote metadata") + } + localAt, err := globalBucketMetadataSys.Update(ctx, bucket, bucketQuotaConfigFile, []byte(`{}`)) + if err != nil || !localAt.After(current.QuotaConfigUpdatedAt) { + t.Fatalf("local future monotonic time: %v %v", localAt, err) + } + deletedAt, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketQuotaConfigFile) + if err != nil || !deletedAt.After(localAt) { + t.Fatalf("adjacent local time: %v %v", deletedAt, err) + } + if q, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err != nil || q.Quota != 0 { + t.Fatalf("deleted quota cache: %+v %v", q, err) + } + }) + }}) +} + +func TestPeerBucketAdoptionRebasesOnlyDefaults(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { + for _, shift := range []time.Duration{-time.Hour, 0, time.Hour} { + t.Run(fmt.Sprintf("%s/%s", backend, shift), func(t *testing.T) { + created := UTCNow().Add(-3 * time.Hour) + meta := newBucketMetadata(bucket) + meta.Created = created + meta.defaultTimestamps() + meta.QuotaConfigUpdatedAt = time.Time{} + meta.TaggingConfigUpdatedAt = created.Add(2 * time.Hour) // real deletion + meta.EncryptionConfigXML = []byte(`AES256`) + meta.EncryptionConfigUpdatedAt = created.Add(2 * time.Hour) + if err := globalBucketMetadataSys.save(t.Context(), meta); err != nil { + t.Fatal(err) + } + if err := globalSiteReplicationSys.PeerBucketMakeWithVersioningHandler(t.Context(), bucket, MakeBucketOptions{CreatedAt: created.Add(shift)}); err != nil { + t.Fatal(err) + } + got, err := readBucketMetadata(t.Context(), obj, bucket) + if err != nil { + t.Fatal(err) + } + if !got.PolicyConfigUpdatedAt.Equal(got.Created) || !got.QuotaConfigUpdatedAt.Equal(got.Created) { + t.Fatal("default turned into tombstone or invalid source") + } + if !got.TaggingConfigUpdatedAt.Equal(meta.TaggingConfigUpdatedAt) || !got.EncryptionConfigUpdatedAt.Equal(meta.EncryptionConfigUpdatedAt) || !bytes.Equal(got.EncryptionConfigXML, meta.EncryptionConfigXML) { + t.Fatal("actual state changed during adoption") + } + }) + } + }}) +} + +func recordBucketConfigPeer(t *testing.T, cred auth.Credentials) func() []madmin.SRBucketMeta { + t.Helper() + var mu sync.Mutex + var events []madmin.SRBucketMeta + remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + var event madmin.SRBucketMeta + if err := json.NewDecoder(r.Body).Decode(&event); err != nil { + t.Error(err) + w.WriteHeader(http.StatusBadRequest) + return + } + mu.Lock() + events = append(events, event) + mu.Unlock() + w.WriteHeader(http.StatusOK) + })) + t.Cleanup(remote.Close) + serviceCred, err := auth.CreateCredentials("metadata-source-time-svc", "metadata-source-time-service-secret") + if err != nil { + t.Fatal(err) + } + serviceCred.ParentUser = cred.AccessKey + globalSiteReplicatorCred.Set(serviceCred.SecretKey) + t.Cleanup(func() { globalSiteReplicatorCred.Set("") }) + if _, err = globalIAMSys.store.AddServiceAccount(t.Context(), serviceCred); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { globalIAMSys.DeleteServiceAccount(context.Background(), serviceCred.AccessKey, false) }) + globalSiteReplicationSys.Lock() + oldEnabled, oldState := globalSiteReplicationSys.enabled, globalSiteReplicationSys.state + globalSiteReplicationSys.enabled = true + globalSiteReplicationSys.state = srState{Name: "metadata-source-test", ServiceAccountAccessKey: serviceCred.AccessKey, Peers: map[string]madmin.PeerInfo{ + globalDeploymentID(): {Name: "local", DeploymentID: globalDeploymentID()}, + "metadata-peer": {Name: "remote", DeploymentID: "metadata-peer", Endpoint: remote.URL}, + }} + globalSiteReplicationSys.Unlock() + t.Cleanup(func() { + globalSiteReplicationSys.Lock() + globalSiteReplicationSys.enabled, globalSiteReplicationSys.state = oldEnabled, oldState + globalSiteReplicationSys.Unlock() + }) + return func() []madmin.SRBucketMeta { + mu.Lock() + defer mu.Unlock() + return append([]madmin.SRBucketMeta(nil), events...) + } +} + +// The final metadata lock is taken after reading the ZIP. Commit another +// writer immediately before that acquisition, using the existing lock wrapper. +type importInterleavingObjectLayer struct { + ObjectLayer + once atomic.Bool + write func() +} + +func (o *importInterleavingObjectLayer) NewNSLock(bucket string, objects ...string) RWLocker { + lock := o.ObjectLayer.NewNSLock(bucket, objects...) + if bucket != minioMetaBucket || len(objects) != 1 || path.Base(objects[0]) != "metadata.lock" { + return lock + } + return metadataObservedRWLocker{RWLocker: lock, onLock: func(context.Context) { + if o.once.CompareAndSwap(false, true) { + o.write() + } + }} +} + +func TestLocalBucketMetadataCommittedEvents(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, endpoints: []string{"PutBucketPolicy", "GetBucketPolicy", "PutBucketVersioning"}, objAPITest: func(obj ObjectLayer, backend, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) { + t.Run(backend, func(t *testing.T) { + ctx := t.Context() + events := recordBucketConfigPeer(t, cred) + putPolicy := func(data []byte) { + t.Helper() + req, err := newTestSignedRequestV4(http.MethodPut, getPutPolicyURL("", bucket), int64(len(data)), bytes.NewReader(data), cred.AccessKey, cred.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req) + if rec.Code != http.StatusNoContent { + t.Fatalf("empty policy PUT: %d %s", rec.Code, rec.Body.String()) + } + } + putPolicy([]byte(`{"Version":"2012-10-17","Statement":[]}`)) + meta, err := loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + got := events() + if len(got) != 1 || got[0].Type != madmin.SRBucketMetaTypePolicy || len(got[0].Policy) != 0 || len(meta.PolicyConfigJSON) != 0 || !meta.PolicyConfigUpdatedAt.Equal(got[0].UpdatedAt) { + t.Fatalf("empty policy event/disk mismatch: %+v", got) + } + req, err := newTestSignedRequestV4(http.MethodGet, getGetPolicyURL("", bucket), 0, nil, cred.AccessKey, cred.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req) + if rec.Code != http.StatusNotFound { + t.Fatalf("empty policy GET: %d %s", rec.Code, rec.Body.String()) + } + corsAdminRequest(t, cred, http.MethodPut, "/set-bucket-quota?bucket="+bucket, []byte(`{}`)) + meta, err = loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + got = events() + last := got[len(got)-1] + if last.Type != madmin.SRBucketMetaTypeQuotaConfig || len(last.Quota) == 0 || !last.UpdatedAt.Equal(meta.QuotaConfigUpdatedAt) { + t.Fatalf("zero quota event: %+v", last) + } + + injectedAt := UTCNow().Add(2 * time.Hour) + interleaving := &importInterleavingObjectLayer{ObjectLayer: obj, write: func() { + data := []byte(`{"quota":2048,"quotatype":"hard"}`) + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketQuotaConfigFile, data, false, false, &injectedAt) + if err != nil { + t.Error(err) + } + }} + setObjectLayer(interleaving) + defer setObjectLayer(obj) + before := len(events()) + rec = corsAdminRequest(t, cred, http.MethodPut, "/import-bucket-metadata", corsZip(t, map[string][]byte{ + bucket + "/" + bucketPolicyConfig: []byte(`{"Version":"2012-10-17","Statement":[]}`), + bucket + "/quota.json": []byte(`{}`), + bucket + "/" + bucketTaggingConfig: []byte(`importedyes`), + })) + report := corsImportReport(t, rec).Buckets[bucket] + if report.Err != "" || report.Policy.Err != "" || report.Quota.Err != "" || report.Tagging.Err != "" { + t.Fatalf("import: %+v", report) + } + meta, err = loadBucketMetadata(ctx, obj, bucket) + if err != nil { + t.Fatal(err) + } + if !meta.QuotaConfigUpdatedAt.After(injectedAt) || !meta.TaggingConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) || !meta.PolicyConfigUpdatedAt.Equal(meta.QuotaConfigUpdatedAt) { + t.Fatalf("import time %v must exceed intervening %v; policy=%v tags=%v", meta.QuotaConfigUpdatedAt, injectedAt, meta.PolicyConfigUpdatedAt, meta.TaggingConfigUpdatedAt) + } + got = events()[before:] + if len(got) != 2 { + t.Fatalf("import events: %+v", got) + } + for _, event := range got { + if !event.UpdatedAt.Equal(meta.QuotaConfigUpdatedAt) { + t.Fatal("import hook used ZIP start time") + } + if event.Type == madmin.SRBucketMetaTypePolicy { + if len(event.Policy) != 0 { + t.Fatal("empty policy import is not deletion") + } + } else if event.Tags == nil || len(event.Quota) == 0 { + t.Fatalf("bulk omitted imported state: %+v", event) + } + } + }) + }}) +} + +func TestPeerBucketMetadataNormalizesBeforeComparison(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) { + t.Run(backend, func(t *testing.T) { + meta, err := loadBucketMetadata(t.Context(), obj, bucket) + if err != nil { + t.Fatal(err) + } + versioning := base64.StdEncoding.EncodeToString([]byte(`Enabledskip/true`)) + lock := base64.StdEncoding.EncodeToString(enabledBucketObjectLockConfig) + item := madmin.SRBucketMeta{Bucket: bucket, ObjectLockConfig: &lock, Versioning: &versioning, UpdatedAt: meta.Created.Add(time.Hour)} + counter := &bucketConfigWriteCounter{ObjectLayer: obj} + setObjectLayer(counter) + defer setObjectLayer(obj) + for range 2 { + rec := applySRBucketMetaViaAdmin(t, cred, item) + if rec.Code != http.StatusOK { + t.Fatalf("bulk: %d %s", rec.Code, rec.Body.String()) + } + } + if counter.writes.Load() != 1 { + t.Fatalf("normalized duplicate persisted %d times", counter.writes.Load()) + } + meta, err = loadBucketMetadata(t.Context(), obj, bucket) + if err != nil || !bytes.Equal(meta.VersioningConfigXML, enabledBucketVersioningConfig) || !meta.VersioningConfigUpdatedAt.Equal(item.UpdatedAt) { + t.Fatalf("normalization: %s %v", meta.VersioningConfigXML, err) + } + // Ordinary local updates use the identical effective document in their + // commit result, even if Object Lock changed since the handler's precheck. + data, _ := base64.StdEncoding.DecodeString(versioning) + result, err := globalBucketMetadataSys.updateAndParseMetadata(t.Context(), bucket, bucketVersioningConfig, data, false, false, nil) + if err != nil || !bytes.Equal(result.meta.VersioningConfigXML, enabledBucketVersioningConfig) { + t.Fatalf("local commit snapshot: %s %v", result.meta.VersioningConfigXML, err) + } + }) + }}) +} diff --git a/cmd/site-replication.go b/cmd/site-replication.go index 31e843e6f..fba20be6b 100644 --- a/cmd/site-replication.go +++ b/cmd/site-replication.go @@ -906,7 +906,7 @@ func enablePeerBucketVersioning(meta *BucketMetadata, lockEnabled bool) error { config, err := versioning.ParseConfig(bytes.NewReader(meta.VersioningConfigXML)) if err != nil || (lockEnabled && (config.Suspended() || config.PrefixesExcluded())) { meta.VersioningConfigXML = enabledBucketVersioningConfig - meta.VersioningConfigUpdatedAt = UTCNow() + meta.VersioningConfigUpdatedAt = localBucketConfigUpdatedAt(*meta, bucketVersioningConfig, UTCNow()) return nil } if config.Enabled() { @@ -917,7 +917,7 @@ func enablePeerBucketVersioning(meta *BucketMetadata, lockEnabled bool) error { if err != nil { return err } - meta.VersioningConfigUpdatedAt = UTCNow() + meta.VersioningConfigUpdatedAt = localBucketConfigUpdatedAt(*meta, bucketVersioningConfig, UTCNow()) return nil } @@ -947,7 +947,9 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con if err != nil { return err } + oldCreated := meta.Created meta.SetCreatedAt(opts.CreatedAt) + rebaseBucketConfigDefaults(&meta, oldCreated) if err = enablePeerBucketVersioning(&meta, opts.LockEnabled || len(meta.ObjectLockConfigXML) != 0); err != nil { return err @@ -958,7 +960,7 @@ func (c *SiteReplicationSys) PeerBucketMakeWithVersioningHandler(ctx context.Con meta.ObjectLockConfigUpdatedAt = meta.Created } } - return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, meta) + return globalBucketMetadataSys.saveMetadata(bgContext(ctx), objAPI, &meta) }() if err != nil { return wrapSRErr(c.annotateErr(makeBucketWithVersion, err)) @@ -1583,24 +1585,18 @@ func (c *SiteReplicationSys) BucketMetaHook(ctx context.Context, item madmin.SRB // PeerBucketVersioningHandler - updates versioning config to local cluster. func (c *SiteReplicationSys) PeerBucketVersioningHandler(ctx context.Context, bucket string, versioning *string, updatedAt time.Time) error { + var data []byte if versioning != nil { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetVersioningConfig(bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - configData, err := base64.StdEncoding.DecodeString(*versioning) + var err error + data, err = base64.StdEncoding.DecodeString(*versioning) if err != nil { return wrapSRErr(err) } - _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketVersioningConfig, configData) - if err != nil { - return wrapSRErr(err) - } - return nil } - + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketVersioningConfig, data, false, false, &updatedAt) + if err != nil { + return wrapSRErr(err) + } return nil } @@ -1650,50 +1646,43 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context return nil } - if item.Policy != nil { - meta.PolicyConfigJSON = item.Policy - meta.PolicyConfigUpdatedAt = item.UpdatedAt + // Presence is separate from content: omitted RawMessage is not a delete; + // explicit JSON null still goes through the existing policy/quota parser. + updates := make(map[string][]byte, len(replicatedBucketConfigs)) + if len(item.Policy) != 0 { + updates[bucketPolicyConfig] = item.Policy } - - if item.Versioning != nil { - configData, err := base64.StdEncoding.DecodeString(*item.Versioning) + if len(item.Quota) != 0 { + updates[bucketQuotaConfigFile] = item.Quota + } + for file, payload := range map[string]*string{ + objectLockConfig: item.ObjectLockConfig, bucketVersioningConfig: item.Versioning, + bucketTaggingConfig: item.Tags, bucketSSEConfig: item.SSEConfig, + } { + if payload != nil { + data, err := base64.StdEncoding.DecodeString(*payload) + if err != nil { + return wrapSRErr(err) + } + updates[file] = data + } + } + if len(updates) != 0 { + if err := ensureBucketMetadataCreated(ctx, objectAPI, &meta); err != nil { + return wrapSRErr(err) + } + } + changed := false + for _, file := range replicatedBucketConfigs { + data, supplied := updates[file] + if !supplied { + continue + } + applied, err := applyBucketConfig(&meta, file, data, item.UpdatedAt) if err != nil { return wrapSRErr(err) } - meta.VersioningConfigXML = configData - meta.VersioningConfigUpdatedAt = item.UpdatedAt - } - - if item.Tags != nil { - configData, err := base64.StdEncoding.DecodeString(*item.Tags) - if err != nil { - return wrapSRErr(err) - } - meta.TaggingConfigXML = configData - meta.TaggingConfigUpdatedAt = item.UpdatedAt - } - - if item.ObjectLockConfig != nil { - configData, err := base64.StdEncoding.DecodeString(*item.ObjectLockConfig) - if err != nil { - return wrapSRErr(err) - } - meta.ObjectLockConfigXML = configData - meta.ObjectLockConfigUpdatedAt = item.UpdatedAt - } - - if item.SSEConfig != nil { - configData, err := base64.StdEncoding.DecodeString(*item.SSEConfig) - if err != nil { - return wrapSRErr(err) - } - meta.EncryptionConfigXML = configData - meta.EncryptionConfigUpdatedAt = item.UpdatedAt - } - - if item.Quota != nil { - meta.QuotaConfigJSON = item.Quota - meta.QuotaConfigUpdatedAt = item.UpdatedAt + changed = changed || applied } if item.Cors != nil { @@ -1702,10 +1691,14 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context if compareCORSReplicationStates(localState, incoming) < 0 { meta.CorsConfigXML = bytes.Clone(corsConfigData) meta.CorsConfigUpdatedAt = item.UpdatedAt + changed = true } } - if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil { + if !changed { + return nil + } + if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, &meta); err != nil { return err } unlock() @@ -1716,62 +1709,35 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context // PeerBucketPolicyHandler - copies/deletes policy to local cluster. func (c *SiteReplicationSys) PeerBucketPolicyHandler(ctx context.Context, bucket string, policy *policy.BucketPolicy, updatedAt time.Time) error { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetPolicyConfig(bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - + var data []byte if policy != nil { - configData, err := json.Marshal(policy) + var err error + data, err = canonicalBucketPolicy(policy) if err != nil { return wrapSRErr(err) } - - _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketPolicyConfig, configData) - if err != nil { - return wrapSRErr(err) - } - return nil } - - // Delete the bucket policy - _, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketPolicyConfig) + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketPolicyConfig, data, false, false, &updatedAt) if err != nil { return wrapSRErr(err) } - return nil } // PeerBucketTaggingHandler - copies/deletes tags to local cluster. func (c *SiteReplicationSys) PeerBucketTaggingHandler(ctx context.Context, bucket string, tags *string, updatedAt time.Time) error { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetTaggingConfig(bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - + var data []byte if tags != nil { - configData, err := base64.StdEncoding.DecodeString(*tags) + var err error + data, err = base64.StdEncoding.DecodeString(*tags) if err != nil { return wrapSRErr(err) } - _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketTaggingConfig, configData) - if err != nil { - return wrapSRErr(err) - } - return nil } - - // Delete the tags - _, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketTaggingConfig) + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketTaggingConfig, data, false, false, &updatedAt) if err != nil { return wrapSRErr(err) } - return nil } @@ -1797,51 +1763,32 @@ func (c *SiteReplicationSys) peerBucketObjectLockConfigItem(ctx context.Context, // PeerBucketObjectLockConfigHandler - sets object lock on local bucket. func (c *SiteReplicationSys) PeerBucketObjectLockConfigHandler(ctx context.Context, bucket string, objectLockData *string, updatedAt time.Time) error { + var data []byte if objectLockData != nil { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetObjectLockConfig(bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - - configData, err := base64.StdEncoding.DecodeString(*objectLockData) + var err error + data, err = base64.StdEncoding.DecodeString(*objectLockData) if err != nil { return wrapSRErr(err) } - _, err = globalBucketMetadataSys.Update(ctx, bucket, objectLockConfig, configData) - if err != nil { - return wrapSRErr(err) - } - return nil } - + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, objectLockConfig, data, false, false, &updatedAt) + if err != nil { + return wrapSRErr(err) + } return nil } // PeerBucketSSEConfigHandler - copies/deletes SSE config to local cluster. func (c *SiteReplicationSys) PeerBucketSSEConfigHandler(ctx context.Context, bucket string, sseConfig *string, updatedAt time.Time) error { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetSSEConfig(bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - + var data []byte if sseConfig != nil { - configData, err := base64.StdEncoding.DecodeString(*sseConfig) + var err error + data, err = base64.StdEncoding.DecodeString(*sseConfig) if err != nil { return wrapSRErr(err) } - _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketSSEConfig, configData) - if err != nil { - return wrapSRErr(err) - } - return nil } - - // Delete sse config - _, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketSSEConfig) + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketSSEConfig, data, false, false, &updatedAt) if err != nil { return wrapSRErr(err) } @@ -2046,7 +1993,7 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket meta.CorsConfigXML = bytes.Clone(configData) meta.CorsConfigUpdatedAt = updatedAt - if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil { + if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, &meta); err != nil { return time.Time{}, err } unlock() @@ -2078,32 +2025,18 @@ func (c *SiteReplicationSys) PeerBucketCorsConfigHandler(ctx context.Context, bu // PeerBucketQuotaConfigHandler - copies/deletes policy to local cluster. func (c *SiteReplicationSys) PeerBucketQuotaConfigHandler(ctx context.Context, bucket string, quota *madmin.BucketQuota, updatedAt time.Time) error { - // skip overwrite if local update is newer than peer update. - if !updatedAt.IsZero() { - if _, updateTm, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket); err == nil && updateTm.After(updatedAt) { - return nil - } - } - + var data []byte if quota != nil { - quotaData, err := json.Marshal(quota) + var err error + data, err = json.Marshal(quota) if err != nil { return wrapSRErr(err) } - - if _, err = globalBucketMetadataSys.Update(ctx, bucket, bucketQuotaConfigFile, quotaData); err != nil { - return wrapSRErr(err) - } - - return nil } - - // Delete the bucket policy - _, err := globalBucketMetadataSys.Delete(ctx, bucket, bucketQuotaConfigFile) + _, err := globalBucketMetadataSys.updateAndParseMetadata(ctx, bucket, bucketQuotaConfigFile, data, false, false, &updatedAt) if err != nil { return wrapSRErr(err) } - return nil }