diff --git a/weed/filer/filer.go b/weed/filer/filer.go index ad3a186a0..39b7f75b0 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -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) + } } /* diff --git a/weed/filer/filer_inode_test.go b/weed/filer/filer_inode_test.go index 41c8883a9..42990a6fb 100644 --- a/weed/filer/filer_inode_test.go +++ b/weed/filer/filer_inode_test.go @@ -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(``), + "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() diff --git a/weed/filer/filer_lazy_remote_test.go b/weed/filer/filer_lazy_remote_test.go index f9b3442e0..15cfe8c40 100644 --- a/weed/filer/filer_lazy_remote_test.go +++ b/weed/filer/filer_lazy_remote_test.go @@ -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 } diff --git a/weed/shell/command_s3_bucket_create.go b/weed/shell/command_s3_bucket_create.go index 5fcc7e7fb..07f0eab1d 100644 --- a/weed/shell/command_s3_bucket_create.go +++ b/weed/shell/command_s3_bucket_create.go @@ -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) diff --git a/weed/shell/command_s3_bucket_create_test.go b/weed/shell/command_s3_bucket_create_test.go new file mode 100644 index 000000000..c409585fa --- /dev/null +++ b/weed/shell/command_s3_bucket_create_test.go @@ -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 +}