From 7ade0c045b580cac6c80d56a2162913b9bd70bb5 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Wed, 2 Sep 2026 00:38:16 +0800 Subject: [PATCH] fix: serialize bucket metadata updates across config types Use one per-bucket metadata.lock for ordinary updates, CORS transitions, and legacy bulk replication. Persist and update the local cache while locked, then release before peer metadata reload fan-out.\n\nRefs: #102 Signed-off-by: Feng Ruohang --- cmd/bucket-metadata-sys.go | 141 ++++++++++++++++++++++--------------- cmd/site-replication.go | 50 +++++++------ 2 files changed, 113 insertions(+), 78 deletions(-) diff --git a/cmd/bucket-metadata-sys.go b/cmd/bucket-metadata-sys.go index fea022736..8cb8e9e17 100644 --- a/cmd/bucket-metadata-sys.go +++ b/cmd/bucket-metadata-sys.go @@ -123,63 +123,74 @@ func (sys *BucketMetadataSys) updateAndParse(ctx context.Context, bucket string, if isMinioMetaBucketName(bucket) { return updatedAt, errInvalidArgument } - - meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, parse) + ctx, unlock, err := lockBucketMetadata(ctx, objAPI, bucket) if err != nil { - if !globalIsErasure && !globalIsDistErasure && errors.Is(err, errVolumeNotFound) { - // Only single drive mode needs this fallback. - meta = newBucketMetadata(bucket) - } else { - return updatedAt, 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 bucketCorsConfig: - meta.CorsConfigXML = configData - meta.CorsConfigUpdatedAt = 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 updatedAt, fmt.Errorf("Error encrypting bucket target metadata %w", err) - } - meta.BucketTargetsConfigUpdatedAt = updatedAt - meta.BucketTargetsConfigMetaUpdatedAt = updatedAt - default: - return updatedAt, fmt.Errorf("Unknown bucket %s metadata update requested %s", bucket, configFile) + return updatedAt, err } - return updatedAt, sys.save(ctx, meta) + err = func() error { + defer unlock() + meta, err := loadBucketMetadataParse(ctx, objAPI, bucket, parse) + if err != nil { + if !globalIsErasure && !globalIsDistErasure && errors.Is(err, errVolumeNotFound) { + // Only single drive mode needs this fallback. + meta = newBucketMetadata(bucket) + } else { + 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 bucketCorsConfig: + meta.CorsConfigXML = configData + meta.CorsConfigUpdatedAt = 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) + } + 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 != nil { + return updatedAt, err + } + globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket) // Do not use caller context here + return updatedAt, nil } func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) error { @@ -192,15 +203,33 @@ func (sys *BucketMetadataSys) save(ctx context.Context, meta BucketMetadata) err return errInvalidArgument } - if err := meta.Save(ctx, objAPI); err != nil { + if err := sys.saveMetadata(ctx, objAPI, meta); err != nil { return err } - sys.Set(meta.Name, meta) globalNotificationSys.LoadBucketMetadata(bgContext(ctx), meta.Name) // Do not use caller context here return nil } +// 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 { + if err := meta.Save(ctx, objAPI); err != nil { + return err + } + sys.Set(meta.Name, meta) + return nil +} + +func lockBucketMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (context.Context, func(), error) { + lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "metadata.lock")) + lkctx, err := lock.GetLock(ctx, globalOperationTimeout) + if err != nil { + return nil, nil, err + } + return lkctx.Context(), func() { lock.Unlock(lkctx) }, nil +} + // Delete delete the bucket metadata for the specified bucket. // must be used by all callers instead of using Update() with nil configData. func (sys *BucketMetadataSys) Delete(ctx context.Context, bucket string, configFile string) (updatedAt time.Time, err error) { diff --git a/cmd/site-replication.go b/cmd/site-replication.go index 399260583..aff230822 100644 --- a/cmd/site-replication.go +++ b/cmd/site-replication.go @@ -1620,13 +1620,17 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context if err = validateCORSReplicationPayload(corsConfigData); err != nil { return wrapSRErr(err) } - var unlock func() - ctx, unlock, err = lockBucketCORSMetadata(ctx, objectAPI, item.Bucket) - if err != nil { - return wrapSRErr(err) - } - defer unlock() } + ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, item.Bucket) + if err != nil { + return wrapSRErr(err) + } + locked := true + defer func() { + if locked { + unlock() + } + }() meta, err := readBucketMetadata(ctx, objectAPI, item.Bucket) if err != nil { @@ -1695,7 +1699,13 @@ func (c *SiteReplicationSys) PeerBucketMetadataUpdateHandler(ctx context.Context } } - return globalBucketMetadataSys.save(ctx, meta) + if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil { + return err + } + unlock() + locked = false + globalNotificationSys.LoadBucketMetadata(bgContext(ctx), item.Bucket) + return nil } // PeerBucketPolicyHandler - copies/deletes policy to local cluster. @@ -1960,18 +1970,6 @@ func newBucketCORSReplicationEvent(bucket string, meta BucketMetadata) (madmin.S }, true } -func lockBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string) (context.Context, func(), error) { - // The lock name is deliberately different from .metadata.bin. Saving the - // metadata locks that object internally, and namespace locks are not - // re-entrant. - lock := objectAPI.NewNSLock(minioMetaBucket, pathJoin(bucketMetaPrefix, bucket, "cors-config.lock")) - lkctx, err := lock.GetLock(ctx, globalOperationTimeout) - if err != nil { - return nil, nil, err - } - return lkctx.Context(), func() { lock.Unlock(lkctx) }, nil -} - func updateLocalBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket string, configData []byte) (time.Time, error) { return applyBucketCORSMetadata(ctx, objectAPI, bucket, configData, time.Time{}, true) } @@ -1984,11 +1982,16 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket return time.Time{}, err } - ctx, unlock, err := lockBucketCORSMetadata(ctx, objectAPI, bucket) + ctx, unlock, err := lockBucketMetadata(ctx, objectAPI, bucket) if err != nil { return time.Time{}, err } - defer unlock() + locked := true + defer func() { + if locked { + unlock() + } + }() var meta BucketMetadata if local { @@ -2026,9 +2029,12 @@ func applyBucketCORSMetadata(ctx context.Context, objectAPI ObjectLayer, bucket meta.CorsConfigXML = bytes.Clone(configData) meta.CorsConfigUpdatedAt = updatedAt - if err = globalBucketMetadataSys.save(ctx, meta); err != nil { + if err = globalBucketMetadataSys.saveMetadata(ctx, objectAPI, meta); err != nil { return time.Time{}, err } + unlock() + locked = false + globalNotificationSys.LoadBucketMetadata(bgContext(ctx), bucket) return updatedAt, nil }