shell: refuse s3.bucket.create on an existing bucket (#11455)

* shell: refuse s3.bucket.create on an existing bucket

CreateEntry without o_excl replaces the bucket entry, dropping every
extended attribute: lifecycle configuration, owner, versioning and the
irreversible Object Lock flag. Send o_excl so a re-run fails with
'bucket already exists' instead of silently resetting the bucket.

* filer: fail exclusive creates when the lookup itself fails

CreateEntry discards FindEntry errors, so an o_excl create hitting a
transient store failure would take the insert path and upsert over the
entry it was meant to preserve. Propagate the lookup error when o_excl
is set; non-exclusive creates keep their existing semantics.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* shell: test s3.bucket.create requests an exclusive create

Exercises the command end to end through a fake filer gRPC server and
asserts the OExcl flag reaches the wire along with the already-exists
error path.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* shell: synchronize captured requests and assert the exact bucket error

The fake filer records CreateEntry requests on the gRPC server goroutine,
so reads need the same mutex; the test also now checks for the exact
"bucket my-bucket already exists" message rather than any error that
mentions existence.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
Chris Lu
2026-09-25 22:04:54 +08:00
committed by GitHub
co-authored by Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
parent 317e756b9a
commit b750853c42
5 changed files with 170 additions and 4 deletions
+7 -1
View File
@@ -265,7 +265,13 @@ func (f *Filer) CreateEntry(ctx context.Context, entry *Entry, existing *Entry,
oldEntry := existing
if oldEntry == nil {
oldEntry, _ = f.FindEntry(ctx, entry.FullPath)
var findErr error
oldEntry, findErr = f.FindEntry(ctx, entry.FullPath)
if o_excl && findErr != nil && !errors.Is(findErr, filer_pb.ErrNotFound) {
// An exclusive create cannot decide whether the path exists when
// the lookup itself failed; proceeding would upsert over it.
return fmt.Errorf("find entry %s: %w", entry.FullPath, findErr)
}
}
/*
+53
View File
@@ -2,11 +2,13 @@ package filer
import (
"context"
"errors"
"os"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -101,6 +103,57 @@ func TestCreateEntryAssignsInodesToAutoCreatedParents(t *testing.T) {
}
}
func TestCreateEntryOExclPreservesExistingEntry(t *testing.T) {
f, store := newTestFilerWithStubStore()
original := &Entry{
FullPath: util.FullPath("/buckets/my-bucket"),
Attr: Attr{Mode: os.ModeDir | 0o777},
Extended: map[string][]byte{
"lifecycle": []byte(`<LifecycleConfiguration/>`),
"owner": []byte("alice"),
},
}
require.NoError(t, store.InsertEntry(context.Background(), original))
replacement := &Entry{
FullPath: util.FullPath("/buckets/my-bucket"),
Attr: Attr{Mode: os.ModeDir | 0o777},
}
err := f.CreateEntry(context.Background(), replacement, original, true, false, nil, false, f.MaxFilenameLength)
require.ErrorIs(t, err, filer_pb.ErrEntryAlreadyExists)
stored, findErr := store.FindEntry(context.Background(), original.FullPath)
require.NoError(t, findErr)
assert.Equal(t, original.Extended, stored.Extended)
}
func TestCreateEntryOExclFailsOnLookupError(t *testing.T) {
f, store := newTestFilerWithStubStore()
original := &Entry{
FullPath: util.FullPath("/buckets/my-bucket"),
Attr: Attr{Mode: os.ModeDir | 0o777},
Extended: map[string][]byte{"owner": []byte("alice")},
}
require.NoError(t, store.InsertEntry(context.Background(), original))
// a failed lookup must not masquerade as "not found": without the check
// the insert path would upsert over the stored bucket entry
store.findErr = errors.New("transient store failure")
err := f.CreateEntry(context.Background(), &Entry{
FullPath: util.FullPath("/buckets/my-bucket"),
Attr: Attr{Mode: os.ModeDir | 0o777},
}, nil, true, false, nil, false, f.MaxFilenameLength)
require.Error(t, err)
assert.NotErrorIs(t, err, filer_pb.ErrEntryAlreadyExists)
store.findErr = nil
stored, findErr := store.FindEntry(context.Background(), original.FullPath)
require.NoError(t, findErr)
assert.Equal(t, original.Extended, stored.Extended)
}
func TestUpdateEntryPreservesExistingInode(t *testing.T) {
f, store := newTestFilerWithStubStore()
+4
View File
@@ -34,6 +34,7 @@ type stubFilerStore struct {
entries map[string]*Entry
kv map[string][]byte
insertErr error
findErr error
deleteErrByPath map[string]error
}
@@ -174,6 +175,9 @@ func (s *stubFilerStore) UpdateEntry(_ context.Context, entry *Entry) error {
func (s *stubFilerStore) FindEntry(_ context.Context, p util.FullPath) (*Entry, error) {
s.mu.Lock()
defer s.mu.Unlock()
if s.findErr != nil {
return nil, s.findErr
}
if e, ok := s.entries[string(p)]; ok {
return e, nil
}
+9 -3
View File
@@ -2,6 +2,7 @@ package shell
import (
"context"
"errors"
"flag"
"fmt"
"io"
@@ -112,11 +113,16 @@ func (c *commandS3BucketCreate) Do(args []string, commandEnv *CommandEnv, writer
entry.Extended[s3_constants.ExtObjectLockEnabledKey] = []byte(s3_constants.ObjectLockEnabled)
}
if _, err := client.CreateEntry(context.Background(), &filer_pb.CreateEntryRequest{
createErr := filer_pb.CreateEntry(context.Background(), client, &filer_pb.CreateEntryRequest{
Directory: filerBucketsPath,
Entry: entry,
}); err != nil {
return err
OExcl: true,
})
if errors.Is(createErr, filer_pb.ErrEntryAlreadyExists) {
return fmt.Errorf("bucket %s already exists", *bucketName)
}
if createErr != nil {
return createErr
}
fmt.Fprintln(writer, "created bucket", *bucketName)
@@ -0,0 +1,97 @@
package shell
import (
"bytes"
"context"
"fmt"
"net"
"os"
"path/filepath"
"sync"
"testing"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
)
type bucketCreateTestFilerServer struct {
filer_pb.UnimplementedSeaweedFilerServer
mu sync.Mutex
createReqs []*filer_pb.CreateEntryRequest
}
func (s *bucketCreateTestFilerServer) GetFilerConfiguration(context.Context, *filer_pb.GetFilerConfigurationRequest) (*filer_pb.GetFilerConfigurationResponse, error) {
return &filer_pb.GetFilerConfigurationResponse{DirBuckets: "/buckets"}, nil
}
func (s *bucketCreateTestFilerServer) CreateEntry(_ context.Context, req *filer_pb.CreateEntryRequest) (*filer_pb.CreateEntryResponse, error) {
s.mu.Lock()
defer s.mu.Unlock()
s.createReqs = append(s.createReqs, req)
return &filer_pb.CreateEntryResponse{
ErrorCode: filer_pb.FilerError_ENTRY_ALREADY_EXISTS,
Error: "entry already exists",
}, nil
}
func (s *bucketCreateTestFilerServer) createRequests() []*filer_pb.CreateEntryRequest {
s.mu.Lock()
defer s.mu.Unlock()
return append([]*filer_pb.CreateEntryRequest(nil), s.createReqs...)
}
func TestS3BucketCreateRequestsExclusiveCreate(t *testing.T) {
filerServer := &bucketCreateTestFilerServer{}
commandEnv, cleanup := newBucketCreateTestCommandEnv(t, filerServer)
defer cleanup()
var output bytes.Buffer
err := (&commandS3BucketCreate{}).Do([]string{"-name", "my-bucket"}, commandEnv, &output)
require.EqualError(t, err, "bucket my-bucket already exists")
reqs := filerServer.createRequests()
require.Len(t, reqs, 1)
req := reqs[0]
assert.True(t, req.OExcl, "bucket create must not replace an existing entry")
assert.Equal(t, "/buckets", req.Directory)
assert.Equal(t, "my-bucket", req.Entry.Name)
}
func newBucketCreateTestCommandEnv(t *testing.T, filerServer filer_pb.SeaweedFilerServer) (*CommandEnv, func()) {
t.Helper()
socketDir, err := os.MkdirTemp("", "swbucket-")
require.NoError(t, err)
t.Cleanup(func() { _ = os.RemoveAll(socketDir) })
socketPath := filepath.Join(socketDir, "filer.sock")
listener, err := net.Listen("unix", socketPath)
require.NoError(t, err)
grpcServer := grpc.NewServer()
filer_pb.RegisterSeaweedFilerServer(grpcServer, filerServer)
go func() {
_ = grpcServer.Serve(listener)
}()
grpcPort := 48000 + os.Getpid()%1000
pb.RegisterLocalGrpcSocket("127.0.0.1", grpcPort, socketPath)
cleanup := func() {
grpcServer.Stop()
_ = listener.Close()
}
return &CommandEnv{
option: &ShellOptions{
FilerAddress: pb.ServerAddress(fmt.Sprintf("127.0.0.1:8888.%d", grpcPort)),
GrpcDialOption: grpc.WithTransportCredentials(insecure.NewCredentials()),
Directory: "/",
},
}, cleanup
}