diff --git a/test/s3/normal/s3_list_buckets_test.go b/test/s3/normal/s3_list_buckets_test.go new file mode 100644 index 000000000..4f85f538d --- /dev/null +++ b/test/s3/normal/s3_list_buckets_test.go @@ -0,0 +1,191 @@ +package example + +import ( + "context" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + v2config "github.com/aws/aws-sdk-go-v2/config" + v2creds "github.com/aws/aws-sdk-go-v2/credentials" + s3v2 "github.com/aws/aws-sdk-go-v2/service/s3" + s3v2types "github.com/aws/aws-sdk-go-v2/service/s3/types" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestListBucketsPaginationAndOwnerIndex exercises ListBuckets pagination +// (max-buckets, continuation-token, prefix) and the owner-index path that +// serves non-admin identities without scanning all buckets. +func TestListBucketsPaginationAndOwnerIndex(t *testing.T) { + if testing.Short() { + t.Skip("Skipping integration test in short mode") + } + + s3Config := `{ + "identities": [ + {"name": "admin", "credentials": [{"accessKey": "admin", "secretKey": "admin"}], "actions": ["Admin"]}, + {"name": "alice", "credentials": [{"accessKey": "alice", "secretKey": "alice_secret"}], + "actions": ["Admin:alice-b1", "Admin:alice-b2", "Admin:alice-b3"]}, + {"name": "carol", "credentials": [{"accessKey": "carol", "secretKey": "carol_secret"}], + "actions": ["List:alice-b1"]}, + {"name": "bob", "credentials": [{"accessKey": "bob", "secretKey": "bob_secret"}], + "actions": ["Read:alice-b1"]} + ] + }` + configPath := filepath.Join(t.TempDir(), "s3.json") + require.NoError(t, os.WriteFile(configPath, []byte(s3Config), 0644)) + + cluster, err := startMiniCluster(t, "-s3.config="+configPath) + require.NoError(t, err) + defer cluster.Stop() + + ctx := context.Background() + newClient := func(key, secret string) *s3v2.Client { + cfg, err := v2config.LoadDefaultConfig(ctx, + v2config.WithRegion("us-east-1"), + v2config.WithCredentialsProvider(v2creds.NewStaticCredentialsProvider(key, secret, "")), + ) + require.NoError(t, err) + return s3v2.NewFromConfig(cfg, func(o *s3v2.Options) { + o.BaseEndpoint = aws.String(cluster.s3Endpoint) + o.UsePathStyle = true + }) + } + admin := newClient("admin", "admin") + alice := newClient("alice", "alice_secret") + carol := newClient("carol", "carol_secret") + bob := newClient("bob", "bob_secret") + + // Non-admin listings use the owner index once the backfill marker exists. + markerURL := fmt.Sprintf("http://127.0.0.1:%d/buckets/.system/owners/.complete", cluster.filerPort) + require.Eventually(t, func() bool { + resp, err := http.Get(markerURL) + if err != nil { + return false + } + defer resp.Body.Close() + return resp.StatusCode == http.StatusOK + }, 15*time.Second, 200*time.Millisecond, "owner index backfill marker") + + aliceBuckets := []string{"alice-b1", "alice-b2", "alice-b3"} + for _, b := range aliceBuckets { + _, err := alice.CreateBucket(ctx, &s3v2.CreateBucketInput{Bucket: aws.String(b)}) + require.NoError(t, err, "alice create %s", b) + } + var adminBuckets []string + for i := 0; i < 5; i++ { + name := fmt.Sprintf("admin-b%d", i) + adminBuckets = append(adminBuckets, name) + _, err := admin.CreateBucket(ctx, &s3v2.CreateBucketInput{Bucket: aws.String(name)}) + require.NoError(t, err, "admin create %s", name) + } + + t.Run("OwnerSeesOwnBuckets", func(t *testing.T) { + out, err := alice.ListBuckets(ctx, &s3v2.ListBucketsInput{}) + require.NoError(t, err) + assert.Equal(t, aliceBuckets, bucketNames(out.Buckets)) + assert.Nil(t, out.ContinuationToken) + }) + + t.Run("OwnerPagination", func(t *testing.T) { + page1, err := alice.ListBuckets(ctx, &s3v2.ListBucketsInput{MaxBuckets: aws.Int32(2)}) + require.NoError(t, err) + assert.Equal(t, aliceBuckets[:2], bucketNames(page1.Buckets)) + require.NotNil(t, page1.ContinuationToken) + + page2, err := alice.ListBuckets(ctx, &s3v2.ListBucketsInput{ + MaxBuckets: aws.Int32(2), + ContinuationToken: page1.ContinuationToken, + }) + require.NoError(t, err) + assert.Equal(t, aliceBuckets[2:], bucketNames(page2.Buckets)) + assert.Nil(t, page2.ContinuationToken) + }) + + t.Run("OwnerPrefix", func(t *testing.T) { + out, err := alice.ListBuckets(ctx, &s3v2.ListBucketsInput{Prefix: aws.String("alice-b2")}) + require.NoError(t, err) + assert.Equal(t, []string{"alice-b2"}, bucketNames(out.Buckets)) + }) + + t.Run("GrantedListVisible", func(t *testing.T) { + out, err := carol.ListBuckets(ctx, &s3v2.ListBucketsInput{}) + require.NoError(t, err) + assert.Equal(t, []string{"alice-b1"}, bucketNames(out.Buckets)) + }) + + t.Run("ReadGrantNotVisible", func(t *testing.T) { + out, err := bob.ListBuckets(ctx, &s3v2.ListBucketsInput{}) + require.NoError(t, err) + assert.Empty(t, bucketNames(out.Buckets)) + }) + + t.Run("AdminPaginatesAll", func(t *testing.T) { + var all []string + var token *string + for pages := 0; ; pages++ { + require.Less(t, pages, 10, "runaway pagination") + out, err := admin.ListBuckets(ctx, &s3v2.ListBucketsInput{MaxBuckets: aws.Int32(3), ContinuationToken: token}) + require.NoError(t, err) + require.LessOrEqual(t, len(out.Buckets), 3) + all = append(all, bucketNames(out.Buckets)...) + if out.ContinuationToken == nil { + break + } + token = out.ContinuationToken + } + assert.Equal(t, append(append([]string{}, adminBuckets...), aliceBuckets...), all) + }) + + t.Run("AdminPrefix", func(t *testing.T) { + out, err := admin.ListBuckets(ctx, &s3v2.ListBucketsInput{Prefix: aws.String("admin-")}) + require.NoError(t, err) + assert.Equal(t, adminBuckets, bucketNames(out.Buckets)) + }) + + t.Run("InvalidMaxBuckets", func(t *testing.T) { + _, err := admin.ListBuckets(ctx, &s3v2.ListBucketsInput{MaxBuckets: aws.Int32(0)}) + require.Error(t, err) + }) + + t.Run("IndexEntriesOnFiler", func(t *testing.T) { + indexURL := fmt.Sprintf("http://127.0.0.1:%d/buckets/.system/owners/alice/?limit=100", cluster.filerPort) + req, err := http.NewRequest(http.MethodGet, indexURL, nil) + require.NoError(t, err) + req.Header.Set("Accept", "application/json") + resp, err := http.DefaultClient.Do(req) + require.NoError(t, err) + defer resp.Body.Close() + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + for _, b := range aliceBuckets { + assert.Contains(t, string(body), b) + } + }) + + t.Run("DeleteRemovesFromListing", func(t *testing.T) { + _, err := alice.DeleteBucket(ctx, &s3v2.DeleteBucketInput{Bucket: aws.String("alice-b2")}) + require.NoError(t, err) + out, err := alice.ListBuckets(ctx, &s3v2.ListBucketsInput{}) + require.NoError(t, err) + assert.Equal(t, []string{"alice-b1", "alice-b3"}, bucketNames(out.Buckets)) + }) + + t.Run("HiddenSystemBucket", func(t *testing.T) { + _, err := admin.HeadBucket(ctx, &s3v2.HeadBucketInput{Bucket: aws.String(".system")}) + require.Error(t, err) + }) +} + +func bucketNames(buckets []s3v2types.Bucket) (names []string) { + for _, b := range buckets { + names = append(names, aws.ToString(b.Name)) + } + return names +} diff --git a/weed/s3api/AmazonS3.xsd b/weed/s3api/AmazonS3.xsd index 8a0136b44..df3de4087 100644 --- a/weed/s3api/AmazonS3.xsd +++ b/weed/s3api/AmazonS3.xsd @@ -578,6 +578,8 @@ + + diff --git a/weed/s3api/auth_credentials.go b/weed/s3api/auth_credentials.go index 0226aea11..427036660 100644 --- a/weed/s3api/auth_credentials.go +++ b/weed/s3api/auth_credentials.go @@ -2359,26 +2359,30 @@ func (iam *IdentityAccessManagement) isActionExplicitlyDeniedByIAM(r *http.Reque return denied } -// VerifyActionPermission checks if the identity is allowed to perform the action on the resource. -// It handles both traditional identities (via Actions) and IAM/STS identities (via Policy). -func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, identity *Identity, action Action, bucket, object string) s3err.ErrorCode { - // Fail closed if identity is nil - if identity == nil { - glog.V(3).Infof("VerifyActionPermission called with nil identity for action %s on %s/%s", action, bucket, object) - return s3err.ErrAccessDenied - } +// authorizationRoute is the mechanism that decides a request/identity pair's +// permissions: the IAM integration, locally attached IAM policies, the +// identity's legacy Actions, or nothing at all. +type authorizationRoute int - // Traditional identities (with Actions from -s3.config) use legacy auth, - // JWT/STS identities (no Actions or having a session token) use IAM authorization. - // IMPORTANT: We MUST prioritize IAM authorization for any request with a session token - // to ensure that session policies are correctly enforced. +const ( + authorizeViaIAMIntegration authorizationRoute = iota + authorizeViaAttachedPolicies + authorizeViaLegacyActions + authorizeDenied +) + +// authorizationRoute picks the mechanism, so every caller routes identically. +// Traditional identities (with Actions from -s3.config) use legacy auth, +// JWT/STS identities (no Actions or having a session token) use IAM +// authorization. A request with a session token must go through the IAM +// integration so session policies are enforced. +func (iam *IdentityAccessManagement) authorizationRoute(r *http.Request, identity *Identity) authorizationRoute { hasSessionToken := r.Header.Get(s3_constants.SeaweedFSSessionTokenHeader) != "" || r.Header.Get("X-Amz-Security-Token") != "" || r.URL.Query().Get("X-Amz-Security-Token") != "" iam.m.RLock() - userGroupNames := iam.userGroups[identity.Name] groupsHavePolicies := false - for _, gn := range userGroupNames { + for _, gn := range iam.userGroups[identity.Name] { if g, ok := iam.groups[gn]; ok && !g.Disabled && len(g.PolicyNames) > 0 { groupsHavePolicies = true break @@ -2388,29 +2392,88 @@ func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, ide hasAttachedPolicies := len(identity.PolicyNames) > 0 || groupsHavePolicies if (len(identity.Actions) == 0 || hasSessionToken || hasAttachedPolicies) && iam.iamIntegration != nil { - return iam.authorizeWithIAM(r, identity, action, bucket, object) + return authorizeViaIAMIntegration + } + if hasAttachedPolicies { + return authorizeViaAttachedPolicies + } + if len(identity.Actions) > 0 { + return authorizeViaLegacyActions + } + return authorizeDenied +} + +// VerifyActionPermission checks if the identity is allowed to perform the action on the resource. +// It handles both traditional identities (via Actions) and IAM/STS identities (via Policy). +func (iam *IdentityAccessManagement) VerifyActionPermission(r *http.Request, identity *Identity, action Action, bucket, object string) s3err.ErrorCode { + // Fail closed if identity is nil + if identity == nil { + glog.V(3).Infof("VerifyActionPermission called with nil identity for action %s on %s/%s", action, bucket, object) + return s3err.ErrAccessDenied } - // Attached IAM policies are authoritative for IAM users. The legacy Actions - // field is a lossy projection that cannot represent deny statements, - // conditions, or fine-grained action differences such as PutObject vs - // DeleteObject. - if hasAttachedPolicies { + switch iam.authorizationRoute(r, identity) { + case authorizeViaIAMIntegration: + return iam.authorizeWithIAM(r, identity, action, bucket, object) + case authorizeViaAttachedPolicies: + // Attached IAM policies are authoritative for IAM users. The legacy Actions + // field is a lossy projection that cannot represent deny statements, + // conditions, or fine-grained action differences such as PutObject vs + // DeleteObject. if iam.evaluateIAMPolicies(r, identity, action, bucket, object) { return s3err.ErrNone } return s3err.ErrAccessDenied - } - - // Traditional actions-based authorization from static S3 config. - if len(identity.Actions) > 0 { + case authorizeViaLegacyActions: if !identity.CanDo(action, bucket, object) { return s3err.ErrAccessDenied } return s3err.ErrNone + default: + return s3err.ErrAccessDenied + } +} + +// canListBucketsFromOwnerIndex reports whether ListBuckets for this identity +// can be served from the bucket owner index instead of scanning /buckets, and +// if so which bucket names its legacy actions may grant beyond ownership. +// +// Admins and identities whose legacy actions can match arbitrary buckets (a +// bare "List" grant, or any wildcard pattern) need the full scan. Identities +// authorized through IAM policies get their owned buckets only, matching the +// AWS behavior of ListBuckets returning the account's buckets. +func (iam *IdentityAccessManagement) canListBucketsFromOwnerIndex(r *http.Request, identity *Identity) (ok bool, granted []string) { + // Fail closed on a nil identity: the scan path filters every bucket out + // without dereferencing it, while the index path would need its name. + if identity == nil || identity.isAdmin() { + return false, nil } - return s3err.ErrAccessDenied + // Identities authorized by IAM policies (or by nothing) cannot have their + // visible set enumerated; they get their owned buckets only. + if iam.authorizationRoute(r, identity) != authorizeViaLegacyActions { + return true, nil + } + + for _, a := range identity.Actions { + act := string(a) + if act == string(s3_constants.ACTION_LIST) { + return false, nil + } + if strings.ContainsAny(act, "*?") { + return false, nil + } + if colon := strings.Index(act, ":"); colon >= 0 { + bucket := act[colon+1:] + if slash := strings.Index(bucket, "/"); slash >= 0 { + bucket = bucket[:slash] + } + if bucket != "" { + granted = append(granted, bucket) + } + } + } + return true, granted } // AuthorizeCopySource verifies the caller is allowed to read the CopyObject / diff --git a/weed/s3api/auth_credentials_subscribe.go b/weed/s3api/auth_credentials_subscribe.go index d0862d8ed..38657a3ed 100644 --- a/weed/s3api/auth_credentials_subscribe.go +++ b/weed/s3api/auth_credentials_subscribe.go @@ -183,6 +183,7 @@ func (s3a *S3ApiServer) onCircuitBreakerConfigChange(dir string, oldEntry *filer // reload bucket metadata func (s3a *S3ApiServer) onBucketMetadataChange(dir string, oldEntry *filer_pb.Entry, newEntry *filer_pb.Entry) error { if dir == s3a.option.BucketsPath { + s3a.maintainBucketOwnerIndex(oldEntry, newEntry) if newEntry != nil { // Update bucket registry (existing functionality) s3a.bucketRegistry.LoadBucketMetadata(newEntry) diff --git a/weed/s3api/bucket_paths.go b/weed/s3api/bucket_paths.go index 44ad76747..f3febef4f 100644 --- a/weed/s3api/bucket_paths.go +++ b/weed/s3api/bucket_paths.go @@ -94,5 +94,10 @@ func (s3a *S3ApiServer) bucketExists(bucket string) (bool, error) { } func (s3a *S3ApiServer) getBucketEntry(bucket string) (*filer_pb.Entry, error) { + // Dot-prefixed entries under /buckets are internal (.system) and can never + // be valid bucket names; don't let them resolve as buckets. + if strings.HasPrefix(bucket, ".") { + return nil, filer_pb.ErrNotFound + } return s3a.getEntry(s3a.option.BucketsPath, bucket) } diff --git a/weed/s3api/s3api_bucket_handlers.go b/weed/s3api/s3api_bucket_handlers.go index 1ce21804c..08dae3d74 100644 --- a/weed/s3api/s3api_bucket_handlers.go +++ b/weed/s3api/s3api_bucket_handlers.go @@ -3,14 +3,16 @@ package s3api import ( "bytes" "context" + "encoding/base64" "encoding/json" "encoding/xml" "errors" "fmt" "io" - "math" "net/http" + "net/url" "sort" + "strconv" "strings" "time" @@ -62,56 +64,133 @@ func (s3a *S3ApiServer) ListBucketsHandler(w http.ResponseWriter, r *http.Reques } } - var response ListAllMyBucketsResult - - entries, _, err := s3a.list(s3a.option.BucketsPath, "", "", false, math.MaxInt32) - - if err != nil { - s3err.WriteErrorResponse(w, r, s3err.ErrInternalError) + maxBuckets, prefix, startAfter, errCode := getListBucketsArgs(r.URL.Query()) + if errCode != s3err.ErrNone { + s3err.WriteErrorResponse(w, r, errCode) return } - var listBuckets ListAllMyBucketsList - for _, entry := range entries { - if entry.IsDirectory { - if strings.HasPrefix(entry.Name, ".") { - continue - } - // Unauthenticated users should not see any buckets - if identity == nil { - continue - } - - // Check if bucket should be visible to this identity - // A bucket is visible if the user owns it OR has explicit permission to list it - isOwner := isBucketOwnedByIdentity(entry, identity) - - // Skip permission check if user is already the owner (optimization) - if !isOwner { - if errCode := s3a.iam.VerifyActionPermission(r, identity, s3_constants.ACTION_LIST, entry.Name, ""); errCode != s3err.ErrNone { - continue - } - } - - listBuckets.Bucket = append(listBuckets.Bucket, ListAllMyBucketsEntry{ - Name: entry.Name, - CreationDate: time.Unix(entry.Attributes.Crtime, 0).UTC(), - }) + var buckets []ListAllMyBucketsEntry + var nextToken string + // Unauthenticated users should not see any buckets + if identity != nil { + var err error + if fromIndex, granted := s3a.iam.canListBucketsFromOwnerIndex(r, identity); fromIndex && s3a.bucketOwnerIndexReady() { + buckets, nextToken, err = s3a.listBucketsFromOwnerIndex(r, identity, granted, prefix, startAfter, maxBuckets) + } else { + buckets, nextToken, err = s3a.scanVisibleBuckets(r, identity, prefix, startAfter, maxBuckets) + } + if err != nil { + glog.Errorf("ListBucketsHandler: %v", err) + s3err.WriteErrorResponse(w, r, s3err.ErrInternalError) + return } } - response = ListAllMyBucketsResult{ + response := ListAllMyBucketsResult{ Owner: CanonicalUser{ ID: identityId, DisplayName: identityId, }, - Buckets: listBuckets, + Buckets: ListAllMyBucketsList{Bucket: buckets}, + ContinuationToken: nextToken, + Prefix: prefix, } glog.V(3).Infof("ListBucketsHandler response: %+v", response) writeSuccessResponseXML(w, r, response) } +const ( + // maxBucketsPerPage caps a single ListBuckets response, matching AWS. + maxBucketsPerPage = 10000 + // bucketScanPageSize is the filer page size while scanning /buckets. + bucketScanPageSize = 1000 +) + +func getListBucketsArgs(values url.Values) (maxBuckets int, prefix string, startAfter string, errCode s3err.ErrorCode) { + maxBuckets = maxBucketsPerPage + if v := values.Get("max-buckets"); v != "" { + n, err := strconv.Atoi(v) + if err != nil || n < 1 || n > maxBucketsPerPage { + return 0, "", "", s3err.ErrInvalidMaxBuckets + } + maxBuckets = n + } + prefix = values.Get("prefix") + if v := values.Get("continuation-token"); v != "" { + name, err := decodeContinuationToken(v) + if err != nil { + return 0, "", "", s3err.ErrInvalidContinuationToken + } + startAfter = name + } + return maxBuckets, prefix, startAfter, s3err.ErrNone +} + +// The continuation token is the last returned bucket name, base64-encoded so +// clients treat it as opaque per the S3 API contract. +func encodeContinuationToken(bucket string) string { + return base64.RawURLEncoding.EncodeToString([]byte(bucket)) +} + +func decodeContinuationToken(token string) (string, error) { + name, err := base64.RawURLEncoding.DecodeString(token) + if err != nil || len(name) == 0 { + return "", fmt.Errorf("invalid continuation token %q", token) + } + return string(name), nil +} + +// bucketVisibleToIdentity reports whether ListBuckets should include the bucket: +// the identity owns it or has explicit permission to list it. +func (s3a *S3ApiServer) bucketVisibleToIdentity(r *http.Request, entry *filer_pb.Entry, identity *Identity) bool { + if isBucketOwnedByIdentity(entry, identity) { + return true + } + return s3a.iam.VerifyActionPermission(r, identity, s3_constants.ACTION_LIST, entry.Name, "") == s3err.ErrNone +} + +// scanVisibleBuckets pages through /buckets and collects up to maxBuckets +// entries visible to the identity, returning a continuation token when more +// remain. It never buffers the full bucket set. +func (s3a *S3ApiServer) scanVisibleBuckets(r *http.Request, identity *Identity, prefix, startAfter string, maxBuckets int) (buckets []ListAllMyBucketsEntry, nextToken string, err error) { + startFrom := startAfter + for len(buckets) <= maxBuckets { + entries, isLast, listErr := s3a.list(s3a.option.BucketsPath, prefix, startFrom, false, bucketScanPageSize) + if listErr != nil { + return nil, "", listErr + } + if len(entries) == 0 { + break + } + for _, entry := range entries { + startFrom = entry.Name + if !entry.IsDirectory || strings.HasPrefix(entry.Name, ".") { + continue + } + if !s3a.bucketVisibleToIdentity(r, entry, identity) { + continue + } + buckets = append(buckets, ListAllMyBucketsEntry{ + Name: entry.Name, + CreationDate: time.Unix(entry.Attributes.Crtime, 0).UTC(), + }) + if len(buckets) > maxBuckets { + break + } + } + if isLast || len(buckets) > maxBuckets { + break + } + } + if len(buckets) > maxBuckets { + buckets = buckets[:maxBuckets] + nextToken = encodeContinuationToken(buckets[maxBuckets-1].Name) + } + return buckets, nextToken, nil +} + // isBucketOwnedByIdentity checks if a bucket entry is owned by the given identity. // Returns true if the identity owns the bucket, false otherwise. // @@ -218,7 +297,11 @@ func (s3a *S3ApiServer) PutBucketHandler(w http.ResponseWriter, r *http.Request) // Bucket already exists: report whether the caller already owns it or the // name is taken / the request conflicts. if exist, err := s3a.exists(s3a.option.BucketsPath, bucket, true); err == nil && exist { - s3err.WriteErrorResponse(w, r, s3a.existingBucketError(r, bucket, currentIdentityId, requestHasACL)) + errCode := s3a.existingBucketError(r, bucket, currentIdentityId, requestHasACL) + if errCode == s3err.ErrBucketAlreadyOwnedByYou { + s3a.healBucketOwnerIndex(bucket, currentIdentityId) + } + s3err.WriteErrorResponse(w, r, errCode) return } @@ -241,7 +324,10 @@ func (s3a *S3ApiServer) PutBucketHandler(w http.ResponseWriter, r *http.Request) // Create the folder for bucket with all settings atomically // This ensures Object Lock configuration is set in the same CreateEntry call, // preventing race conditions where the bucket exists without Object Lock enabled + var bucketCrtime int64 if err := s3a.mkdir(s3a.option.BucketsPath, bucket, func(entry *filer_pb.Entry) { + bucketCrtime = entry.Attributes.Crtime + // Set bucket owner setBucketOwner(r)(entry) @@ -282,7 +368,11 @@ func (s3a *S3ApiServer) PutBucketHandler(w http.ResponseWriter, r *http.Request) // return the appropriate already-exists error instead of InternalError. if exist, checkErr := s3a.exists(s3a.option.BucketsPath, bucket, true); checkErr == nil && exist { glog.V(3).Infof("PutBucketHandler: bucket %s was created concurrently", bucket) - s3err.WriteErrorResponse(w, r, s3a.existingBucketError(r, bucket, currentIdentityId, requestHasACL)) + errCode := s3a.existingBucketError(r, bucket, currentIdentityId, requestHasACL) + if errCode == s3err.ErrBucketAlreadyOwnedByYou { + s3a.healBucketOwnerIndex(bucket, currentIdentityId) + } + s3err.WriteErrorResponse(w, r, errCode) return } glog.Errorf("PutBucketHandler mkdir: %v", err) @@ -306,10 +396,34 @@ func (s3a *S3ApiServer) PutBucketHandler(w http.ResponseWriter, r *http.Request) s3a.bucketConfigCache.RemoveNegativeCache(bucket) } + // Index the bucket under its owner; the metadata subscription would catch + // up eventually, but writing here keeps ListBuckets read-your-writes. + if currentIdentityId != "" { + if err := s3a.addBucketToOwnerIndex(currentIdentityId, bucket, bucketCrtime); err != nil { + glog.Warningf("PutBucketHandler: owner index add %s/%s: %v", currentIdentityId, bucket, err) + } + } + w.Header().Set("Location", "/"+bucket) writeSuccessResponseEmpty(w, r) } +// healBucketOwnerIndex re-creates the owner index entry for a bucket the +// caller owns, repairing holes left by crashes between bucket creation and +// index write. +func (s3a *S3ApiServer) healBucketOwnerIndex(bucket, identityId string) { + if identityId == "" { + return + } + config, errCode := s3a.getBucketConfig(bucket) + if errCode != s3err.ErrNone || bucketEntryOwner(config.Entry) != identityId { + return + } + if err := s3a.addBucketToOwnerIndex(identityId, bucket, config.Entry.Attributes.Crtime); err != nil { + glog.V(1).Infof("owner index heal %s/%s: %v", identityId, bucket, err) + } +} + func (s3a *S3ApiServer) DeleteBucketHandler(w http.ResponseWriter, r *http.Request) { bucket, _ := s3_constants.GetBucketAndObject(r) @@ -369,6 +483,12 @@ func (s3a *S3ApiServer) DeleteBucketHandler(w http.ResponseWriter, r *http.Reque return } + if owner := bucketEntryOwner(bucketConfig.Entry); owner != "" { + if err := s3a.removeBucketFromOwnerIndex(owner, bucket); err != nil { + glog.Warningf("DeleteBucketHandler: owner index remove %s/%s: %v", owner, bucket, err) + } + } + err = s3a.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error { deleteCollectionRequest := &filer_pb.DeleteCollectionRequest{ Collection: s3a.getCollectionName(bucket), @@ -589,7 +709,12 @@ func (s3a *S3ApiServer) autoCreateBucket(r *http.Request, bucket string) error { return fmt.Errorf("auto-create bucket %s: %w", bucket, ErrAutoCreatePermissionDenied) } - if err := s3a.mkdir(s3a.option.BucketsPath, bucket, setBucketOwner(r)); err != nil { + identityId := s3_constants.GetIdentityNameFromContext(r) + var bucketCrtime int64 + if err := s3a.mkdir(s3a.option.BucketsPath, bucket, func(entry *filer_pb.Entry) { + bucketCrtime = entry.Attributes.Crtime + setBucketOwner(r)(entry) + }); err != nil { // In case of a race condition where another request created the bucket // in the meantime, check for existence before returning an error. if exist, err2 := s3a.exists(s3a.option.BucketsPath, bucket, true); err2 != nil { @@ -607,6 +732,11 @@ func (s3a *S3ApiServer) autoCreateBucket(r *http.Request, bucket string) error { glog.Warningf("autoCreateBucket: failed to set owner for existing bucket %s: %v", bucket, updateErr) } else { glog.V(1).Infof("Set owner for existing bucket %s (created by concurrent request)", bucket) + if owner := bucketEntryOwner(entry); owner != "" { + if indexErr := s3a.addBucketToOwnerIndex(owner, bucket, entry.Attributes.Crtime); indexErr != nil { + glog.Warningf("autoCreateBucket: owner index add %s/%s: %v", owner, bucket, indexErr) + } + } } } } else { @@ -626,6 +756,12 @@ func (s3a *S3ApiServer) autoCreateBucket(r *http.Request, bucket string) error { s3a.bucketConfigCache.RemoveNegativeCache(bucket) } + if identityId != "" { + if err := s3a.addBucketToOwnerIndex(identityId, bucket, bucketCrtime); err != nil { + glog.Warningf("autoCreateBucket: owner index add %s/%s: %v", identityId, bucket, err) + } + } + glog.V(1).Infof("Auto-created bucket %s", bucket) return nil } diff --git a/weed/s3api/s3api_bucket_handlers_list_test.go b/weed/s3api/s3api_bucket_handlers_list_test.go new file mode 100644 index 000000000..aa6f43c6f --- /dev/null +++ b/weed/s3api/s3api_bucket_handlers_list_test.go @@ -0,0 +1,123 @@ +package s3api + +import ( + "net/http" + "net/url" + "strings" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" +) + +func TestGetListBucketsArgs(t *testing.T) { + tests := []struct { + name string + query string + wantMaxBuckets int + wantPrefix string + wantStartAfter string + wantErr s3err.ErrorCode + }{ + {name: "defaults", query: "", wantMaxBuckets: maxBucketsPerPage, wantErr: s3err.ErrNone}, + {name: "max-buckets", query: "max-buckets=25", wantMaxBuckets: 25, wantErr: s3err.ErrNone}, + {name: "max-buckets upper bound", query: "max-buckets=10000", wantMaxBuckets: 10000, wantErr: s3err.ErrNone}, + {name: "max-buckets zero", query: "max-buckets=0", wantErr: s3err.ErrInvalidMaxBuckets}, + {name: "max-buckets negative", query: "max-buckets=-1", wantErr: s3err.ErrInvalidMaxBuckets}, + {name: "max-buckets too large", query: "max-buckets=10001", wantErr: s3err.ErrInvalidMaxBuckets}, + {name: "max-buckets not a number", query: "max-buckets=abc", wantErr: s3err.ErrInvalidMaxBuckets}, + {name: "prefix", query: "prefix=team-", wantMaxBuckets: maxBucketsPerPage, wantPrefix: "team-", wantErr: s3err.ErrNone}, + {name: "continuation token", query: "continuation-token=" + url.QueryEscape(encodeContinuationToken("bucket-42")), + wantMaxBuckets: maxBucketsPerPage, wantStartAfter: "bucket-42", wantErr: s3err.ErrNone}, + {name: "bad continuation token", query: "continuation-token=%25%25", wantErr: s3err.ErrInvalidContinuationToken}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + values, err := url.ParseQuery(tt.query) + if err != nil { + t.Fatalf("parse query: %v", err) + } + maxBuckets, prefix, startAfter, errCode := getListBucketsArgs(values) + if errCode != tt.wantErr { + t.Fatalf("errCode = %v, want %v", errCode, tt.wantErr) + } + if errCode != s3err.ErrNone { + return + } + if maxBuckets != tt.wantMaxBuckets { + t.Errorf("maxBuckets = %d, want %d", maxBuckets, tt.wantMaxBuckets) + } + if prefix != tt.wantPrefix { + t.Errorf("prefix = %q, want %q", prefix, tt.wantPrefix) + } + if startAfter != tt.wantStartAfter { + t.Errorf("startAfter = %q, want %q", startAfter, tt.wantStartAfter) + } + }) + } +} + +func TestCanListBucketsFromOwnerIndex(t *testing.T) { + iam := &IdentityAccessManagement{} + r := func() *http.Request { + req, _ := http.NewRequest(http.MethodGet, "/", nil) + return req + } + + tests := []struct { + name string + identity *Identity + wantOk bool + wantGranted []string + }{ + {name: "nil identity scans", identity: nil, wantOk: false}, + {name: "admin scans", identity: &Identity{Name: "root", Actions: []Action{"Admin"}}, wantOk: false}, + {name: "bare List scans", identity: &Identity{Name: "u", Actions: []Action{"List"}}, wantOk: false}, + {name: "wildcard scans", identity: &Identity{Name: "u", Actions: []Action{"List:team-*"}}, wantOk: false}, + {name: "no actions is owned-only", identity: &Identity{Name: "u"}, wantOk: true}, + {name: "attached policies is owned-only", identity: &Identity{Name: "u", Actions: []Action{"List:b1"}, PolicyNames: []string{"p"}}, wantOk: true}, + {name: "named grants enumerate", identity: &Identity{Name: "u", Actions: []Action{"Read:b1", "List:b2", "Admin:b3/prefix", "Write"}}, + wantOk: true, wantGranted: []string{"b1", "b2", "b3"}}, + // grants come back in action order; resolveGrantedBuckets sorts and dedups + {name: "unsorted grants keep action order", identity: &Identity{Name: "u", Actions: []Action{"List:c1", "List:a1", "List:b1"}}, + wantOk: true, wantGranted: []string{"c1", "a1", "b1"}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ok, granted := iam.canListBucketsFromOwnerIndex(r(), tt.identity) + if ok != tt.wantOk { + t.Fatalf("ok = %v, want %v", ok, tt.wantOk) + } + if strings.Join(granted, ",") != strings.Join(tt.wantGranted, ",") { + t.Errorf("granted = %v, want %v", granted, tt.wantGranted) + } + }) + } + + t.Run("session token routes to policies when integrated", func(t *testing.T) { + integrated := &IdentityAccessManagement{iamIntegration: &S3IAMIntegration{}} + req := r() + req.Header.Set("X-Amz-Security-Token", "tok") + ok, granted := integrated.canListBucketsFromOwnerIndex(req, &Identity{Name: "u", Actions: []Action{"List:b1"}}) + if !ok || granted != nil { + t.Errorf("ok = %v granted = %v, want owned-only", ok, granted) + } + }) +} + +func TestContinuationTokenRoundTrip(t *testing.T) { + for _, name := range []string{"a", "bucket-42", "with.dots-and-dashes", "0123456789"} { + decoded, err := decodeContinuationToken(encodeContinuationToken(name)) + if err != nil { + t.Fatalf("decode(encode(%q)): %v", name, err) + } + if decoded != name { + t.Errorf("round trip %q -> %q", name, decoded) + } + } + if _, err := decodeContinuationToken(""); err == nil { + t.Error("empty token should not decode") + } + if _, err := decodeContinuationToken("!!!"); err == nil { + t.Error("non-base64 token should not decode") + } +} diff --git a/weed/s3api/s3api_bucket_owner_index.go b/weed/s3api/s3api_bucket_owner_index.go new file mode 100644 index 000000000..db8dbc0a8 --- /dev/null +++ b/weed/s3api/s3api_bucket_owner_index.go @@ -0,0 +1,316 @@ +package s3api + +import ( + "net/http" + "net/url" + "slices" + "strings" + "time" + + "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/util" +) + +// The owner index maps bucket owner -> owned buckets so ListBuckets can serve +// an identity's buckets without scanning the global /buckets directory: +// +// /.system/owners// +// +// Index entries are zero-length files whose Crtime mirrors the bucket's +// creation time. The bucket handlers write the index synchronously, the +// /buckets metadata subscription reconciles changes made elsewhere (weed +// shell, another gateway, direct filer operations), and a startup backfill +// indexes pre-existing buckets. The .complete marker is written once the +// backfill has finished; until it exists readers keep scanning /buckets. +const ( + s3SystemFolder = ".system" + bucketOwnerIndexFolder = "owners" + bucketOwnerIndexReadyMarker = ".complete" +) + +func (s3a *S3ApiServer) bucketOwnerIndexPath() string { + return s3a.option.BucketsPath + "/" + s3SystemFolder + "/" + bucketOwnerIndexFolder +} + +func (s3a *S3ApiServer) bucketOwnerDir(owner string) string { + return s3a.bucketOwnerIndexPath() + "/" + escapeOwnerName(owner) +} + +// escapeOwnerName makes an owner id safe to use as a single directory name. +// url.PathEscape escapes "/", so no owner name can traverse out of the index. +// It leaves dots alone, so a leading dot is encoded by hand: "." and ".." +// must not resolve as path segments, and other dot-names would collide with +// index-internal entries like the ready marker. PathEscape never emits a +// literal "%2E" itself, so the encoding stays collision-free. +func escapeOwnerName(owner string) string { + escaped := url.PathEscape(owner) + if strings.HasPrefix(escaped, ".") { + escaped = "%2E" + escaped[1:] + } + return escaped +} + +func bucketEntryOwner(entry *filer_pb.Entry) string { + if entry == nil || entry.Extended == nil { + return "" + } + return string(entry.Extended[s3_constants.AmzIdentityId]) +} + +func (s3a *S3ApiServer) addBucketToOwnerIndex(owner, bucket string, crtime int64) error { + return s3a.mkFile(s3a.bucketOwnerDir(owner), bucket, nil, func(entry *filer_pb.Entry) { + entry.Attributes.Crtime = crtime + }) +} + +func (s3a *S3ApiServer) removeBucketFromOwnerIndex(owner, bucket string) error { + return s3a.rm(s3a.bucketOwnerDir(owner), bucket, false, false) +} + +// maintainBucketOwnerIndex applies a /buckets metadata event to the owner +// index, covering bucket changes made outside this gateway's handlers. Every +// operation is idempotent, so overlap with the handlers' synchronous writes +// and with other gateways is harmless. +func (s3a *S3ApiServer) maintainBucketOwnerIndex(oldEntry, newEntry *filer_pb.Entry) { + var oldOwner, newOwner string + if oldEntry != nil && (!oldEntry.IsDirectory || strings.HasPrefix(oldEntry.Name, ".")) { + oldEntry = nil + } + if newEntry != nil && (!newEntry.IsDirectory || strings.HasPrefix(newEntry.Name, ".")) { + newEntry = nil + } + if oldEntry != nil { + oldOwner = bucketEntryOwner(oldEntry) + } + if newEntry != nil { + newOwner = bucketEntryOwner(newEntry) + } + + unchanged := oldEntry != nil && newEntry != nil && oldEntry.Name == newEntry.Name && oldOwner == newOwner + if unchanged { + return + } + if oldOwner != "" { + if err := s3a.removeBucketFromOwnerIndex(oldOwner, oldEntry.Name); err != nil { + glog.V(1).Infof("owner index: remove %s/%s: %v", oldOwner, oldEntry.Name, err) + } + } + if newOwner != "" { + if err := s3a.addBucketToOwnerIndex(newOwner, newEntry.Name, newEntry.Attributes.Crtime); err != nil { + glog.V(1).Infof("owner index: add %s/%s: %v", newOwner, newEntry.Name, err) + } + } +} + +// bucketOwnerIndexReady reports whether the backfill marker exists, caching a +// positive answer for the life of the process. +func (s3a *S3ApiServer) bucketOwnerIndexReady() bool { + if s3a.ownerIndexReady.Load() { + return true + } + if exists, err := s3a.exists(s3a.bucketOwnerIndexPath(), bucketOwnerIndexReadyMarker, false); err == nil && exists { + s3a.ownerIndexReady.Store(true) + return true + } + return false +} + +// startBucketOwnerIndexBackfill brings the owner index up to date with the +// buckets that existed before it was introduced, then writes the ready +// marker. Multiple gateways may race the backfill; every step is idempotent. +func (s3a *S3ApiServer) startBucketOwnerIndexBackfill() { + util.RetryUntil("bucketOwnerIndexBackfill", func() error { + if s3a.bucketOwnerIndexReady() { + return nil + } + return s3a.backfillBucketOwnerIndex() + }, func(err error) bool { + glog.V(1).Infof("bucket owner index backfill: %v", err) + return true + }) +} + +// listBucketsFromOwnerIndex serves ListBuckets from the owner index: the +// identity's owned buckets merged with the buckets its legacy actions name +// explicitly. Cost is proportional to the identity's own buckets, not to the +// global bucket count. +func (s3a *S3ApiServer) listBucketsFromOwnerIndex(r *http.Request, identity *Identity, grantedNames []string, prefix, startAfter string, maxBuckets int) ([]ListAllMyBucketsEntry, string, error) { + granted := s3a.resolveGrantedBuckets(r, identity, grantedNames, prefix, startAfter) + var listOwned bucketPageLister + if identity.Name != "" { + ownerDir := s3a.bucketOwnerDir(identity.Name) + listOwned = func(startFrom string) ([]*filer_pb.Entry, bool, error) { + return s3a.list(ownerDir, prefix, startFrom, false, bucketScanPageSize) + } + } + return listMergedBuckets(listOwned, granted, startAfter, maxBuckets) +} + +// resolveGrantedBuckets maps action-granted bucket names to listing entries, +// dropping duplicates, names outside the requested page, buckets that no +// longer exist, and names the identity cannot actually list. +func (s3a *S3ApiServer) resolveGrantedBuckets(r *http.Request, identity *Identity, names []string, prefix, startAfter string) (granted []ListAllMyBucketsEntry) { + slices.Sort(names) + names = slices.Compact(names) + for _, name := range names { + if name <= startAfter || !strings.HasPrefix(name, prefix) || strings.HasPrefix(name, ".") { + continue + } + config, errCode := s3a.getBucketConfig(name) + if errCode != s3err.ErrNone { + if errCode != s3err.ErrNoSuchBucket { + glog.V(1).Infof("owner index: resolve granted bucket %s: %v", name, errCode) + } + continue + } + if !s3a.bucketVisibleToIdentity(r, config.Entry, identity) { + continue + } + granted = append(granted, ListAllMyBucketsEntry{ + Name: name, + CreationDate: time.Unix(config.Entry.Attributes.Crtime, 0).UTC(), + }) + } + return granted +} + +type bucketPageLister func(startFrom string) (entries []*filer_pb.Entry, isLast bool, err error) + +// listMergedBuckets merges the sorted owned-bucket index stream with the +// pre-filtered, sorted granted entries, returning up to maxBuckets and a +// continuation token when more remain. +func listMergedBuckets(listOwned bucketPageLister, granted []ListAllMyBucketsEntry, startAfter string, maxBuckets int) (buckets []ListAllMyBucketsEntry, nextToken string, err error) { + gi := 0 + startFrom := startAfter + if listOwned != nil { + owned: + for len(buckets) <= maxBuckets { + entries, isLast, listErr := listOwned(startFrom) + if listErr != nil { + return nil, "", listErr + } + if len(entries) == 0 { + break + } + for _, entry := range entries { + startFrom = entry.Name + if entry.IsDirectory || strings.HasPrefix(entry.Name, ".") { + continue + } + for gi < len(granted) && granted[gi].Name < entry.Name { + buckets = append(buckets, granted[gi]) + gi++ + if len(buckets) > maxBuckets { + break owned + } + } + if gi < len(granted) && granted[gi].Name == entry.Name { + gi++ // also owned; keep the index entry + } + buckets = append(buckets, ListAllMyBucketsEntry{ + Name: entry.Name, + CreationDate: time.Unix(entry.Attributes.Crtime, 0).UTC(), + }) + if len(buckets) > maxBuckets { + break owned + } + } + if isLast { + break + } + } + } + for gi < len(granted) && len(buckets) <= maxBuckets { + buckets = append(buckets, granted[gi]) + gi++ + } + if len(buckets) > maxBuckets { + buckets = buckets[:maxBuckets] + nextToken = encodeContinuationToken(buckets[maxBuckets-1].Name) + } + return buckets, nextToken, nil +} + +func (s3a *S3ApiServer) backfillBucketOwnerIndex() error { + startFrom := "" + for { + pageStart := startFrom + entries, isLast, err := s3a.list(s3a.option.BucketsPath, "", startFrom, false, bucketScanPageSize) + if err != nil { + return err + } + if len(entries) == 0 { + break + } + var indexed []*filer_pb.Entry + for _, entry := range entries { + startFrom = entry.Name + if !entry.IsDirectory || strings.HasPrefix(entry.Name, ".") { + continue + } + if owner := bucketEntryOwner(entry); owner != "" { + if err := s3a.addBucketToOwnerIndex(owner, entry.Name, entry.Attributes.Crtime); err != nil { + return err + } + indexed = append(indexed, entry) + } + } + // A bucket deleted while this page was in flight would leave a permanent + // phantom record: the index write lands after the delete already tried to + // clean up. Re-list the page and drop records for buckets that vanished; + // a delete after this point sees the record and removes it itself. + if err := s3a.dropVanishedFromOwnerIndex(pageStart, indexed); err != nil { + return err + } + if isLast { + break + } + } + if err := s3a.mkFile(s3a.bucketOwnerIndexPath(), bucketOwnerIndexReadyMarker, nil, nil); err != nil { + return err + } + s3a.ownerIndexReady.Store(true) + glog.V(0).Infof("bucket owner index ready") + return nil +} + +func (s3a *S3ApiServer) dropVanishedFromOwnerIndex(startFrom string, indexed []*filer_pb.Entry) error { + if len(indexed) == 0 { + return nil + } + last := indexed[len(indexed)-1].Name + present := make(map[string]bool, len(indexed)) + from := startFrom +scan: + for { + entries, isLast, err := s3a.list(s3a.option.BucketsPath, "", from, false, bucketScanPageSize) + if err != nil { + return err + } + if len(entries) == 0 { + break + } + for _, entry := range entries { + if entry.Name > last { + break scan + } + from = entry.Name + present[entry.Name] = true + } + if isLast { + break + } + } + for _, entry := range indexed { + if present[entry.Name] { + continue + } + if err := s3a.removeBucketFromOwnerIndex(bucketEntryOwner(entry), entry.Name); err != nil { + return err + } + } + return nil +} diff --git a/weed/s3api/s3api_bucket_owner_index_test.go b/weed/s3api/s3api_bucket_owner_index_test.go new file mode 100644 index 000000000..bfc8e8c34 --- /dev/null +++ b/weed/s3api/s3api_bucket_owner_index_test.go @@ -0,0 +1,162 @@ +package s3api + +import ( + "strings" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" +) + +func TestEscapeOwnerName(t *testing.T) { + tests := []struct { + owner string + want string + }{ + {"alice", "alice"}, + {"arn:aws:iam::123:user/bob", "arn:aws:iam::123:user%2Fbob"}, + {".", "%2E"}, + {"..", "%2E."}, + {".complete", "%2Ecomplete"}, + {"../../../etc", "%2E.%2F..%2F..%2Fetc"}, + {"%2E", "%252E"}, + } + seen := make(map[string]string) + for _, tt := range tests { + got := escapeOwnerName(tt.owner) + if got != tt.want { + t.Errorf("escapeOwnerName(%q) = %q, want %q", tt.owner, got, tt.want) + } + // No result may resolve to a different directory level or collide + // with index-internal dot-names. + if strings.ContainsRune(got, '/') || strings.HasPrefix(got, ".") { + t.Errorf("escapeOwnerName(%q) = %q is not a safe single path segment", tt.owner, got) + } + if prev, dup := seen[got]; dup { + t.Errorf("escapeOwnerName collision: %q and %q both map to %q", prev, tt.owner, got) + } + seen[got] = tt.owner + } +} + +func ownedPages(pages ...[]string) bucketPageLister { + // Serves a static sorted stream of owned bucket names, honoring startFrom. + var all []string + for _, p := range pages { + all = append(all, p...) + } + return func(startFrom string) ([]*filer_pb.Entry, bool, error) { + var out []*filer_pb.Entry + for _, name := range all { + if name <= startFrom { + continue + } + out = append(out, &filer_pb.Entry{ + Name: name, + Attributes: &filer_pb.FuseAttributes{Crtime: 42}, + }) + if len(out) == 2 { // small pages to exercise paging + return out, false, nil + } + } + return out, true, nil + } +} + +func names(buckets []ListAllMyBucketsEntry) []string { + var out []string + for _, b := range buckets { + out = append(out, b.Name) + } + return out +} + +func TestListMergedBuckets(t *testing.T) { + granted := func(n ...string) (out []ListAllMyBucketsEntry) { + for _, name := range n { + out = append(out, ListAllMyBucketsEntry{Name: name}) + } + return out + } + + t.Run("interleave and dedup", func(t *testing.T) { + buckets, token, err := listMergedBuckets(ownedPages([]string{"a", "c", "e"}), granted("b", "c", "f"), "", 10) + if err != nil { + t.Fatal(err) + } + want := []string{"a", "b", "c", "e", "f"} + if strings.Join(names(buckets), ",") != strings.Join(want, ",") { + t.Errorf("got %v, want %v", names(buckets), want) + } + if token != "" { + t.Errorf("unexpected token %q", token) + } + }) + + t.Run("truncation and resume", func(t *testing.T) { + buckets, token, err := listMergedBuckets(ownedPages([]string{"a", "c", "e"}), granted("b", "f"), "", 3) + if err != nil { + t.Fatal(err) + } + if strings.Join(names(buckets), ",") != "a,b,c" { + t.Errorf("page 1 = %v", names(buckets)) + } + startAfter, err := decodeContinuationToken(token) + if err != nil { + t.Fatalf("token: %v", err) + } + if startAfter != "c" { + t.Errorf("token resumes after %q, want c", startAfter) + } + // granted entries for page 2 are pre-filtered by startAfter upstream + buckets, token, err = listMergedBuckets(ownedPages([]string{"a", "c", "e"}), granted("f"), startAfter, 3) + if err != nil { + t.Fatal(err) + } + if strings.Join(names(buckets), ",") != "e,f" { + t.Errorf("page 2 = %v", names(buckets)) + } + if token != "" { + t.Errorf("unexpected token %q on final page", token) + } + }) + + t.Run("no owned stream", func(t *testing.T) { + buckets, token, err := listMergedBuckets(nil, granted("a", "b"), "", 1) + if err != nil { + t.Fatal(err) + } + if strings.Join(names(buckets), ",") != "a" || token == "" { + t.Errorf("got %v token %q", names(buckets), token) + } + }) + + t.Run("owned only", func(t *testing.T) { + buckets, token, err := listMergedBuckets(ownedPages([]string{"a", "b", "c", "d", "e"}), nil, "b", 2) + if err != nil { + t.Fatal(err) + } + if strings.Join(names(buckets), ",") != "c,d" { + t.Errorf("got %v", names(buckets)) + } + if startAfter, _ := decodeContinuationToken(token); startAfter != "d" { + t.Errorf("token resumes after %q, want d", startAfter) + } + }) +} + +func TestBucketEntryOwner(t *testing.T) { + if owner := bucketEntryOwner(nil); owner != "" { + t.Errorf("nil entry owner = %q", owner) + } + if owner := bucketEntryOwner(&filer_pb.Entry{Name: "b"}); owner != "" { + t.Errorf("ownerless entry owner = %q", owner) + } + entry := &filer_pb.Entry{ + Name: "b", + Extended: map[string][]byte{s3_constants.AmzIdentityId: []byte("alice")}, + } + if owner := bucketEntryOwner(entry); owner != "alice" { + t.Errorf("owner = %q, want alice", owner) + } +} diff --git a/weed/s3api/s3api_server.go b/weed/s3api/s3api_server.go index 77f650b5a..96028a3e4 100644 --- a/weed/s3api/s3api_server.go +++ b/weed/s3api/s3api_server.go @@ -12,6 +12,7 @@ import ( "slices" "strings" "sync" + "sync/atomic" "time" "github.com/gorilla/mux" @@ -111,6 +112,9 @@ type S3ApiServer struct { // is nil in this commit; a follow-up wires in an in-memory chunk cache. readerCache *filer.ReaderCache + // ownerIndexReady caches the presence of the owner index backfill marker. + ownerIndexReady atomic.Bool + versionsHealQueue *versionsHealQueue versionsReconcilerStop func() } @@ -452,6 +456,9 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl // Start bucket size metrics collection in background go s3ApiServer.startBucketSizeMetricsLoop(context.Background()) + // Bring the bucket owner index up to date with pre-existing buckets + go s3ApiServer.startBucketOwnerIndexBackfill() + // Start the versioning reconciler that drains stranded .versions/ // pointer-to-missing-file states without waiting for a client GET. s3ApiServer.versionsReconcilerStop = s3ApiServer.startVersioningReconciler() diff --git a/weed/s3api/s3api_xsd_generated.go b/weed/s3api/s3api_xsd_generated.go index 79300cf4f..acc8d78e7 100644 --- a/weed/s3api/s3api_xsd_generated.go +++ b/weed/s3api/s3api_xsd_generated.go @@ -1077,6 +1077,10 @@ type ListAllMyBucketsResult struct { XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ ListAllMyBucketsResult"` Owner CanonicalUser `xml:"Owner"` Buckets ListAllMyBucketsList `xml:"Buckets"` + // ContinuationToken is set when the listing is truncated; passing it back + // via ?continuation-token= resumes the listing after the last bucket returned. + ContinuationToken string `xml:"ContinuationToken,omitempty"` + Prefix string `xml:"Prefix,omitempty"` } type ListBucket struct { diff --git a/weed/s3api/s3err/s3api_errors.go b/weed/s3api/s3err/s3api_errors.go index e4281feab..6174ef61f 100644 --- a/weed/s3api/s3err/s3api_errors.go +++ b/weed/s3api/s3err/s3api_errors.go @@ -64,6 +64,8 @@ const ( ErrInvalidDigest ErrBadDigest ErrInvalidMaxKeys + ErrInvalidMaxBuckets + ErrInvalidContinuationToken ErrInvalidMaxUploads ErrInvalidMaxParts ErrInvalidMaxDeleteObjects @@ -225,6 +227,16 @@ var errorCodeResponse = map[ErrorCode]APIError{ Description: "Argument maxKeys must be an integer between 0 and 2147483647", HTTPStatusCode: http.StatusBadRequest, }, + ErrInvalidMaxBuckets: { + Code: "InvalidArgument", + Description: "Argument max-buckets must be an integer between 1 and 10000", + HTTPStatusCode: http.StatusBadRequest, + }, + ErrInvalidContinuationToken: { + Code: "InvalidArgument", + Description: "The continuation token provided is incorrect", + HTTPStatusCode: http.StatusBadRequest, + }, ErrInvalidMaxParts: { Code: "InvalidArgument", Description: "Argument max-parts must be an integer between 0 and 2147483647", diff --git a/weed/shell/command_s3_bucket_list.go b/weed/shell/command_s3_bucket_list.go index bb55fc013..94eb2f8ba 100644 --- a/weed/shell/command_s3_bucket_list.go +++ b/weed/shell/command_s3_bucket_list.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "math" + "strings" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" @@ -59,7 +60,7 @@ func (c *commandS3BucketList) Do(args []string, commandEnv *CommandEnv, writer i } err = filer_pb.List(context.Background(), commandEnv, filerBucketsPath, "", func(entry *filer_pb.Entry, isLast bool) error { - if !entry.IsDirectory { + if !entry.IsDirectory || strings.HasPrefix(entry.Name, ".") { return nil } collection := getCollectionName(commandEnv, entry.Name)