From 17af32f3ffb010e7547827da84647f8c1e5df990 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Thu, 2 Jul 2026 21:11:41 -0700 Subject: [PATCH] s3: paginate ListBuckets and serve it from a bucket owner index (#10214) * s3: paginate ListBuckets with max-buckets, continuation-token, and prefix ListBuckets buffered every bucket entry into one slice and one XML body, which falls over with very large bucket counts. Page through the filer listing instead, cap each response at 10000 buckets like AWS, and honor max-buckets, prefix, and an opaque keyset continuation-token. * s3: maintain a bucket owner index under /buckets/.system/owners Map each bucket owner to its buckets as zero-length entries at /buckets/.system/owners//, with Crtime mirroring the bucket's creation time. The bucket handlers write the index synchronously, the /buckets metadata subscription reconciles changes made elsewhere (weed shell, other gateways, direct filer operations), and a startup backfill indexes pre-existing buckets before writing a ready marker. Owner names are path-escaped so no identity name can escape the index directory. * s3: serve ListBuckets from the bucket owner index Once the owner index is ready, non-admin identities list their owned buckets straight from it, merged with any buckets their legacy actions name explicitly, so ListBuckets costs O(own buckets) instead of a scan of the global /buckets directory. Admins, identities with a bare List grant or wildcard action patterns, and policy-authorized identities whose grants cannot be enumerated keep the paged scan; policy-routed identities get their owned buckets, matching AWS ListBuckets returning only the caller's buckets. * s3: keep dot-prefixed names under /buckets out of bucket surfaces Dot-prefixed entries (.system) can never be valid bucket names, so refuse to resolve them as buckets and skip them in the shell bucket listing, matching what ListBuckets and the admin UI already do. * test: cover ListBuckets pagination and the owner index end to end * s3: fail closed on a nil identity when routing ListBuckets * s3: decide the IAM authorization mechanism in one place VerifyActionPermission and the ListBuckets owner-index routing each re-derived the session-token / attached-policy / legacy-actions split; extract the decision so the two cannot drift. * s3: heal the owner index on concurrent bucket recreation too The mkdir-lost-the-race path answers BucketAlreadyOwnedByYou just like the up-front existence check, so give it the same index repair. * s3: drop owner-index records for buckets deleted during backfill A bucket removed between the backfill reading its page and writing the index record became a permanent phantom in its owner's listing: the delete's own cleanup ran before the record existed. After indexing each page, re-list the same name range and remove records whose bucket is gone; deletes landing after the re-list find the record and remove it themselves. * s3: add ContinuationToken and Prefix to the ListBuckets schema Keep AmazonS3.xsd aligned with the generated ListAllMyBucketsResult so a regeneration does not drop the pagination fields. --- test/s3/normal/s3_list_buckets_test.go | 191 +++++++++++ weed/s3api/AmazonS3.xsd | 2 + weed/s3api/auth_credentials.go | 113 +++++-- weed/s3api/auth_credentials_subscribe.go | 1 + weed/s3api/bucket_paths.go | 5 + weed/s3api/s3api_bucket_handlers.go | 212 +++++++++--- weed/s3api/s3api_bucket_handlers_list_test.go | 123 +++++++ weed/s3api/s3api_bucket_owner_index.go | 316 ++++++++++++++++++ weed/s3api/s3api_bucket_owner_index_test.go | 162 +++++++++ weed/s3api/s3api_server.go | 7 + weed/s3api/s3api_xsd_generated.go | 4 + weed/s3api/s3err/s3api_errors.go | 12 + weed/shell/command_s3_bucket_list.go | 3 +- 13 files changed, 1087 insertions(+), 64 deletions(-) create mode 100644 test/s3/normal/s3_list_buckets_test.go create mode 100644 weed/s3api/s3api_bucket_handlers_list_test.go create mode 100644 weed/s3api/s3api_bucket_owner_index.go create mode 100644 weed/s3api/s3api_bucket_owner_index_test.go 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)