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/<owner>/<bucket>, 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.
This commit is contained in:
Chris Lu
2026-07-02 21:11:41 -07:00
committed by GitHub
parent c4f0b12a9a
commit 17af32f3ff
13 changed files with 1087 additions and 64 deletions
+191
View File
@@ -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
}
+2
View File
@@ -578,6 +578,8 @@
<xsd:sequence>
<xsd:element name="Owner" type="tns:CanonicalUser"/>
<xsd:element name="Buckets" type="tns:ListAllMyBucketsList"/>
<xsd:element name="ContinuationToken" type="xsd:string" minOccurs="0"/>
<xsd:element name="Prefix" type="xsd:string" minOccurs="0"/>
</xsd:sequence>
</xsd:complexType>
+88 -25
View File
@@ -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 /
+1
View File
@@ -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)
+5
View File
@@ -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)
}
+174 -38
View File
@@ -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
}
@@ -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")
}
}
+316
View File
@@ -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:
//
// <BucketsPath>/.system/owners/<escaped owner>/<bucket>
//
// 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
}
+162
View File
@@ -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)
}
}
+7
View File
@@ -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()
+4
View File
@@ -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 {
+12
View File
@@ -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",
+2 -1
View File
@@ -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)