mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 14:02:00 +02:00
s3: lazy, bounded, entry-free per-bucket caches (#10213)
* s3: stop warming every bucket's config at startup Listing all buckets in BucketRegistry.init() made S3 gateway startup O(buckets) and pinned every bucket's metadata and config resident, which does not scale past a few hundred thousand buckets. Both caches already have lazy miss paths, so load on first access instead and let the metadata subscription refresh only entries already resident; cold buckets cost one filer round-trip on their first request. * s3: bound the per-bucket caches with LRU eviction The bucket config cache, bucket registry, and their negative caches were plain maps that only ever grew: the config cache TTL made Get miss but never evicted the entry, and the not-found sets grew on every probe of a nonexistent bucket name. Cap all four at 65536 entries with LRU eviction so a gateway keeps its hot working set and evicted buckets reload from the filer on next access. * s3: cache parsed bucket config instead of the full filer entry Each cached BucketConfig retained the whole bucket entry (extended attribute map plus raw content bytes) alongside the fields parsed from it, roughly doubling per-bucket cache cost and keeping data the read path never looks at. Parse everything up front in newBucketConfigFromEntry - now also the creator identity, tags, encryption config, and stored lifecycle XML - and drop the entry. updateBucketConfig now reads the entry fresh from the filer and diffs the mapped extended attributes against it, so the patch is computed against current state instead of a cached copy; the config clone helpers that existed for that path go away. * s3: dedup cold bucket-registry loads per bucket The registry's notFound lock doubled as the load serializer, holding one global mutex across the filer round-trip so first-touch requests for different buckets queued behind each other; the cache fill also happened after the lock was released, so two concurrent misses for the same bucket could both reach the filer. Replace it with a singleflight per bucket that fills the cache inside the flight: different buckets load concurrently, the same bucket loads once.
This commit is contained in:
1 parent
3089480c30
commit
c4f0b12a9a
19 files changed
+328
-623
No files matched your search
@@ -133,6 +133,7 @@ require (
|
||||
github.com/go-ldap/ldap/v3 v3.4.13
|
||||
github.com/golang-jwt/jwt/v5 v5.3.1
|
||||
github.com/google/flatbuffers/go v0.0.0-20230108230133-3b8644d32c50
|
||||
github.com/hashicorp/golang-lru/v2 v2.0.7
|
||||
github.com/hashicorp/raft v1.7.3
|
||||
github.com/hashicorp/raft-boltdb/v2 v2.3.1
|
||||
github.com/hashicorp/vault/api v1.23.0
|
||||
@@ -219,7 +220,6 @@ require (
|
||||
github.com/hashicorp/go-secure-stdlib/parseutil v0.2.0 // indirect
|
||||
github.com/hashicorp/go-secure-stdlib/strutil v0.1.2 // indirect
|
||||
github.com/hashicorp/go-sockaddr v1.0.7 // indirect
|
||||
github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect
|
||||
github.com/hashicorp/hcl v1.0.1-vault-7 // indirect
|
||||
github.com/internxt/rclone-adapter v0.0.0-20260331173834-036f908d0160 // indirect
|
||||
github.com/jackc/pgpassfile v1.0.0 // indirect
|
||||
|
||||
@@ -54,13 +54,9 @@ func TestUnscopedIdentitiesGetDistinctAccounts(t *testing.T) {
|
||||
func TestCheckAccessByOwnershipDeniesNonOwner(t *testing.T) {
|
||||
adminOwner := AccountAdmin.Id
|
||||
s3a := &S3ApiServer{
|
||||
bucketRegistry: &BucketRegistry{
|
||||
metadataCache: map[string]*BucketMetaData{
|
||||
"b": {Name: "b", Owner: &s3.Owner{ID: &adminOwner}},
|
||||
},
|
||||
notFound: map[string]struct{}{},
|
||||
},
|
||||
bucketRegistry: NewBucketRegistry(nil),
|
||||
}
|
||||
s3a.bucketRegistry.setMetadataCache(&BucketMetaData{Name: "b", Owner: &s3.Owner{ID: &adminOwner}})
|
||||
|
||||
nonOwner := httptest.NewRequest(http.MethodGet, "/b?ownershipControls=", nil)
|
||||
nonOwner.Header.Set(s3_constants.AmzAccountId, "alice")
|
||||
|
||||
@@ -3,7 +3,6 @@ package s3api
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
@@ -211,28 +210,24 @@ func (s3a *S3ApiServer) updateBucketConfigCacheFromEntry(entry *filer_pb.Entry)
|
||||
|
||||
bucket := entry.Name
|
||||
|
||||
glog.V(3).Infof("updateBucketConfigCacheFromEntry: called for bucket %s, ExtObjectLockEnabledKey=%s",
|
||||
bucket, string(entry.Extended[s3_constants.ExtObjectLockEnabledKey]))
|
||||
|
||||
// Create new bucket config from the entry. populateBucketConfigDerivedFields
|
||||
// is the single source of truth for mapping Entry.Extended → cached
|
||||
// fields (incl. LifecycleTTL), so a meta-log Put/DeleteBucketLifecycle
|
||||
// here can't leave a stale resolver in cache.
|
||||
config := &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: entry,
|
||||
}
|
||||
s3a.populateBucketConfigDerivedFields(config)
|
||||
|
||||
// Update timestamp
|
||||
config.LastModified = time.Now()
|
||||
|
||||
// Update cache
|
||||
glog.V(3).Infof("updateBucketConfigCacheFromEntry: updating cache for bucket %s, ObjectLockConfig=%+v", bucket, config.ObjectLockConfig)
|
||||
s3a.bucketConfigCache.Set(bucket, config)
|
||||
// Remove from negative cache since bucket now exists
|
||||
// This is important for buckets created via weed shell or other external means
|
||||
s3a.bucketConfigCache.RemoveNegativeCache(bucket)
|
||||
|
||||
// Only refresh buckets already resident in the cache; cold buckets
|
||||
// lazy-load on first access so the cache holds this gateway's working
|
||||
// set, not every bucket in the cluster.
|
||||
if !s3a.bucketConfigCache.Contains(bucket) {
|
||||
return
|
||||
}
|
||||
|
||||
// newBucketConfigFromEntry is the single source of truth for mapping
|
||||
// Entry.Extended → cached fields (incl. LifecycleTTL), so a meta-log
|
||||
// Put/DeleteBucketLifecycle here can't leave a stale resolver in cache.
|
||||
config := s3a.newBucketConfigFromEntry(bucket, entry)
|
||||
|
||||
glog.V(3).Infof("updateBucketConfigCacheFromEntry: refreshing cache for bucket %s, ObjectLockConfig=%+v", bucket, config.ObjectLockConfig)
|
||||
s3a.bucketConfigCache.Set(bucket, config)
|
||||
}
|
||||
|
||||
// invalidateBucketConfigCache removes a bucket from the configuration cache
|
||||
|
||||
@@ -1,18 +1,16 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"math"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/aws/aws-sdk-go/service/s3"
|
||||
lru "github.com/hashicorp/golang-lru/v2"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
|
||||
"golang.org/x/sync/singleflight"
|
||||
)
|
||||
|
||||
var loadBucketMetadataFromFiler = func(r *BucketRegistry, bucketName string) (*BucketMetaData, error) {
|
||||
@@ -46,54 +44,35 @@ type BucketMetaData struct {
|
||||
}
|
||||
|
||||
type BucketRegistry struct {
|
||||
metadataCache map[string]*BucketMetaData
|
||||
metadataCacheLock sync.RWMutex
|
||||
metadataCache *lru.Cache[string, *BucketMetaData]
|
||||
|
||||
notFound map[string]struct{}
|
||||
notFoundLock sync.RWMutex
|
||||
s3a *S3ApiServer
|
||||
notFound *lru.Cache[string, struct{}]
|
||||
// loadGroup deduplicates concurrent filer loads of the same bucket
|
||||
// without serializing loads of different buckets
|
||||
loadGroup singleflight.Group
|
||||
s3a *S3ApiServer
|
||||
}
|
||||
|
||||
// NewBucketRegistry creates a lazy registry: nothing is listed at startup,
|
||||
// buckets load from the filer on first access and stay fresh via the
|
||||
// metadata subscription.
|
||||
func NewBucketRegistry(s3a *S3ApiServer) *BucketRegistry {
|
||||
br := &BucketRegistry{
|
||||
metadataCache: make(map[string]*BucketMetaData),
|
||||
notFound: make(map[string]struct{}),
|
||||
metadataCache, _ := lru.New[string, *BucketMetaData](bucketCacheCapacity)
|
||||
notFound, _ := lru.New[string, struct{}](bucketCacheCapacity)
|
||||
return &BucketRegistry{
|
||||
metadataCache: metadataCache,
|
||||
notFound: notFound,
|
||||
s3a: s3a,
|
||||
}
|
||||
err := br.init()
|
||||
if err != nil {
|
||||
glog.Fatal("init bucket registry failed", err)
|
||||
return nil
|
||||
}
|
||||
return br
|
||||
}
|
||||
|
||||
func (r *BucketRegistry) init() error {
|
||||
var bucketCount int
|
||||
err := filer_pb.List(context.Background(), r.s3a, r.s3a.option.BucketsPath, "", func(entry *filer_pb.Entry, isLast bool) error {
|
||||
if entry != nil && strings.HasPrefix(entry.Name, ".") {
|
||||
return nil
|
||||
}
|
||||
r.LoadBucketMetadata(entry)
|
||||
// Also warm the bucket config cache with Object Lock and versioning settings
|
||||
// This ensures cache consistency across multi-filer clusters after restart
|
||||
r.s3a.updateBucketConfigCacheFromEntry(entry)
|
||||
bucketCount++
|
||||
return nil
|
||||
}, "", false, math.MaxUint32)
|
||||
if err != nil {
|
||||
glog.Errorf("BucketRegistry.init: failed to list buckets: %v", err)
|
||||
return err
|
||||
}
|
||||
glog.V(1).Infof("BucketRegistry.init: warmed config cache for %d buckets", bucketCount)
|
||||
return nil
|
||||
}
|
||||
|
||||
// LoadBucketMetadata refreshes a bucket already resident in the cache from a
|
||||
// subscription event. Cold buckets are left to lazy-load on first access so
|
||||
// the cache holds only this gateway's working set.
|
||||
func (r *BucketRegistry) LoadBucketMetadata(entry *filer_pb.Entry) {
|
||||
bucketMetadata := buildBucketMetadata(r.s3a.iam, entry)
|
||||
r.metadataCacheLock.Lock()
|
||||
r.metadataCache[entry.Name] = bucketMetadata
|
||||
r.metadataCacheLock.Unlock()
|
||||
if r.metadataCache.Contains(entry.Name) {
|
||||
r.metadataCache.Add(entry.Name, buildBucketMetadata(r.s3a.iam, entry))
|
||||
}
|
||||
// Remove from notFound cache since bucket now exists
|
||||
r.unMarkNotFound(entry.Name)
|
||||
}
|
||||
@@ -163,69 +142,58 @@ func (r *BucketRegistry) RemoveBucketMetadata(entry *filer_pb.Entry) {
|
||||
}
|
||||
|
||||
func (r *BucketRegistry) GetBucketMetadata(bucketName string) (*BucketMetaData, s3err.ErrorCode) {
|
||||
r.metadataCacheLock.RLock()
|
||||
bucketMetadata, ok := r.metadataCache[bucketName]
|
||||
r.metadataCacheLock.RUnlock()
|
||||
bucketMetadata, ok := r.metadataCache.Get(bucketName)
|
||||
if ok {
|
||||
return bucketMetadata, s3err.ErrNone
|
||||
}
|
||||
|
||||
r.notFoundLock.RLock()
|
||||
_, ok = r.notFound[bucketName]
|
||||
r.notFoundLock.RUnlock()
|
||||
if ok {
|
||||
if r.notFound.Contains(bucketName) {
|
||||
return nil, s3err.ErrNoSuchBucket
|
||||
}
|
||||
|
||||
bucketMetadata, errCode := r.LoadBucketMetadataFromFiler(bucketName)
|
||||
if errCode != s3err.ErrNone {
|
||||
return nil, errCode
|
||||
}
|
||||
|
||||
r.setMetadataCache(bucketMetadata)
|
||||
r.unMarkNotFound(bucketName)
|
||||
return bucketMetadata, s3err.ErrNone
|
||||
return r.LoadBucketMetadataFromFiler(bucketName)
|
||||
}
|
||||
|
||||
// LoadBucketMetadataFromFiler loads the bucket from the filer; concurrent
|
||||
// calls for the same bucket share one load, and the cache is filled inside
|
||||
// the flight so a bucket is fetched only once.
|
||||
func (r *BucketRegistry) LoadBucketMetadataFromFiler(bucketName string) (*BucketMetaData, s3err.ErrorCode) {
|
||||
r.notFoundLock.Lock()
|
||||
defer r.notFoundLock.Unlock()
|
||||
metadata, err, _ := r.loadGroup.Do(bucketName, func() (interface{}, error) {
|
||||
//check if already exists
|
||||
if bucketMetaData, ok := r.metadataCache.Get(bucketName); ok {
|
||||
return bucketMetaData, nil
|
||||
}
|
||||
|
||||
//check if already exists
|
||||
r.metadataCacheLock.RLock()
|
||||
bucketMetaData, ok := r.metadataCache[bucketName]
|
||||
r.metadataCacheLock.RUnlock()
|
||||
if ok {
|
||||
return bucketMetaData, s3err.ErrNone
|
||||
}
|
||||
|
||||
//if not exists, load from filer
|
||||
bucketMetadata, err := loadBucketMetadataFromFiler(r, bucketName)
|
||||
//if not exists, load from filer
|
||||
bucketMetadata, err := loadBucketMetadataFromFiler(r, bucketName)
|
||||
if err != nil {
|
||||
if err == filer_pb.ErrNotFound {
|
||||
// The bucket doesn't actually exist and should no longer loaded from the filer
|
||||
r.notFound.Add(bucketName, struct{}{})
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
r.setMetadataCache(bucketMetadata)
|
||||
r.unMarkNotFound(bucketName)
|
||||
return bucketMetadata, nil
|
||||
})
|
||||
if err != nil {
|
||||
if err == filer_pb.ErrNotFound {
|
||||
// The bucket doesn't actually exist and should no longer loaded from the filer
|
||||
r.notFound[bucketName] = struct{}{}
|
||||
return nil, s3err.ErrNoSuchBucket
|
||||
}
|
||||
return nil, s3err.ErrInternalError
|
||||
}
|
||||
return bucketMetadata, s3err.ErrNone
|
||||
return metadata.(*BucketMetaData), s3err.ErrNone
|
||||
}
|
||||
|
||||
func (r *BucketRegistry) setMetadataCache(metadata *BucketMetaData) {
|
||||
r.metadataCacheLock.Lock()
|
||||
defer r.metadataCacheLock.Unlock()
|
||||
r.metadataCache[metadata.Name] = metadata
|
||||
r.metadataCache.Add(metadata.Name, metadata)
|
||||
}
|
||||
|
||||
func (r *BucketRegistry) removeMetadataCache(bucket string) {
|
||||
r.metadataCacheLock.Lock()
|
||||
defer r.metadataCacheLock.Unlock()
|
||||
delete(r.metadataCache, bucket)
|
||||
r.metadataCache.Remove(bucket)
|
||||
}
|
||||
|
||||
func (r *BucketRegistry) unMarkNotFound(bucket string) {
|
||||
r.notFoundLock.Lock()
|
||||
defer r.notFoundLock.Unlock()
|
||||
delete(r.notFound, bucket)
|
||||
r.notFound.Remove(bucket)
|
||||
}
|
||||
@@ -79,7 +79,8 @@ var (
|
||||
}
|
||||
|
||||
//load filer is
|
||||
loadFilerBucket = make(map[string]int, 1)
|
||||
loadFilerBucket = make(map[string]int, 1)
|
||||
loadFilerBucketLock sync.Mutex
|
||||
//override `loadBucketMetadataFromFiler` to avoid really load from filer
|
||||
)
|
||||
|
||||
@@ -177,17 +178,15 @@ func TestBuildBucketMetadata(t *testing.T) {
|
||||
func TestGetBucketMetadata(t *testing.T) {
|
||||
loadBucketMetadataFromFiler = func(r *BucketRegistry, bucketName string) (*BucketMetaData, error) {
|
||||
time.Sleep(time.Second)
|
||||
loadFilerBucketLock.Lock()
|
||||
loadFilerBucket[bucketName] = loadFilerBucket[bucketName] + 1
|
||||
loadFilerBucketLock.Unlock()
|
||||
return &BucketMetaData{
|
||||
Name: bucketName,
|
||||
}, nil
|
||||
}
|
||||
|
||||
br := &BucketRegistry{
|
||||
metadataCache: make(map[string]*BucketMetaData),
|
||||
notFound: make(map[string]struct{}),
|
||||
s3a: nil,
|
||||
}
|
||||
br := NewBucketRegistry(nil)
|
||||
|
||||
//start 40 goroutine for
|
||||
var wg sync.WaitGroup
|
||||
|
||||
@@ -7,6 +7,7 @@ import (
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
|
||||
)
|
||||
|
||||
@@ -17,21 +18,19 @@ func (s3a *S3ApiServer) isTableBucket(bucket string) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
// Check cache first
|
||||
if s3a.bucketRegistry != nil {
|
||||
s3a.bucketRegistry.metadataCacheLock.RLock()
|
||||
if metadata, ok := s3a.bucketRegistry.metadataCache[bucket]; ok {
|
||||
s3a.bucketRegistry.metadataCacheLock.RUnlock()
|
||||
return metadata.IsTableBucket
|
||||
metadata, errCode := s3a.bucketRegistry.GetBucketMetadata(bucket)
|
||||
if errCode != s3err.ErrNone {
|
||||
if errCode != s3err.ErrNoSuchBucket {
|
||||
glog.V(1).Infof("bucket lookup failed for %s: %v", bucket, errCode)
|
||||
}
|
||||
return false
|
||||
}
|
||||
s3a.bucketRegistry.metadataCacheLock.RUnlock()
|
||||
return metadata.IsTableBucket
|
||||
}
|
||||
|
||||
entry, err := s3a.getEntry(s3a.option.BucketsPath, bucket)
|
||||
if err == nil && entry != nil {
|
||||
if s3a.bucketRegistry != nil {
|
||||
s3a.bucketRegistry.LoadBucketMetadata(entry)
|
||||
}
|
||||
return s3tables.IsTableBucketEntry(entry)
|
||||
}
|
||||
|
||||
|
||||
@@ -337,24 +337,16 @@ func (s3a *S3ApiServer) CleanupBucketKMSCache(bucketName string) int {
|
||||
func (s3a *S3ApiServer) CleanupAllBucketKMSCaches() int {
|
||||
totalCleaned := 0
|
||||
|
||||
// Access the bucket config cache safely
|
||||
if s3a.bucketConfigCache != nil {
|
||||
s3a.bucketConfigCache.mutex.RLock()
|
||||
bucketNames := make([]string, 0, len(s3a.bucketConfigCache.cache))
|
||||
for bucketName := range s3a.bucketConfigCache.cache {
|
||||
bucketNames = append(bucketNames, bucketName)
|
||||
}
|
||||
s3a.bucketConfigCache.mutex.RUnlock()
|
||||
|
||||
// Clean up each bucket's KMS cache
|
||||
for _, bucketName := range bucketNames {
|
||||
// Clean up each cached bucket's KMS cache
|
||||
for _, bucketName := range s3a.bucketConfigCache.cache.Keys() {
|
||||
cleaned := s3a.CleanupBucketKMSCache(bucketName)
|
||||
totalCleaned += cleaned
|
||||
}
|
||||
}
|
||||
|
||||
if totalCleaned > 0 {
|
||||
glog.V(2).Infof("Cleaned up %d expired KMS keys across %d bucket caches", totalCleaned, len(s3a.bucketConfigCache.cache))
|
||||
glog.V(2).Infof("Cleaned up %d expired KMS keys across %d bucket caches", totalCleaned, s3a.bucketConfigCache.cache.Len())
|
||||
}
|
||||
return totalCleaned
|
||||
}
|
||||
|
||||
@@ -349,13 +349,16 @@ func buildAccessControlList(accountManager AccountManager, grants []*s3.Grant, o
|
||||
|
||||
// GetAcpGrants return grants parsed from entry
|
||||
func GetAcpGrants(entryExtended map[string][]byte) []*s3.Grant {
|
||||
acpBytes, ok := entryExtended[s3_constants.ExtAmzAclKey]
|
||||
if ok && len(acpBytes) > 0 {
|
||||
var grants []*s3.Grant
|
||||
err := json.Unmarshal(acpBytes, &grants)
|
||||
if err == nil {
|
||||
return grants
|
||||
}
|
||||
return parseAclGrants(entryExtended[s3_constants.ExtAmzAclKey])
|
||||
}
|
||||
|
||||
func parseAclGrants(acpBytes []byte) []*s3.Grant {
|
||||
if len(acpBytes) == 0 {
|
||||
return nil
|
||||
}
|
||||
var grants []*s3.Grant
|
||||
if err := json.Unmarshal(acpBytes, &grants); err == nil {
|
||||
return grants
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
+154
-325
@@ -12,6 +12,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/aws/aws-sdk-go/service/s3"
|
||||
lru "github.com/hashicorp/golang-lru/v2"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
@@ -26,14 +27,18 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
)
|
||||
|
||||
// BucketConfig represents cached bucket configuration
|
||||
// BucketConfig represents cached bucket configuration. Only fields parsed
|
||||
// from the bucket's filer entry are retained — not the entry itself — so a
|
||||
// cached config stays small; the write path (updateBucketConfig) re-reads
|
||||
// the entry from the filer.
|
||||
type BucketConfig struct {
|
||||
Name string
|
||||
Versioning string // "Enabled", "Suspended", or ""
|
||||
Ownership string
|
||||
ACL []byte
|
||||
Owner string
|
||||
IsPublicRead bool // Cached flag to avoid JSON parsing on every request
|
||||
IdentityId string // identity that created the bucket
|
||||
IsPublicRead bool // Cached flag to avoid JSON parsing on every request
|
||||
CORS *cors.CORSConfiguration
|
||||
ObjectLockConfig *ObjectLockConfiguration // Cached parsed Object Lock configuration
|
||||
BucketPolicy *policy_engine.PolicyDocument // Cached bucket policy for performance
|
||||
@@ -45,9 +50,14 @@ type BucketConfig struct {
|
||||
// lifecycle worker reads bucket entries directly off the meta-log
|
||||
// rather than this cache.
|
||||
LifecycleTTL *LifecycleTTLResolver
|
||||
KMSKeyCache *BucketKMSCache // Per-bucket KMS key cache for SSE-KMS operations
|
||||
LastModified time.Time
|
||||
Entry *filer_pb.Entry
|
||||
// LifecycleXML is the stored lifecycle configuration as served by
|
||||
// GetBucketLifecycle, with its companion transition minimum object size.
|
||||
LifecycleXML []byte
|
||||
LifecycleTransitionMinSize string
|
||||
Tags map[string]string // parsed from entry content
|
||||
Encryption *s3_pb.EncryptionConfiguration // parsed from entry content
|
||||
KMSKeyCache *BucketKMSCache // Per-bucket KMS key cache for SSE-KMS operations
|
||||
LastModified time.Time
|
||||
}
|
||||
|
||||
// BucketKMSCache represents per-bucket KMS key caching for SSE-KMS operations
|
||||
@@ -183,15 +193,21 @@ func (bkc *BucketKMSCache) Clear() {
|
||||
bkc.cache = make(map[string]*BucketKMSCacheEntry)
|
||||
}
|
||||
|
||||
// bucketCacheCapacity bounds the per-bucket caches (bucket config cache,
|
||||
// bucket registry, and their negative caches). Only the hot working set
|
||||
// stays resident; evicted buckets reload from the filer on next access.
|
||||
const bucketCacheCapacity = 65536
|
||||
|
||||
// BucketConfigCache provides caching for bucket configurations
|
||||
// Cache entries are automatically updated/invalidated through metadata subscription events,
|
||||
// so TTL serves as a safety fallback rather than the primary consistency mechanism
|
||||
// so TTL serves as a safety fallback rather than the primary consistency mechanism.
|
||||
// Both caches are size-capped LRUs so a gateway serving millions of buckets
|
||||
// only keeps its hot working set resident.
|
||||
type BucketConfigCache struct {
|
||||
cache map[string]*BucketConfig
|
||||
negativeCache map[string]time.Time // Cache for non-existent buckets
|
||||
mutex sync.RWMutex
|
||||
ttl time.Duration // Safety fallback TTL; real-time consistency maintained via events
|
||||
negativeTTL time.Duration // TTL for negative cache entries
|
||||
cache *lru.Cache[string, *BucketConfig]
|
||||
negativeCache *lru.Cache[string, time.Time] // Cache for non-existent buckets
|
||||
ttl time.Duration // Safety fallback TTL; real-time consistency maintained via events
|
||||
negativeTTL time.Duration // TTL for negative cache entries
|
||||
}
|
||||
|
||||
// BucketMetadata represents the complete metadata for a bucket
|
||||
@@ -247,9 +263,11 @@ func NewBucketConfigCache(ttl time.Duration) *BucketConfigCache {
|
||||
negativeTTL = 30 * time.Second // Minimum 30 seconds for negative cache
|
||||
}
|
||||
|
||||
cache, _ := lru.New[string, *BucketConfig](bucketCacheCapacity)
|
||||
negativeCache, _ := lru.New[string, time.Time](bucketCacheCapacity)
|
||||
return &BucketConfigCache{
|
||||
cache: make(map[string]*BucketConfig),
|
||||
negativeCache: make(map[string]time.Time),
|
||||
cache: cache,
|
||||
negativeCache: negativeCache,
|
||||
ttl: ttl,
|
||||
negativeTTL: negativeTTL,
|
||||
}
|
||||
@@ -257,10 +275,7 @@ func NewBucketConfigCache(ttl time.Duration) *BucketConfigCache {
|
||||
|
||||
// Get retrieves bucket configuration from cache
|
||||
func (bcc *BucketConfigCache) Get(bucket string) (*BucketConfig, bool) {
|
||||
bcc.mutex.RLock()
|
||||
defer bcc.mutex.RUnlock()
|
||||
|
||||
config, exists := bcc.cache[bucket]
|
||||
config, exists := bcc.cache.Get(bucket)
|
||||
if !exists {
|
||||
return nil, false
|
||||
}
|
||||
@@ -275,60 +290,47 @@ func (bcc *BucketConfigCache) Get(bucket string) (*BucketConfig, bool) {
|
||||
|
||||
// Set stores bucket configuration in cache
|
||||
func (bcc *BucketConfigCache) Set(bucket string, config *BucketConfig) {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
config.LastModified = time.Now()
|
||||
bcc.cache[bucket] = config
|
||||
bcc.cache.Add(bucket, config)
|
||||
}
|
||||
|
||||
// Contains reports whether the bucket is resident in the cache, regardless of TTL
|
||||
func (bcc *BucketConfigCache) Contains(bucket string) bool {
|
||||
return bcc.cache.Contains(bucket)
|
||||
}
|
||||
|
||||
// Remove removes bucket configuration from cache
|
||||
func (bcc *BucketConfigCache) Remove(bucket string) {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
delete(bcc.cache, bucket)
|
||||
bcc.cache.Remove(bucket)
|
||||
}
|
||||
|
||||
// Clear clears all cached configurations
|
||||
func (bcc *BucketConfigCache) Clear() {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
bcc.cache = make(map[string]*BucketConfig)
|
||||
bcc.negativeCache = make(map[string]time.Time)
|
||||
bcc.cache.Purge()
|
||||
bcc.negativeCache.Purge()
|
||||
}
|
||||
|
||||
// IsNegativelyCached checks if a bucket is in the negative cache (doesn't exist)
|
||||
func (bcc *BucketConfigCache) IsNegativelyCached(bucket string) bool {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
if cachedTime, exists := bcc.negativeCache[bucket]; exists {
|
||||
if cachedTime, exists := bcc.negativeCache.Get(bucket); exists {
|
||||
// Check if the negative cache entry is still valid
|
||||
if time.Since(cachedTime) < bcc.negativeTTL {
|
||||
return true
|
||||
}
|
||||
// Entry expired, remove it
|
||||
delete(bcc.negativeCache, bucket)
|
||||
bcc.negativeCache.Remove(bucket)
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// SetNegativeCache marks a bucket as non-existent in the negative cache
|
||||
func (bcc *BucketConfigCache) SetNegativeCache(bucket string) {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
bcc.negativeCache[bucket] = time.Now()
|
||||
bcc.negativeCache.Add(bucket, time.Now())
|
||||
}
|
||||
|
||||
// RemoveNegativeCache removes a bucket from the negative cache
|
||||
func (bcc *BucketConfigCache) RemoveNegativeCache(bucket string) {
|
||||
bcc.mutex.Lock()
|
||||
defer bcc.mutex.Unlock()
|
||||
|
||||
delete(bcc.negativeCache, bucket)
|
||||
bcc.negativeCache.Remove(bucket)
|
||||
}
|
||||
|
||||
// loadBucketPolicyFromExtended loads and parses bucket policy from entry extended attributes
|
||||
@@ -377,12 +379,7 @@ func (s3a *S3ApiServer) getBucketConfig(bucket string) (*BucketConfig, s3err.Err
|
||||
return nil, s3err.ErrInternalError
|
||||
}
|
||||
|
||||
config := &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: entry,
|
||||
IsPublicRead: false, // Explicitly default to false for private buckets
|
||||
}
|
||||
s3a.populateBucketConfigDerivedFields(config)
|
||||
config := s3a.newBucketConfigFromEntry(bucket, entry)
|
||||
|
||||
// Cache the result
|
||||
s3a.bucketConfigCache.Set(bucket, config)
|
||||
@@ -390,33 +387,22 @@ func (s3a *S3ApiServer) getBucketConfig(bucket string) (*BucketConfig, s3err.Err
|
||||
return config, s3err.ErrNone
|
||||
}
|
||||
|
||||
// populateBucketConfigDerivedFields fills every field on BucketConfig that is
|
||||
// derived from Entry.Extended / Entry.Content (versioning flag, ACL, owner,
|
||||
// object lock, bucket policy, CORS, lifecycle TTL resolver). It is the
|
||||
// single source of truth for that mapping; callers that take a fresh
|
||||
// BucketConfig (getBucketConfig, updateBucketConfig after the user's update
|
||||
// fn runs, the meta-log subscription cache refresher) all funnel through
|
||||
// here so a missed field can't silently keep stale data — e.g. a stale
|
||||
// LifecycleTTL after a Put/DeleteBucketLifecycle would keep stamping the
|
||||
// old policy's irreversible volume TTL onto new writes.
|
||||
func (s3a *S3ApiServer) populateBucketConfigDerivedFields(config *BucketConfig) {
|
||||
// Reset every derived field so stale values from a previous Entry
|
||||
// don't survive a clear (e.g. DELETE policy → BucketPolicy=nil).
|
||||
config.Versioning = ""
|
||||
config.Ownership = ""
|
||||
config.ACL = nil
|
||||
config.Owner = ""
|
||||
config.IsPublicRead = false
|
||||
config.ObjectLockConfig = nil
|
||||
config.BucketPolicy = nil
|
||||
config.LifecycleTTL = nil
|
||||
config.CORS = nil
|
||||
|
||||
entry := config.Entry
|
||||
// newBucketConfigFromEntry builds a BucketConfig from the bucket's filer
|
||||
// entry, parsing Entry.Extended / Entry.Content into the cached fields
|
||||
// (versioning flag, ACL, owner, object lock, bucket policy, CORS, tags,
|
||||
// encryption, lifecycle TTL resolver). It is the single source of truth for
|
||||
// that mapping; the read path (getBucketConfig), the write path
|
||||
// (updateBucketConfig), and the meta-log subscription cache refresher all
|
||||
// funnel through here so a missed field can't silently keep stale data —
|
||||
// e.g. a stale LifecycleTTL after a Put/DeleteBucketLifecycle would keep
|
||||
// stamping the old policy's irreversible volume TTL onto new writes.
|
||||
func (s3a *S3ApiServer) newBucketConfigFromEntry(bucket string, entry *filer_pb.Entry) *BucketConfig {
|
||||
config := &BucketConfig{
|
||||
Name: bucket,
|
||||
}
|
||||
if entry == nil {
|
||||
return
|
||||
return config
|
||||
}
|
||||
bucket := config.Name
|
||||
|
||||
if entry.Extended != nil {
|
||||
if versioning, exists := entry.Extended[s3_constants.ExtVersioningKey]; exists {
|
||||
@@ -433,31 +419,37 @@ func (s3a *S3ApiServer) populateBucketConfigDerivedFields(config *BucketConfig)
|
||||
if owner, exists := entry.Extended[s3_constants.ExtAmzOwnerKey]; exists {
|
||||
config.Owner = string(owner)
|
||||
}
|
||||
if identityId, exists := entry.Extended[s3_constants.AmzIdentityId]; exists {
|
||||
config.IdentityId = string(identityId)
|
||||
}
|
||||
if objectLockConfig, found := LoadObjectLockConfigurationFromExtended(entry); found {
|
||||
config.ObjectLockConfig = objectLockConfig
|
||||
}
|
||||
config.BucketPolicy = loadBucketPolicyFromExtended(entry, bucket)
|
||||
|
||||
if lifecycleXML, exists := entry.Extended[bucketLifecycleConfigurationXMLKey]; exists && len(lifecycleXML) > 0 {
|
||||
config.LifecycleXML = lifecycleXML
|
||||
config.LifecycleTransitionMinSize = string(entry.Extended[bucketLifecycleTransitionMinimumObjectSizeKey])
|
||||
}
|
||||
|
||||
// The lifecycle TTL fast path is opt-in per bucket: a volume TTL
|
||||
// stamped at write time can't honor a later policy change (rule
|
||||
// removed or lengthened) the way worker-driven expiration does,
|
||||
// so it stays off unless explicitly enabled. Skip the XML parse
|
||||
// entirely when off. nil on parse error so the PUT path falls
|
||||
// through to "no TTL" rather than rejecting writes.
|
||||
if bytes.Equal(entry.Extended[s3_constants.ExtLifecycleTtlFastPathKey], []byte("true")) {
|
||||
if xmlBytes, ok := entry.Extended[bucketLifecycleConfigurationXMLKey]; ok && len(xmlBytes) > 0 {
|
||||
if rules, err := lifecycle_xml.ParseCanonical(xmlBytes); err == nil {
|
||||
// Object Lock requires versioning, so an ObjectLockConfig
|
||||
// implies the bucket is versioned even when the explicit
|
||||
// Versioning header was never written. BucketIsVersioned
|
||||
// in this file uses the same OR — keep them aligned.
|
||||
versioned := config.Versioning == s3_constants.VersioningEnabled ||
|
||||
config.Versioning == s3_constants.VersioningSuspended ||
|
||||
config.ObjectLockConfig != nil
|
||||
config.LifecycleTTL = NewLifecycleTTLResolver(rules, versioned)
|
||||
} else {
|
||||
glog.V(1).Infof("populateBucketConfigDerivedFields: bucket %s lifecycle xml parse: %v", bucket, err)
|
||||
}
|
||||
if bytes.Equal(entry.Extended[s3_constants.ExtLifecycleTtlFastPathKey], []byte("true")) && len(config.LifecycleXML) > 0 {
|
||||
if rules, err := lifecycle_xml.ParseCanonical(config.LifecycleXML); err == nil {
|
||||
// Object Lock requires versioning, so an ObjectLockConfig
|
||||
// implies the bucket is versioned even when the explicit
|
||||
// Versioning header was never written. BucketIsVersioned
|
||||
// in this file uses the same OR — keep them aligned.
|
||||
versioned := config.Versioning == s3_constants.VersioningEnabled ||
|
||||
config.Versioning == s3_constants.VersioningSuspended ||
|
||||
config.ObjectLockConfig != nil
|
||||
config.LifecycleTTL = NewLifecycleTTLResolver(rules, versioned)
|
||||
} else {
|
||||
glog.V(1).Infof("newBucketConfigFromEntry: bucket %s lifecycle xml parse: %v", bucket, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -465,75 +457,57 @@ func (s3a *S3ApiServer) populateBucketConfigDerivedFields(config *BucketConfig)
|
||||
// Sync bucket policy to the policy engine for evaluation.
|
||||
s3a.syncBucketPolicyToEngine(bucket, config.BucketPolicy)
|
||||
|
||||
// Parse CORS configuration directly from the entry's Content field.
|
||||
// This avoids a redundant RPC call since we already have the entry.
|
||||
config.CORS = parseCORSFromEntryContent(entry.Content)
|
||||
}
|
||||
|
||||
// updateBucketConfig updates bucket configuration and invalidates cache
|
||||
func (s3a *S3ApiServer) updateBucketConfig(bucket string, updateFn func(*BucketConfig) error) s3err.ErrorCode {
|
||||
config, errCode := s3a.getBucketConfig(bucket)
|
||||
if errCode != s3err.ErrNone {
|
||||
return errCode
|
||||
// Parse tags, CORS, and encryption from the entry's Content field.
|
||||
if len(entry.Content) > 0 {
|
||||
var protoMetadata s3_pb.BucketMetadata
|
||||
if err := proto.Unmarshal(entry.Content, &protoMetadata); err != nil {
|
||||
glog.Errorf("newBucketConfigFromEntry: failed to unmarshal metadata for bucket %s: %v", bucket, err)
|
||||
} else {
|
||||
config.Tags = protoMetadata.Tags
|
||||
config.CORS = corsConfigFromProto(protoMetadata.Cors)
|
||||
config.Encryption = protoMetadata.Encryption
|
||||
}
|
||||
}
|
||||
|
||||
nextConfig := cloneBucketConfig(config)
|
||||
if nextConfig == nil {
|
||||
glog.Errorf("updateBucketConfig: failed to clone config for bucket %s", bucket)
|
||||
return config
|
||||
}
|
||||
|
||||
// updateBucketConfig updates bucket configuration and invalidates cache.
|
||||
// It reads the bucket entry fresh from the filer so the patch diff is
|
||||
// computed against current state rather than a possibly stale cached copy.
|
||||
func (s3a *S3ApiServer) updateBucketConfig(bucket string, updateFn func(*BucketConfig) error) s3err.ErrorCode {
|
||||
entry, err := s3a.getBucketEntry(bucket)
|
||||
if err != nil {
|
||||
if errors.Is(err, filer_pb.ErrNotFound) {
|
||||
if s3a.bucketConfigCache != nil {
|
||||
s3a.bucketConfigCache.SetNegativeCache(bucket)
|
||||
}
|
||||
return s3err.ErrNoSuchBucket
|
||||
}
|
||||
glog.Errorf("updateBucketConfig: failed to get bucket entry for %s: %v", bucket, err)
|
||||
return s3err.ErrInternalError
|
||||
}
|
||||
|
||||
config := s3a.newBucketConfigFromEntry(bucket, entry)
|
||||
|
||||
// Apply update function
|
||||
if err := updateFn(nextConfig); err != nil {
|
||||
if err := updateFn(config); err != nil {
|
||||
glog.Errorf("updateBucketConfig: update function failed for bucket %s: %v", bucket, err)
|
||||
return s3err.ErrInternalError
|
||||
}
|
||||
|
||||
// Prepare extended attributes
|
||||
if nextConfig.Entry == nil {
|
||||
glog.Errorf("updateBucketConfig: missing bucket entry for %s", bucket)
|
||||
oldExt := entry.GetExtended()
|
||||
newExt := make(map[string][]byte, len(oldExt))
|
||||
for k, v := range oldExt {
|
||||
newExt[k] = v
|
||||
}
|
||||
if err := applyBucketConfigToExtended(config, newExt); err != nil {
|
||||
glog.Errorf("updateBucketConfig: failed to serialize config for bucket %s: %v", bucket, err)
|
||||
return s3err.ErrInternalError
|
||||
}
|
||||
if nextConfig.Entry.Extended == nil {
|
||||
nextConfig.Entry.Extended = make(map[string][]byte)
|
||||
}
|
||||
|
||||
// Update extended attributes
|
||||
if nextConfig.Versioning != "" {
|
||||
nextConfig.Entry.Extended[s3_constants.ExtVersioningKey] = []byte(nextConfig.Versioning)
|
||||
} else {
|
||||
delete(nextConfig.Entry.Extended, s3_constants.ExtVersioningKey)
|
||||
}
|
||||
if nextConfig.Ownership != "" {
|
||||
nextConfig.Entry.Extended[s3_constants.ExtOwnershipKey] = []byte(nextConfig.Ownership)
|
||||
} else {
|
||||
delete(nextConfig.Entry.Extended, s3_constants.ExtOwnershipKey)
|
||||
}
|
||||
if nextConfig.ACL != nil {
|
||||
nextConfig.Entry.Extended[s3_constants.ExtAmzAclKey] = nextConfig.ACL
|
||||
} else {
|
||||
delete(nextConfig.Entry.Extended, s3_constants.ExtAmzAclKey)
|
||||
}
|
||||
if nextConfig.Owner != "" {
|
||||
nextConfig.Entry.Extended[s3_constants.ExtAmzOwnerKey] = []byte(nextConfig.Owner)
|
||||
} else {
|
||||
delete(nextConfig.Entry.Extended, s3_constants.ExtAmzOwnerKey)
|
||||
}
|
||||
// Update Object Lock configuration
|
||||
if nextConfig.ObjectLockConfig != nil {
|
||||
glog.V(3).Infof("updateBucketConfig: storing Object Lock config for bucket %s: %+v", bucket, nextConfig.ObjectLockConfig)
|
||||
if err := StoreObjectLockConfigurationInExtended(nextConfig.Entry, nextConfig.ObjectLockConfig); err != nil {
|
||||
glog.Errorf("updateBucketConfig: failed to store Object Lock configuration for bucket %s: %v", bucket, err)
|
||||
return s3err.ErrInternalError
|
||||
}
|
||||
glog.V(3).Infof("updateBucketConfig: stored Object Lock config in extended attributes for bucket %s, key=%s, value=%s",
|
||||
bucket, s3_constants.ExtObjectLockEnabledKey, string(nextConfig.Entry.Extended[s3_constants.ExtObjectLockEnabledKey]))
|
||||
}
|
||||
|
||||
// Patch only the changed/removed extended keys, leaving Entry.content
|
||||
// untouched so a concurrent content write (e.g. encryption) is preserved.
|
||||
oldExt := config.Entry.GetExtended()
|
||||
newExt := nextConfig.Entry.Extended
|
||||
set := make(map[string][]byte)
|
||||
for k, v := range newExt {
|
||||
if ov, ok := oldExt[k]; !ok || !bytes.Equal(ov, v) {
|
||||
@@ -552,8 +526,8 @@ func (s3a *S3ApiServer) updateBucketConfig(bucket string, updateFn func(*BucketC
|
||||
return s3err.ErrInternalError
|
||||
}
|
||||
|
||||
// Invalidate rather than cache nextConfig: its content may be stale relative
|
||||
// to a concurrent content write. The next read re-fetches the merged entry.
|
||||
// Invalidate rather than cache the updated config: it may be stale relative
|
||||
// to a concurrent write. The next read re-fetches the merged entry.
|
||||
if s3a.bucketConfigCache != nil {
|
||||
s3a.bucketConfigCache.Remove(bucket)
|
||||
s3a.bucketConfigCache.RemoveNegativeCache(bucket)
|
||||
@@ -562,6 +536,34 @@ func (s3a *S3ApiServer) updateBucketConfig(bucket string, updateFn func(*BucketC
|
||||
return s3err.ErrNone
|
||||
}
|
||||
|
||||
// applyBucketConfigToExtended maps the persisted BucketConfig fields onto an
|
||||
// extended-attribute map; keys not derived from BucketConfig are left alone.
|
||||
func applyBucketConfigToExtended(config *BucketConfig, ext map[string][]byte) error {
|
||||
setOrDelete := func(key string, value []byte) {
|
||||
if len(value) > 0 {
|
||||
ext[key] = value
|
||||
} else {
|
||||
delete(ext, key)
|
||||
}
|
||||
}
|
||||
setOrDelete(s3_constants.ExtVersioningKey, []byte(config.Versioning))
|
||||
setOrDelete(s3_constants.ExtOwnershipKey, []byte(config.Ownership))
|
||||
setOrDelete(s3_constants.ExtAmzAclKey, config.ACL)
|
||||
setOrDelete(s3_constants.ExtAmzOwnerKey, []byte(config.Owner))
|
||||
setOrDelete(bucketLifecycleConfigurationXMLKey, config.LifecycleXML)
|
||||
if len(config.LifecycleXML) > 0 && config.LifecycleTransitionMinSize != "" {
|
||||
ext[bucketLifecycleTransitionMinimumObjectSizeKey] = []byte(config.LifecycleTransitionMinSize)
|
||||
} else {
|
||||
delete(ext, bucketLifecycleTransitionMinimumObjectSizeKey)
|
||||
}
|
||||
// Object Lock, once enabled, is never deleted; a nil config leaves any
|
||||
// existing keys untouched.
|
||||
if config.ObjectLockConfig != nil {
|
||||
return StoreObjectLockConfigurationInExtended(&filer_pb.Entry{Extended: ext}, config.ObjectLockConfig)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// patchBucketEntry applies a field-level PATCH_EXTENDED mutation to the bucket's
|
||||
// entry via ObjectTransaction, routed to the bucket's owner filer so its per-path
|
||||
// lock serializes concurrent config writes cluster-wide rather than racing
|
||||
@@ -598,136 +600,6 @@ func (s3a *S3ApiServer) patchBucketEntry(bucket string, m *filer_pb.ObjectMutati
|
||||
return s3a.WithFilerClient(false, txn)
|
||||
}
|
||||
|
||||
func cloneBucketConfig(config *BucketConfig) *BucketConfig {
|
||||
if config == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
cloned := *config
|
||||
if config.ACL != nil {
|
||||
cloned.ACL = append([]byte(nil), config.ACL...)
|
||||
}
|
||||
if config.Entry != nil {
|
||||
cloned.Entry = proto.Clone(config.Entry).(*filer_pb.Entry)
|
||||
}
|
||||
if config.CORS != nil {
|
||||
cloned.CORS = cloneCORSConfiguration(config.CORS)
|
||||
}
|
||||
if config.ObjectLockConfig != nil {
|
||||
cloned.ObjectLockConfig = cloneObjectLockConfiguration(config.ObjectLockConfig)
|
||||
}
|
||||
if config.BucketPolicy != nil {
|
||||
cloned.BucketPolicy = cloneBucketPolicy(config.BucketPolicy)
|
||||
}
|
||||
|
||||
return &cloned
|
||||
}
|
||||
|
||||
func cloneCORSConfiguration(config *cors.CORSConfiguration) *cors.CORSConfiguration {
|
||||
if config == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
cloned := &cors.CORSConfiguration{
|
||||
CORSRules: make([]cors.CORSRule, len(config.CORSRules)),
|
||||
}
|
||||
for i, rule := range config.CORSRules {
|
||||
cloned.CORSRules[i] = cors.CORSRule{
|
||||
AllowedHeaders: append([]string(nil), rule.AllowedHeaders...),
|
||||
AllowedMethods: append([]string(nil), rule.AllowedMethods...),
|
||||
AllowedOrigins: append([]string(nil), rule.AllowedOrigins...),
|
||||
ExposeHeaders: append([]string(nil), rule.ExposeHeaders...),
|
||||
ID: rule.ID,
|
||||
}
|
||||
if rule.MaxAgeSeconds != nil {
|
||||
maxAge := *rule.MaxAgeSeconds
|
||||
cloned.CORSRules[i].MaxAgeSeconds = &maxAge
|
||||
}
|
||||
}
|
||||
|
||||
return cloned
|
||||
}
|
||||
|
||||
func cloneObjectLockConfiguration(config *ObjectLockConfiguration) *ObjectLockConfiguration {
|
||||
if config == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
cloned := &ObjectLockConfiguration{
|
||||
XMLNS: config.XMLNS,
|
||||
XMLName: config.XMLName,
|
||||
ObjectLockEnabled: config.ObjectLockEnabled,
|
||||
}
|
||||
if config.Rule != nil {
|
||||
cloned.Rule = &ObjectLockRule{
|
||||
XMLName: config.Rule.XMLName,
|
||||
}
|
||||
if config.Rule.DefaultRetention != nil {
|
||||
cloned.Rule.DefaultRetention = &DefaultRetention{
|
||||
XMLName: config.Rule.DefaultRetention.XMLName,
|
||||
Mode: config.Rule.DefaultRetention.Mode,
|
||||
Days: config.Rule.DefaultRetention.Days,
|
||||
Years: config.Rule.DefaultRetention.Years,
|
||||
DaysSet: config.Rule.DefaultRetention.DaysSet,
|
||||
YearsSet: config.Rule.DefaultRetention.YearsSet,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return cloned
|
||||
}
|
||||
|
||||
func cloneBucketPolicy(policyDoc *policy_engine.PolicyDocument) *policy_engine.PolicyDocument {
|
||||
if policyDoc == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
cloned := &policy_engine.PolicyDocument{
|
||||
Version: policyDoc.Version,
|
||||
Statement: make([]policy_engine.PolicyStatement, len(policyDoc.Statement)),
|
||||
}
|
||||
for i, statement := range policyDoc.Statement {
|
||||
cloned.Statement[i] = clonePolicyStatement(statement)
|
||||
}
|
||||
|
||||
return cloned
|
||||
}
|
||||
|
||||
func clonePolicyStatement(statement policy_engine.PolicyStatement) policy_engine.PolicyStatement {
|
||||
cloned := policy_engine.PolicyStatement{
|
||||
Sid: statement.Sid,
|
||||
Effect: statement.Effect,
|
||||
Action: cloneStringOrStringSlice(statement.Action),
|
||||
NotResource: cloneStringOrStringSlicePtr(statement.NotResource),
|
||||
Principal: policy_engine.ClonePolicyPrincipal(statement.Principal),
|
||||
NotPrincipal: policy_engine.ClonePolicyPrincipal(statement.NotPrincipal),
|
||||
Resource: cloneStringOrStringSlicePtr(statement.Resource),
|
||||
}
|
||||
if statement.Condition != nil {
|
||||
cloned.Condition = make(policy_engine.PolicyConditions, len(statement.Condition))
|
||||
for operator, operands := range statement.Condition {
|
||||
copiedOperands := make(map[string]policy_engine.StringOrStringSlice, len(operands))
|
||||
for key, value := range operands {
|
||||
copiedOperands[key] = cloneStringOrStringSlice(value)
|
||||
}
|
||||
cloned.Condition[operator] = copiedOperands
|
||||
}
|
||||
}
|
||||
return cloned
|
||||
}
|
||||
|
||||
func cloneStringOrStringSlice(value policy_engine.StringOrStringSlice) policy_engine.StringOrStringSlice {
|
||||
return policy_engine.CloneStringOrStringSlice(value)
|
||||
}
|
||||
|
||||
func cloneStringOrStringSlicePtr(value *policy_engine.StringOrStringSlice) *policy_engine.StringOrStringSlice {
|
||||
if value == nil {
|
||||
return nil
|
||||
}
|
||||
cloned := policy_engine.CloneStringOrStringSlice(*value)
|
||||
return &cloned
|
||||
}
|
||||
|
||||
// isVersioningEnabled checks if versioning is enabled for a bucket (with caching)
|
||||
func (s3a *S3ApiServer) isVersioningEnabled(bucket string) (bool, error) {
|
||||
config, errCode := s3a.getBucketConfig(bucket)
|
||||
@@ -832,21 +704,6 @@ func (s3a *S3ApiServer) setBucketOwnership(bucket, ownership string) s3err.Error
|
||||
})
|
||||
}
|
||||
|
||||
// parseCORSFromEntryContent parses CORS configuration directly from an entry's Content field.
|
||||
// This avoids a separate RPC call when the entry is already available (e.g., from a
|
||||
// subscription event or a prior getBucketEntry call).
|
||||
func parseCORSFromEntryContent(content []byte) *cors.CORSConfiguration {
|
||||
if len(content) == 0 {
|
||||
return nil
|
||||
}
|
||||
var protoMetadata s3_pb.BucketMetadata
|
||||
if err := proto.Unmarshal(content, &protoMetadata); err != nil {
|
||||
glog.Errorf("parseCORSFromEntryContent: failed to unmarshal protobuf metadata: %v", err)
|
||||
return nil
|
||||
}
|
||||
return corsConfigFromProto(protoMetadata.Cors)
|
||||
}
|
||||
|
||||
// getCORSConfiguration retrieves CORS configuration with caching
|
||||
func (s3a *S3ApiServer) getCORSConfiguration(bucket string) (*cors.CORSConfiguration, s3err.ErrorCode) {
|
||||
config, errCode := s3a.getBucketConfig(bucket)
|
||||
@@ -989,13 +846,16 @@ func (s3a *S3ApiServer) getBucketMetadata(bucket string) (*BucketMetadata, error
|
||||
return nil, fmt.Errorf("bucket directory not found %s", bucket)
|
||||
}
|
||||
|
||||
// Try to get from positive cache
|
||||
// Build from the cached parsed config; copy the tags so callers
|
||||
// can't mutate the cached map.
|
||||
if config, found := s3a.bucketConfigCache.Get(bucket); found {
|
||||
// Extract metadata from cached config
|
||||
if metadata, err := s3a.extractMetadataFromConfig(config); err == nil {
|
||||
return metadata, nil
|
||||
metadata := NewBucketMetadata()
|
||||
for k, v := range config.Tags {
|
||||
metadata.Tags[k] = v
|
||||
}
|
||||
// If extraction fails, fall through to direct load
|
||||
metadata.CORS = config.CORS
|
||||
metadata.Encryption = config.Encryption
|
||||
return metadata, nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1003,37 +863,6 @@ func (s3a *S3ApiServer) getBucketMetadata(bucket string) (*BucketMetadata, error
|
||||
return s3a.loadBucketMetadataFromFiler(bucket)
|
||||
}
|
||||
|
||||
// extractMetadataFromConfig extracts BucketMetadata from cached BucketConfig
|
||||
func (s3a *S3ApiServer) extractMetadataFromConfig(config *BucketConfig) (*BucketMetadata, error) {
|
||||
if config == nil || config.Entry == nil {
|
||||
return NewBucketMetadata(), nil
|
||||
}
|
||||
|
||||
// Parse metadata from entry content if available
|
||||
if len(config.Entry.Content) > 0 {
|
||||
var protoMetadata s3_pb.BucketMetadata
|
||||
if err := proto.Unmarshal(config.Entry.Content, &protoMetadata); err != nil {
|
||||
glog.Errorf("extractMetadataFromConfig: failed to unmarshal protobuf metadata for bucket %s: %v", config.Name, err)
|
||||
return nil, err
|
||||
}
|
||||
// Convert protobuf to structured metadata
|
||||
metadata := &BucketMetadata{
|
||||
Tags: protoMetadata.Tags,
|
||||
CORS: corsConfigFromProto(protoMetadata.Cors),
|
||||
Encryption: protoMetadata.Encryption,
|
||||
}
|
||||
return metadata, nil
|
||||
}
|
||||
|
||||
// Fallback: create metadata from cached CORS config
|
||||
metadata := NewBucketMetadata()
|
||||
if config.CORS != nil {
|
||||
metadata.CORS = config.CORS
|
||||
}
|
||||
|
||||
return metadata, nil
|
||||
}
|
||||
|
||||
// loadBucketMetadataFromFiler loads bucket metadata directly from the filer
|
||||
func (s3a *S3ApiServer) loadBucketMetadataFromFiler(bucket string) (*BucketMetadata, error) {
|
||||
// Validate bucket name to prevent path traversal attacks
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
)
|
||||
|
||||
func TestBucketConfigStubs(t *testing.T) {
|
||||
@@ -17,10 +16,7 @@ func TestBucketConfigStubs(t *testing.T) {
|
||||
iam: &IdentityAccessManagement{isAuthEnabled: true},
|
||||
bucketConfigCache: NewBucketConfigCache(time.Minute),
|
||||
}
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{Name: bucket},
|
||||
})
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{Name: bucket})
|
||||
|
||||
listCases := []struct {
|
||||
name string
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/stretchr/testify/assert"
|
||||
@@ -22,17 +21,12 @@ func TestUpdateBucketConfigDoesNotMutateCacheOnPersistFailure(t *testing.T) {
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Versioning: "",
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: bucket,
|
||||
IsDirectory: true,
|
||||
Extended: map[string][]byte{},
|
||||
},
|
||||
})
|
||||
|
||||
// This test server only has in-memory IAM state and no filer connection, so
|
||||
// updateBucketConfig is expected to fail during the persist step with an
|
||||
// internal error. The assertion below verifies that the cached config stays
|
||||
// unchanged when that write path fails.
|
||||
// updateBucketConfig is expected to fail with an internal error when it
|
||||
// reads the bucket entry. The assertion below verifies that the cached
|
||||
// config stays unchanged when the write path fails.
|
||||
errCode := s3a.updateBucketConfig(bucket, func(config *BucketConfig) error {
|
||||
config.Versioning = s3_constants.VersioningEnabled
|
||||
return nil
|
||||
@@ -43,5 +37,4 @@ func TestUpdateBucketConfigDoesNotMutateCacheOnPersistFailure(t *testing.T) {
|
||||
config, found := s3a.bucketConfigCache.Get(bucket)
|
||||
require.True(t, found)
|
||||
assert.Empty(t, config.Versioning)
|
||||
assert.NotContains(t, config.Entry.Extended, s3_constants.ExtVersioningKey)
|
||||
}
|
||||
@@ -474,7 +474,7 @@ func (s3a *S3ApiServer) checkBucket(r *http.Request, bucket string) s3err.ErrorC
|
||||
if s3a.iam.isEnabled() {
|
||||
return s3err.ErrNone
|
||||
}
|
||||
if !s3a.hasAccess(r, config.Entry) {
|
||||
if !s3a.hasAccess(r, config.IdentityId) {
|
||||
return s3err.ErrAccessDenied
|
||||
}
|
||||
return s3err.ErrNone
|
||||
@@ -648,23 +648,23 @@ func (s3a *S3ApiServer) handleAutoCreateBucket(w http.ResponseWriter, r *http.Re
|
||||
return true
|
||||
}
|
||||
|
||||
func (s3a *S3ApiServer) hasAccess(r *http.Request, entry *filer_pb.Entry) bool {
|
||||
// hasAccess checks the caller against the identity recorded at bucket
|
||||
// creation; buckets with no recorded identity are open to any caller.
|
||||
func (s3a *S3ApiServer) hasAccess(r *http.Request, bucketIdentityId string) bool {
|
||||
// Check if user is properly authenticated as admin through IAM system
|
||||
if s3a.isUserAdmin(r) {
|
||||
return true
|
||||
}
|
||||
|
||||
if entry.Extended == nil {
|
||||
if bucketIdentityId == "" {
|
||||
return true
|
||||
}
|
||||
|
||||
// Get authenticated identity from context (secure, cannot be spoofed)
|
||||
identityId := s3_constants.GetIdentityNameFromContext(r)
|
||||
if id, ok := entry.Extended[s3_constants.AmzIdentityId]; ok {
|
||||
if identityId != string(id) {
|
||||
glog.V(3).Infof("hasAccess: %s != %s (entry.Extended = %v)", identityId, id, entry.Extended)
|
||||
return false
|
||||
}
|
||||
if identityId != bucketIdentityId {
|
||||
glog.V(3).Infof("hasAccess: %s != %s", identityId, bucketIdentityId)
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
@@ -803,8 +803,8 @@ func (s3a *S3ApiServer) GetBucketAclHandler(w http.ResponseWriter, r *http.Reque
|
||||
// return any stored ACL, defaulting to the owner's full-control grant.
|
||||
ownerId := r.Header.Get(s3_constants.AmzAccountId)
|
||||
var storedGrants []*s3.Grant
|
||||
if bucketConfig, errCode := s3a.getBucketConfig(bucket); errCode == s3err.ErrNone && bucketConfig.Entry != nil {
|
||||
storedGrants = GetAcpGrants(bucketConfig.Entry.Extended)
|
||||
if bucketConfig, errCode := s3a.getBucketConfig(bucket); errCode == s3err.ErrNone {
|
||||
storedGrants = parseAclGrants(bucketConfig.ACL)
|
||||
if bucketConfig.Owner != "" {
|
||||
ownerId = bucketConfig.Owner
|
||||
}
|
||||
|
||||
@@ -11,7 +11,6 @@ import (
|
||||
|
||||
"github.com/aws/aws-sdk-go/service/s3"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/policy_engine"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
)
|
||||
@@ -22,10 +21,7 @@ func newMiscTestServer(t *testing.T, bucket string) *S3ApiServer {
|
||||
iam: &IdentityAccessManagement{isAuthEnabled: true},
|
||||
bucketConfigCache: NewBucketConfigCache(time.Minute),
|
||||
}
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{Name: bucket},
|
||||
})
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{Name: bucket})
|
||||
return s3a
|
||||
}
|
||||
|
||||
@@ -133,18 +129,14 @@ func TestPutBucketRequestPaymentRequesterRejected(t *testing.T) {
|
||||
func TestPutBucketOwnershipControlsRejectsRuleWithoutObjectOwnership(t *testing.T) {
|
||||
ownerID := AccountAdmin.Id
|
||||
s3a := &S3ApiServer{
|
||||
bucketRegistry: &BucketRegistry{
|
||||
metadataCache: map[string]*BucketMetaData{
|
||||
"b": {
|
||||
Name: "b",
|
||||
Owner: &s3.Owner{
|
||||
ID: &ownerID,
|
||||
},
|
||||
},
|
||||
},
|
||||
notFound: map[string]struct{}{},
|
||||
},
|
||||
bucketRegistry: NewBucketRegistry(nil),
|
||||
}
|
||||
s3a.bucketRegistry.setMetadataCache(&BucketMetaData{
|
||||
Name: "b",
|
||||
Owner: &s3.Owner{
|
||||
ID: &ownerID,
|
||||
},
|
||||
})
|
||||
body := `<OwnershipControls><Rule></Rule></OwnershipControls>`
|
||||
req := newBucketRequest(http.MethodPut, "b", "ownershipControls=", body)
|
||||
req.Header.Set(s3_constants.AmzAccountId, AccountAdmin.Id)
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
@@ -28,48 +27,27 @@ func (s3a *S3ApiServer) getStoredBucketLifecycleConfiguration(bucket string) ([]
|
||||
if errCode != s3err.ErrNone {
|
||||
return nil, "", false, errCode
|
||||
}
|
||||
if config.Entry == nil || config.Entry.Extended == nil {
|
||||
if len(config.LifecycleXML) == 0 {
|
||||
return nil, "", false, s3err.ErrNone
|
||||
}
|
||||
|
||||
lifecycleXML, found := config.Entry.Extended[bucketLifecycleConfigurationXMLKey]
|
||||
if !found || len(lifecycleXML) == 0 {
|
||||
return nil, "", false, s3err.ErrNone
|
||||
}
|
||||
transitionMinimumObjectSize := normalizeBucketLifecycleTransitionMinimumObjectSize(config.LifecycleTransitionMinSize)
|
||||
|
||||
transitionMinimumObjectSize := normalizeBucketLifecycleTransitionMinimumObjectSize(
|
||||
string(config.Entry.Extended[bucketLifecycleTransitionMinimumObjectSizeKey]),
|
||||
)
|
||||
|
||||
return append([]byte(nil), lifecycleXML...), transitionMinimumObjectSize, true, s3err.ErrNone
|
||||
return append([]byte(nil), config.LifecycleXML...), transitionMinimumObjectSize, true, s3err.ErrNone
|
||||
}
|
||||
|
||||
func (s3a *S3ApiServer) storeBucketLifecycleConfiguration(bucket string, lifecycleXML []byte, transitionMinimumObjectSize string) s3err.ErrorCode {
|
||||
return s3a.updateBucketConfig(bucket, func(config *BucketConfig) error {
|
||||
if config.Entry == nil {
|
||||
return fmt.Errorf("bucket %s is missing its filer entry", bucket)
|
||||
}
|
||||
if config.Entry.Extended == nil {
|
||||
config.Entry.Extended = make(map[string][]byte)
|
||||
}
|
||||
|
||||
config.Entry.Extended[bucketLifecycleConfigurationXMLKey] = append([]byte(nil), lifecycleXML...)
|
||||
config.Entry.Extended[bucketLifecycleTransitionMinimumObjectSizeKey] = []byte(
|
||||
normalizeBucketLifecycleTransitionMinimumObjectSize(transitionMinimumObjectSize),
|
||||
)
|
||||
|
||||
config.LifecycleXML = append([]byte(nil), lifecycleXML...)
|
||||
config.LifecycleTransitionMinSize = normalizeBucketLifecycleTransitionMinimumObjectSize(transitionMinimumObjectSize)
|
||||
return nil
|
||||
})
|
||||
}
|
||||
|
||||
func (s3a *S3ApiServer) clearStoredBucketLifecycleConfiguration(bucket string) s3err.ErrorCode {
|
||||
return s3a.updateBucketConfig(bucket, func(config *BucketConfig) error {
|
||||
if config.Entry == nil {
|
||||
return fmt.Errorf("bucket %s is missing its filer entry", bucket)
|
||||
}
|
||||
|
||||
delete(config.Entry.Extended, bucketLifecycleConfigurationXMLKey)
|
||||
delete(config.Entry.Extended, bucketLifecycleTransitionMinimumObjectSizeKey)
|
||||
config.LifecycleXML = nil
|
||||
config.LifecycleTransitionMinSize = ""
|
||||
return nil
|
||||
})
|
||||
}
|
||||
@@ -9,7 +9,6 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -23,13 +22,9 @@ func TestGetBucketLifecycleConfigurationHandlerUsesStoredLifecycleConfig(t *test
|
||||
s3a.option = &S3ApiServerOption{BucketsPath: "/buckets"}
|
||||
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{
|
||||
Extended: map[string][]byte{
|
||||
bucketLifecycleConfigurationXMLKey: []byte(lifecycleXML),
|
||||
bucketLifecycleTransitionMinimumObjectSizeKey: []byte("varies_by_storage_class"),
|
||||
},
|
||||
},
|
||||
Name: bucket,
|
||||
LifecycleXML: []byte(lifecycleXML),
|
||||
LifecycleTransitionMinSize: "varies_by_storage_class",
|
||||
})
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/"+bucket+"?lifecycle", nil)
|
||||
@@ -51,12 +46,8 @@ func TestGetBucketLifecycleConfigurationHandlerDefaultsTransitionMinimumObjectSi
|
||||
s3a.option = &S3ApiServerOption{BucketsPath: "/buckets"}
|
||||
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{
|
||||
Extended: map[string][]byte{
|
||||
bucketLifecycleConfigurationXMLKey: []byte(lifecycleXML),
|
||||
},
|
||||
},
|
||||
Name: bucket,
|
||||
LifecycleXML: []byte(lifecycleXML),
|
||||
})
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/"+bucket+"?lifecycle", nil)
|
||||
@@ -76,10 +67,7 @@ func TestPutBucketLifecycleConfigurationHandlerRejectsOversizedBody(t *testing.T
|
||||
s3a := newTestS3ApiServerWithMemoryIAM(t, nil)
|
||||
s3a.option = &S3ApiServerOption{BucketsPath: "/buckets"}
|
||||
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{},
|
||||
})
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{Name: bucket})
|
||||
|
||||
req := httptest.NewRequest(http.MethodPut, "/"+bucket+"?lifecycle", strings.NewReader(strings.Repeat("x", maxBucketLifecycleConfigurationSize+1)))
|
||||
req = mux.SetURLVars(req, map[string]string{"bucket": bucket})
|
||||
@@ -97,10 +85,7 @@ func TestPutBucketLifecycleConfigurationHandlerMapsReadErrorsToInvalidRequest(t
|
||||
s3a := newTestS3ApiServerWithMemoryIAM(t, nil)
|
||||
s3a.option = &S3ApiServerOption{BucketsPath: "/buckets"}
|
||||
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{
|
||||
Name: bucket,
|
||||
Entry: &filer_pb.Entry{},
|
||||
})
|
||||
s3a.bucketConfigCache.Set(bucket, &BucketConfig{Name: bucket})
|
||||
|
||||
req := httptest.NewRequest(http.MethodPut, "/"+bucket+"?lifecycle", nil)
|
||||
req = mux.SetURLVars(req, map[string]string{"bucket": bucket})
|
||||
@@ -124,4 +109,3 @@ func (f failingReadCloser) Read(_ []byte) (int, error) {
|
||||
func (f failingReadCloser) Close() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -614,13 +614,8 @@ func TestPostPolicyBucketHandler_PolicyViolationReturns403(t *testing.T) {
|
||||
}
|
||||
// Pre-populate the bucket registry so validateTableBucketObjectPath sees
|
||||
// a non-table bucket without needing a live filer connection.
|
||||
s3a.bucketRegistry = &BucketRegistry{
|
||||
metadataCache: map[string]*BucketMetaData{
|
||||
testBucket: {Name: testBucket, IsTableBucket: false},
|
||||
},
|
||||
notFound: make(map[string]struct{}),
|
||||
s3a: s3a,
|
||||
}
|
||||
s3a.bucketRegistry = NewBucketRegistry(s3a)
|
||||
s3a.bucketRegistry.setMetadataCache(&BucketMetaData{Name: testBucket, IsTableBucket: false})
|
||||
|
||||
now := time.Now().UTC()
|
||||
amzDate := now.Format(iso8601Format)
|
||||
|
||||
@@ -194,7 +194,7 @@ func BenchmarkLifecycleTTLResolver_Resolve_FiveRulesNoMatch(b *testing.B) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPopulateBucketConfigDerivedFields_RefreshesLifecycleTTL(t *testing.T) {
|
||||
func TestNewBucketConfigFromEntry_RefreshesLifecycleTTL(t *testing.T) {
|
||||
// Regression: storeBucketLifecycleConfiguration only updates
|
||||
// Entry.Extended; if the cache-refresh path doesn't re-derive
|
||||
// LifecycleTTL, an Add → Update → Delete dance would leave a stale
|
||||
@@ -205,26 +205,27 @@ func TestPopulateBucketConfigDerivedFields_RefreshesLifecycleTTL(t *testing.T) {
|
||||
xmlAdd := []byte(`<LifecycleConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ID>r</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><Expiration><Days>7</Days></Expiration></Rule></LifecycleConfiguration>`)
|
||||
xmlReplace := []byte(`<LifecycleConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ID>r</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><Expiration><Days>30</Days></Expiration></Rule></LifecycleConfiguration>`)
|
||||
|
||||
cfg := &BucketConfig{Name: "bk", Entry: &filer_pb.Entry{Extended: map[string][]byte{
|
||||
ext := map[string][]byte{
|
||||
s3_constants.ExtLifecycleTtlFastPathKey: []byte("true"),
|
||||
}}}
|
||||
}
|
||||
entry := &filer_pb.Entry{Extended: ext}
|
||||
|
||||
// 1) No XML yet → no resolver.
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
cfg := s.newBucketConfigFromEntry("bk", entry)
|
||||
if cfg.LifecycleTTL != nil {
|
||||
t.Fatalf("no XML must yield nil resolver, got %v", cfg.LifecycleTTL)
|
||||
}
|
||||
|
||||
// 2) Add: 7d.
|
||||
cfg.Entry.Extended[bucketLifecycleConfigurationXMLKey] = xmlAdd
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
ext[bucketLifecycleConfigurationXMLKey] = xmlAdd
|
||||
cfg = s.newBucketConfigFromEntry("bk", entry)
|
||||
if got := cfg.LifecycleTTL.Resolve("logs/foo", 1); got != 7*86400 {
|
||||
t.Fatalf("after add, want 7d, got %d", got)
|
||||
}
|
||||
|
||||
// 3) Replace: 30d. The previous resolver must NOT linger.
|
||||
cfg.Entry.Extended[bucketLifecycleConfigurationXMLKey] = xmlReplace
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
ext[bucketLifecycleConfigurationXMLKey] = xmlReplace
|
||||
cfg = s.newBucketConfigFromEntry("bk", entry)
|
||||
if got := cfg.LifecycleTTL.Resolve("logs/foo", 1); got != 30*86400 {
|
||||
t.Fatalf("after replace, want 30d, got %d (stale resolver?)", got)
|
||||
}
|
||||
@@ -232,14 +233,14 @@ func TestPopulateBucketConfigDerivedFields_RefreshesLifecycleTTL(t *testing.T) {
|
||||
// 4) Delete: nil resolver. The most dangerous regression — leaving
|
||||
// the old resolver here would keep stamping irreversible volume
|
||||
// TTL onto writes after the policy was removed.
|
||||
delete(cfg.Entry.Extended, bucketLifecycleConfigurationXMLKey)
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
delete(ext, bucketLifecycleConfigurationXMLKey)
|
||||
cfg = s.newBucketConfigFromEntry("bk", entry)
|
||||
if cfg.LifecycleTTL != nil {
|
||||
t.Fatalf("after delete, want nil resolver, got %v", cfg.LifecycleTTL)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPopulateBucketConfigDerivedFields_ObjectLockTreatedAsVersioned(t *testing.T) {
|
||||
func TestNewBucketConfigFromEntry_ObjectLockTreatedAsVersioned(t *testing.T) {
|
||||
// Object Lock requires versioning, so a bucket with ObjectLock but
|
||||
// no explicit Versioning header is still effectively versioned —
|
||||
// volume TTL would expire all noncurrent versions as a unit. The
|
||||
@@ -247,15 +248,11 @@ func TestPopulateBucketConfigDerivedFields_ObjectLockTreatedAsVersioned(t *testi
|
||||
// treat ObjectLockConfig != nil as versioned.
|
||||
s := &S3ApiServer{}
|
||||
xml := []byte(`<LifecycleConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ID>r</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><Expiration><Days>7</Days></Expiration></Rule></LifecycleConfiguration>`)
|
||||
cfg := &BucketConfig{
|
||||
Name: "bk",
|
||||
Entry: &filer_pb.Entry{Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte(s3_constants.ObjectLockEnabled),
|
||||
s3_constants.ExtLifecycleTtlFastPathKey: []byte("true"),
|
||||
bucketLifecycleConfigurationXMLKey: xml,
|
||||
}},
|
||||
}
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
cfg := s.newBucketConfigFromEntry("bk", &filer_pb.Entry{Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte(s3_constants.ObjectLockEnabled),
|
||||
s3_constants.ExtLifecycleTtlFastPathKey: []byte("true"),
|
||||
bucketLifecycleConfigurationXMLKey: xml,
|
||||
}})
|
||||
if cfg.ObjectLockConfig == nil {
|
||||
t.Fatal("test setup: ObjectLockConfig should be parsed")
|
||||
}
|
||||
@@ -264,7 +261,7 @@ func TestPopulateBucketConfigDerivedFields_ObjectLockTreatedAsVersioned(t *testi
|
||||
}
|
||||
}
|
||||
|
||||
func TestPopulateBucketConfigDerivedFields_TtlFastPathOptIn(t *testing.T) {
|
||||
func TestNewBucketConfigFromEntry_TtlFastPathOptIn(t *testing.T) {
|
||||
// The fast path is opt-in per bucket: lifecycle XML alone must not
|
||||
// stamp volume TTL on writes. Only the explicit flag builds the
|
||||
// resolver.
|
||||
@@ -272,17 +269,18 @@ func TestPopulateBucketConfigDerivedFields_TtlFastPathOptIn(t *testing.T) {
|
||||
xml := []byte(`<LifecycleConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ID>r</ID><Status>Enabled</Status><Filter><Prefix>logs/</Prefix></Filter><Expiration><Days>7</Days></Expiration></Rule></LifecycleConfiguration>`)
|
||||
|
||||
// XML present but flag off → nil resolver (worker drives expiration).
|
||||
cfg := &BucketConfig{Name: "bk", Entry: &filer_pb.Entry{Extended: map[string][]byte{
|
||||
ext := map[string][]byte{
|
||||
bucketLifecycleConfigurationXMLKey: xml,
|
||||
}}}
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
}
|
||||
entry := &filer_pb.Entry{Extended: ext}
|
||||
cfg := s.newBucketConfigFromEntry("bk", entry)
|
||||
if cfg.LifecycleTTL != nil {
|
||||
t.Fatalf("fast path off must yield nil resolver, got %v", cfg.LifecycleTTL)
|
||||
}
|
||||
|
||||
// Flag on → resolver applies the rule.
|
||||
cfg.Entry.Extended[s3_constants.ExtLifecycleTtlFastPathKey] = []byte("true")
|
||||
s.populateBucketConfigDerivedFields(cfg)
|
||||
ext[s3_constants.ExtLifecycleTtlFastPathKey] = []byte("true")
|
||||
cfg = s.newBucketConfigFromEntry("bk", entry)
|
||||
if got := cfg.LifecycleTTL.Resolve("logs/foo", 1); got != 7*86400 {
|
||||
t.Fatalf("fast path on, want 7d, got %d", got)
|
||||
}
|
||||
|
||||
@@ -18,18 +18,15 @@ func TestVeeamObjectLockBugFix(t *testing.T) {
|
||||
// The old code would immediately return NoSuchObjectLockConfiguration
|
||||
// The new code correctly checks if Object Lock is enabled before returning an error
|
||||
|
||||
bucketConfig := &BucketConfig{
|
||||
Name: "test-bucket",
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Extended: nil, // This is the key - no extended attributes
|
||||
},
|
||||
entry := &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Extended: nil, // This is the key - no extended attributes
|
||||
}
|
||||
|
||||
// Simulate the isObjectLockEnabledForBucket logic
|
||||
enabled := false
|
||||
if bucketConfig.Entry.Extended != nil {
|
||||
if enabledBytes, exists := bucketConfig.Entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
if entry.Extended != nil {
|
||||
if enabledBytes, exists := entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
enabled = string(enabledBytes) == s3_constants.ObjectLockEnabled || string(enabledBytes) == "true"
|
||||
}
|
||||
}
|
||||
@@ -41,20 +38,17 @@ func TestVeeamObjectLockBugFix(t *testing.T) {
|
||||
t.Run("Fix verification: bucket with Object Lock enabled via boolean flag", func(t *testing.T) {
|
||||
// This verifies the fix works when Object Lock is enabled via boolean flag
|
||||
|
||||
bucketConfig := &BucketConfig{
|
||||
entry := &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte("true"),
|
||||
},
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte("true"),
|
||||
},
|
||||
}
|
||||
|
||||
// Simulate the isObjectLockEnabledForBucket logic
|
||||
enabled := false
|
||||
if bucketConfig.Entry.Extended != nil {
|
||||
if enabledBytes, exists := bucketConfig.Entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
if entry.Extended != nil {
|
||||
if enabledBytes, exists := entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
enabled = string(enabledBytes) == s3_constants.ObjectLockEnabled || string(enabledBytes) == "true"
|
||||
}
|
||||
}
|
||||
@@ -66,20 +60,17 @@ func TestVeeamObjectLockBugFix(t *testing.T) {
|
||||
t.Run("Fix verification: bucket with Object Lock enabled via Enabled constant", func(t *testing.T) {
|
||||
// Test using the s3_constants.ObjectLockEnabled constant
|
||||
|
||||
bucketConfig := &BucketConfig{
|
||||
entry := &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: "test-bucket",
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte(s3_constants.ObjectLockEnabled),
|
||||
},
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.ExtObjectLockEnabledKey: []byte(s3_constants.ObjectLockEnabled),
|
||||
},
|
||||
}
|
||||
|
||||
// Simulate the isObjectLockEnabledForBucket logic
|
||||
enabled := false
|
||||
if bucketConfig.Entry.Extended != nil {
|
||||
if enabledBytes, exists := bucketConfig.Entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
if entry.Extended != nil {
|
||||
if enabledBytes, exists := entry.Extended[s3_constants.ExtObjectLockEnabledKey]; exists {
|
||||
enabled = string(enabledBytes) == s3_constants.ObjectLockEnabled || string(enabledBytes) == "true"
|
||||
}
|
||||
}
|
||||
|
||||
@@ -139,16 +139,13 @@ func TestSetObjectOwnerFromRequest(t *testing.T) {
|
||||
// Setup bucket registry with mock behavior
|
||||
if !tt.bucketRegistryNil {
|
||||
// Create a minimal BucketRegistry with overridden GetBucketMetadata
|
||||
s3a.bucketRegistry = &BucketRegistry{
|
||||
metadataCache: map[string]*BucketMetaData{},
|
||||
notFound: map[string]struct{}{},
|
||||
}
|
||||
s3a.bucketRegistry = NewBucketRegistry(nil)
|
||||
|
||||
// Pre-populate the cache with test metadata
|
||||
if tt.bucketMetadata != nil && tt.bucketMetadataError == s3err.ErrNone {
|
||||
s3a.bucketRegistry.metadataCache["test-bucket"] = tt.bucketMetadata
|
||||
s3a.bucketRegistry.metadataCache.Add("test-bucket", tt.bucketMetadata)
|
||||
} else if tt.bucketMetadataError == s3err.ErrNoSuchBucket {
|
||||
s3a.bucketRegistry.notFound["test-bucket"] = struct{}{}
|
||||
s3a.bucketRegistry.notFound.Add("test-bucket", struct{}{})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user