diff --git a/auth/acl.go b/auth/acl.go index 0b0bba31..be9d6871 100644 --- a/auth/acl.go +++ b/auth/acl.go @@ -246,7 +246,7 @@ func ParseACLOutput(data []byte, owner string) (GetBucketAclOutput, error) { }, nil } -func UpdateACL(input *PutBucketAclInput, acl ACL, iam IAMService, isAdmin bool) ([]byte, error) { +func UpdateACL(input *PutBucketAclInput, acl ACL, iam IAMService) ([]byte, error) { if input == nil { return nil, s3err.GetAPIError(s3err.ErrInvalidRequest) } diff --git a/backend/azure/azure.go b/backend/azure/azure.go index f6edd601..c567d67b 100644 --- a/backend/azure/azure.go +++ b/backend/azure/azure.go @@ -157,7 +157,7 @@ func (az *Azure) CreateBucket(ctx context.Context, input *s3.CreateBucketInput, string(keyOwnership): backend.GetPtrFromString(encodeBytes([]byte(input.ObjectOwnership))), } - acct, ok := ctx.Value("account").(auth.Account) + acct, ok := ctx.Value("bucket-owner").(auth.Account) if !ok { acct = auth.Account{} } diff --git a/backend/posix/posix.go b/backend/posix/posix.go index ee1eaa0e..060dac7c 100644 --- a/backend/posix/posix.go +++ b/backend/posix/posix.go @@ -369,7 +369,7 @@ func (p *Posix) HeadBucket(_ context.Context, input *s3.HeadBucketInput) (*s3.He } func (p *Posix) CreateBucket(ctx context.Context, input *s3.CreateBucketInput, acl []byte) error { - acct, ok := ctx.Value("account").(auth.Account) + acct, ok := ctx.Value("bucket-owner").(auth.Account) if !ok { acct = auth.Account{} } diff --git a/cmd/versitygw/admin.go b/cmd/versitygw/admin.go index dea58845..01eb9d51 100644 --- a/cmd/versitygw/admin.go +++ b/cmd/versitygw/admin.go @@ -19,16 +19,20 @@ import ( "crypto/sha256" "crypto/tls" "encoding/hex" + "encoding/json" "encoding/xml" + "errors" "fmt" "io" "net/http" "os" + "strings" "text/tabwriter" "time" "github.com/aws/aws-sdk-go-v2/aws" v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4" + "github.com/aws/aws-sdk-go-v2/service/s3/types" "github.com/aws/smithy-go" "github.com/urfave/cli/v2" "github.com/versity/versitygw/auth" @@ -169,6 +173,66 @@ func adminCommand() *cli.Command { Usage: "Lists all the gateway buckets and owners.", Action: listBuckets, }, + { + Name: "create-bucket", + Usage: "Create a new bucket with owner", + Action: createBucket, + Flags: []cli.Flag{ + &cli.StringFlag{ + Name: "owner", + Usage: "access key id of the bucket owner", + Required: true, + Aliases: []string{"o"}, + }, + &cli.StringFlag{ + Name: "bucket", + Usage: "bucket name", + Required: true, + }, + &cli.StringFlag{ + Name: "acl", + Usage: "canned ACL to apply to the bucket", + }, + &cli.StringFlag{ + Name: "grant-full-control", + Usage: "Allows grantee the read, write, read ACP, and write ACP permissions on the bucket.", + }, + &cli.StringFlag{ + Name: "grant-read", + Usage: "Allows grantee to list the objects in the bucket.", + }, + &cli.StringFlag{ + Name: "grant-read-acp", + Usage: "Allows grantee to read the bucket ACL.", + }, + &cli.StringFlag{ + Name: "grant-write", + Usage: `Allows grantee to create new objects in the bucket. + For the bucket and object owners of existing objects, also allows deletions and overwrites of those objects.`, + }, + &cli.StringFlag{ + Name: "grant-write-acp", + Usage: "Allows grantee to write the ACL for the applicable bucket.", + }, + &cli.StringFlag{ + Name: "create-bucket-configuration", + Usage: "bucket configuration (LocationConstraint, Tags)", + }, + &cli.BoolFlag{ + Name: "object-lock-enabled-for-bucket", + Usage: "enable object lock for the bucket", + }, + &cli.BoolFlag{ + Name: "no-object-lock-enabled-for-bucket", + Usage: "disable object lock for the bucket", + }, + &cli.StringFlag{ + Name: "object-ownership", + Usage: "bucket object ownership setting", + Value: "", + }, + }, + }, }, Flags: []cli.Flag{ // TODO: create a configuration file for this @@ -442,6 +506,246 @@ func listUsers(ctx *cli.Context) error { return nil } +type createBucketInput struct { + LocationConstraint *string + Tags []types.Tag +} + +// parseCreateBucketPayload parses the +func parseCreateBucketPayload(input string) ([]byte, error) { + input = strings.TrimSpace(input) + if input == "" { + return []byte{}, nil + } + + // try to parse as json, if the input starts with '{' + if input[0] == '{' { + var raw createBucketInput + err := json.Unmarshal([]byte(input), &raw) + if err != nil { + return nil, fmt.Errorf("invalid JSON input: %w", err) + } + + return xml.Marshal(s3response.CreateBucketConfiguration{ + LocationConstraint: raw.LocationConstraint, + TagSet: raw.Tags, + }) + } + + var config s3response.CreateBucketConfiguration + + // parse as string - shorthand syntax + inputParts, err := splitTopLevel(input) + if err != nil { + return nil, err + } + for _, part := range inputParts { + part = strings.TrimSpace(part) + if strings.HasPrefix(part, "LocationConstraint=") { + locConstraint := strings.TrimPrefix(part, "LocationConstraint=") + config.LocationConstraint = &locConstraint + } else if strings.HasPrefix(part, "Tags=") { + tags, err := parseTagging(strings.TrimPrefix(part, "Tags=")) + if err != nil { + return nil, err + } + + config.TagSet = tags + } else { + return nil, fmt.Errorf("invalid component: %v", part) + } + } + + return xml.Marshal(config) +} + +var errInvalidTagsSyntax = errors.New("invalid tags syntax") + +// splitTopLevel splits a shorthand configuration string into top-level components. +// The function splits only on commas that are not nested inside '{}' or '[]'. +func splitTopLevel(s string) ([]string, error) { + var parts []string + start := 0 + depth := 0 + + for i, r := range s { + switch r { + case '{', '[': + depth++ + case '}', ']': + depth-- + case ',': + if depth == 0 { + parts = append(parts, s[start:i]) + start = i + 1 + } + } + } + + if depth != 0 { + return nil, errors.New("invalid string format") + } + + // add last segment + if start < len(s) { + parts = append(parts, s[start:]) + } + + return parts, nil +} + +// parseTagging parses a tag set expressed in shorthand syntax into AWS CLI tags. +// Expected format: +// +// [{Key=string,Value=string},{Key=string,Value=string}] +// +// The function validates bracket structure, splits tag objects at the top level, +// and delegates individual tag parsing to parseTag. It returns an error if the +// syntax is invalid or if any tag entry cannot be parsed. +func parseTagging(input string) ([]types.Tag, error) { + if len(input) < 2 { + return nil, errInvalidTagsSyntax + } + + if input[0] != '[' || input[len(input)-1] != ']' { + return nil, errInvalidTagsSyntax + } + // strip [] + input = input[1 : len(input)-1] + + tagComponents, err := splitTopLevel(input) + if err != nil { + return nil, errInvalidTagsSyntax + } + result := make([]types.Tag, 0, len(tagComponents)) + for _, tagComponent := range tagComponents { + tagComponent = strings.TrimSpace(tagComponent) + tag, err := parseTag(tagComponent) + if err != nil { + return nil, err + } + + result = append(result, tag) + } + + return result, nil +} + +// parseTag parses a single tag definition in shorthand form. +// Expected format: +// +// {Key=string,Value=string} +func parseTag(input string) (types.Tag, error) { + input = strings.TrimSpace(input) + + if len(input) < 2 { + return types.Tag{}, errInvalidTagsSyntax + } + + if input[0] != '{' || input[len(input)-1] != '}' { + return types.Tag{}, errInvalidTagsSyntax + } + + // strip {} + input = input[1 : len(input)-1] + + components := strings.Split(input, ",") + if len(components) != 2 { + return types.Tag{}, errInvalidTagsSyntax + } + + var key, value string + + for _, c := range components { + c = strings.TrimSpace(c) + + switch { + case strings.HasPrefix(c, "Key="): + key = strings.TrimPrefix(c, "Key=") + case strings.HasPrefix(c, "Value="): + value = strings.TrimPrefix(c, "Value=") + default: + return types.Tag{}, errInvalidTagsSyntax + } + } + + if key == "" { + return types.Tag{}, errInvalidTagsSyntax + } + + return types.Tag{ + Key: &key, + Value: &value, + }, nil +} + +func createBucket(ctx *cli.Context) error { + bucket, owner := ctx.String("bucket"), ctx.String("owner") + + payload, err := parseCreateBucketPayload(ctx.String("create-bucket-configuration")) + if err != nil { + return fmt.Errorf("invalid create bucket configuration: %w", err) + } + + hashedPayload := sha256.Sum256(payload) + hexPayload := hex.EncodeToString(hashedPayload[:]) + + headers := map[string]string{ + "x-amz-content-sha256": hexPayload, + "x-vgw-owner": owner, + "x-amz-acl": ctx.String("acl"), + "x-amz-grant-full-control": ctx.String("grant-full-control"), + "x-amz-grant-read": ctx.String("grant-read"), + "x-amz-grant-read-acp": ctx.String("grant-read-acp"), + "x-amz-grant-write": ctx.String("grant-write"), + "x-amz-grant-write-acp": ctx.String("grant-write-acp"), + "x-amz-object-ownership": ctx.String("object-ownership"), + } + + if ctx.Bool("object-lock-enabled-for-bucket") { + headers["x-amz-bucket-object-lock-enabled"] = "true" + } + if ctx.Bool("no-object-lock-enabled-for-bucket") { + headers["x-amz-bucket-object-lock-enabled"] = "false" + } + + req, err := http.NewRequestWithContext(ctx.Context, http.MethodPatch, fmt.Sprintf("%s/%s/create", adminEndpoint, bucket), bytes.NewReader(payload)) + if err != nil { + return err + } + + for key, value := range headers { + if value != "" { + req.Header.Set(key, value) + } + } + + signer := v4.NewSigner() + err = signer.SignHTTP(req.Context(), aws.Credentials{AccessKeyID: adminAccess, SecretAccessKey: adminSecret}, req, hexPayload, "s3", adminRegion, time.Now()) + if err != nil { + return fmt.Errorf("failed to sign the request: %w", err) + } + + client := initHTTPClient() + + resp, err := client.Do(req) + if err != nil { + return fmt.Errorf("failed to send the request: %w", err) + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + return err + } + + if resp.StatusCode >= 400 { + return parseApiError(body) + } + + return nil +} + const ( // account table formatting minwidth int = 2 // minimal cell width including any padding diff --git a/cmd/versitygw/main.go b/cmd/versitygw/main.go index 67626d94..399f061b 100644 --- a/cmd/versitygw/main.go +++ b/cmd/versitygw/main.go @@ -821,7 +821,7 @@ func runGateway(ctx context.Context, be backend.Backend) error { opts = append(opts, s3api.WithAdminDebug()) } - admSrv = s3api.NewAdminServer(be, middlewares.RootUserConfig{Access: rootUserAccess, Secret: rootUserSecret}, admPort, region, iam, loggers.AdminLogger, opts...) + admSrv = s3api.NewAdminServer(be, middlewares.RootUserConfig{Access: rootUserAccess, Secret: rootUserSecret}, admPort, region, iam, loggers.AdminLogger, srv.Router.Ctrl, opts...) } if !quiet { diff --git a/metrics/actions.go b/metrics/actions.go index 07559f9c..823c3818 100644 --- a/metrics/actions.go +++ b/metrics/actions.go @@ -125,6 +125,7 @@ var ( ActionAdminChangeBucketOwner = "admin_ChangeBucketOwner" ActionAdminListUsers = "admin_ListUsers" ActionAdminListBuckets = "admin_ListBuckets" + ActionAdminCreateBucket = "admin_CreateBucket" ) func init() { diff --git a/s3api/admin-router.go b/s3api/admin-router.go index b0534b96..40727f02 100644 --- a/s3api/admin-router.go +++ b/s3api/admin-router.go @@ -24,10 +24,12 @@ import ( "github.com/versity/versitygw/s3log" ) -type S3AdminRouter struct{} +type S3AdminRouter struct { + s3api controllers.S3ApiController +} func (ar *S3AdminRouter) Init(app *fiber.App, be backend.Backend, iam auth.IAMService, logger s3log.AuditLogger, root middlewares.RootUserConfig, region string, debug bool, corsAllowOrigin string) { - ctrl := controllers.NewAdminController(iam, be, logger) + ctrl := controllers.NewAdminController(iam, be, logger, ar.s3api) services := &controllers.Services{ Logger: logger, } @@ -103,4 +105,14 @@ func (ar *S3AdminRouter) Init(app *fiber.App, be backend.Backend, iam auth.IAMSe middlewares.ApplyDefaultCORSPreflight(corsAllowOrigin), middlewares.ApplyDefaultCORS(corsAllowOrigin), ) + + app.Patch("/:bucket/create", + controllers.ProcessHandlers(ctrl.CreateBucket, metrics.ActionAdminListBuckets, services, + middlewares.VerifyV4Signature(root, iam, region, false, true), + middlewares.IsAdmin(metrics.ActionAdminCreateBucket), + )) + app.Options("/:bucket/create", + middlewares.ApplyDefaultCORSPreflight(corsAllowOrigin), + middlewares.ApplyDefaultCORS(corsAllowOrigin), + ) } diff --git a/s3api/admin-server.go b/s3api/admin-server.go index 16fdec96..082d612f 100644 --- a/s3api/admin-server.go +++ b/s3api/admin-server.go @@ -38,11 +38,13 @@ type S3AdminServer struct { corsAllowOrigin string } -func NewAdminServer(be backend.Backend, root middlewares.RootUserConfig, port, region string, iam auth.IAMService, l s3log.AuditLogger, opts ...AdminOpt) *S3AdminServer { +func NewAdminServer(be backend.Backend, root middlewares.RootUserConfig, port, region string, iam auth.IAMService, l s3log.AuditLogger, ctrl controllers.S3ApiController, opts ...AdminOpt) *S3AdminServer { server := &S3AdminServer{ backend: be, - router: new(S3AdminRouter), - port: port, + router: &S3AdminRouter{ + s3api: ctrl, + }, + port: port, } for _, opt := range opts { diff --git a/s3api/controllers/admin.go b/s3api/controllers/admin.go index a85caf0e..4f86b7d6 100644 --- a/s3api/controllers/admin.go +++ b/s3api/controllers/admin.go @@ -28,13 +28,14 @@ import ( ) type AdminController struct { - iam auth.IAMService - be backend.Backend - l s3log.AuditLogger + iam auth.IAMService + be backend.Backend + l s3log.AuditLogger + s3api S3ApiController } -func NewAdminController(iam auth.IAMService, be backend.Backend, l s3log.AuditLogger) AdminController { - return AdminController{iam: iam, be: be, l: l} +func NewAdminController(iam auth.IAMService, be backend.Backend, l s3log.AuditLogger, s3api S3ApiController) AdminController { + return AdminController{iam: iam, be: be, l: l, s3api: s3api} } func (c AdminController) CreateUser(ctx *fiber.Ctx) (*Response, error) { @@ -161,3 +162,39 @@ func (c AdminController) ListBuckets(ctx *fiber.Ctx) (*Response, error) { MetaOpts: &MetaOptions{}, }, err } + +func (c AdminController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { + owner := ctx.Get("x-vgw-owner") + if owner == "" { + return &Response{ + MetaOpts: &MetaOptions{}, + }, s3err.GetAPIError(s3err.ErrAdminEmptyBucketOwnerHeader) + } + + acc, err := c.iam.GetUserAccount(owner) + if err != nil { + if err == auth.ErrNoSuchUser { + err = s3err.GetAPIError(s3err.ErrAdminUserNotFound) + } + + return &Response{ + MetaOpts: &MetaOptions{}, + }, err + } + + // store the owner access key id in context + ctx.Context().SetUserValue("bucket-owner", acc) + + _, err = c.s3api.CreateBucket(ctx) + if err != nil { + return &Response{ + MetaOpts: &MetaOptions{}, + }, err + } + + return &Response{ + MetaOpts: &MetaOptions{ + Status: http.StatusCreated, + }, + }, nil +} diff --git a/s3api/controllers/admin_test.go b/s3api/controllers/admin_test.go index 73ef184b..a540f193 100644 --- a/s3api/controllers/admin_test.go +++ b/s3api/controllers/admin_test.go @@ -21,6 +21,7 @@ import ( "net/http" "testing" + "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/stretchr/testify/assert" "github.com/versity/versitygw/auth" "github.com/versity/versitygw/backend" @@ -32,9 +33,10 @@ import ( func TestNewAdminController(t *testing.T) { type args struct { - iam auth.IAMService - be backend.Backend - l s3log.AuditLogger + iam auth.IAMService + be backend.Backend + l s3log.AuditLogger + s3api S3ApiController } tests := []struct { name string @@ -49,7 +51,7 @@ func TestNewAdminController(t *testing.T) { } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - got := NewAdminController(tt.args.iam, tt.args.be, tt.args.l) + got := NewAdminController(tt.args.iam, tt.args.be, tt.args.l, tt.args.s3api) assert.Equal(t, got, tt.want) }) } @@ -577,3 +579,126 @@ func TestAdminController_ListBuckets(t *testing.T) { }) } } + +func TestAdminController_CreateBucket(t *testing.T) { + tests := []struct { + name string + input testInput + output testOutput + }{ + { + name: "empty owner header", + output: testOutput{ + response: &Response{ + MetaOpts: &MetaOptions{}, + }, + err: s3err.GetAPIError(s3err.ErrAdminEmptyBucketOwnerHeader), + }, + }, + { + name: "fails to get user account", + input: testInput{ + extraMockErr: s3err.GetAPIError(s3err.ErrInternalError), + headers: map[string]string{ + "x-vgw-owner": "access", + }, + }, + output: testOutput{ + response: &Response{ + MetaOpts: &MetaOptions{}, + }, + err: s3err.GetAPIError(s3err.ErrInternalError), + }, + }, + { + name: "user not found", + input: testInput{ + extraMockErr: auth.ErrNoSuchUser, + headers: map[string]string{ + "x-vgw-owner": "access", + }, + }, + output: testOutput{ + response: &Response{ + MetaOpts: &MetaOptions{}, + }, + err: s3err.GetAPIError(s3err.ErrAdminUserNotFound), + }, + }, + { + name: "backend returns error", + input: testInput{ + headers: map[string]string{ + "x-vgw-owner": "access", + }, + locals: map[utils.ContextKey]any{ + utils.ContextKeyAccount: auth.Account{ + Access: "test-user", + Role: "admin", + }, + }, + beErr: s3err.GetAPIError(s3err.ErrAdminMethodNotSupported), + }, + output: testOutput{ + response: &Response{ + MetaOpts: &MetaOptions{}, + }, + err: s3err.GetAPIError(s3err.ErrAdminMethodNotSupported), + }, + }, + { + name: "successful response", + input: testInput{ + headers: map[string]string{ + "x-vgw-owner": "access", + }, + locals: map[utils.ContextKey]any{ + utils.ContextKeyAccount: auth.Account{ + Access: "test-user", + Role: "admin", + }, + }, + }, + output: testOutput{ + response: &Response{ + MetaOpts: &MetaOptions{ + Status: http.StatusCreated, + }, + }, + }, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + iam := &IAMServiceMock{ + GetUserAccountFunc: func(access string) (auth.Account, error) { + return auth.Account{}, tt.input.extraMockErr + }, + } + be := &BackendMock{ + CreateBucketFunc: func(contextMoqParam context.Context, createBucketInput *s3.CreateBucketInput, defaultACL []byte) error { + return tt.input.beErr + }, + } + + s3api := New(be, iam, nil, nil, nil, false, "") + + ctrl := AdminController{ + iam: iam, + be: be, + s3api: s3api, + } + + testController( + t, + ctrl.CreateBucket, + tt.output.response, + tt.output.err, + ctxInputs{ + locals: tt.input.locals, + headers: tt.input.headers, + }, + ) + }) + } +} diff --git a/s3api/controllers/bucket-put.go b/s3api/controllers/bucket-put.go index 292ce361..bd46dddd 100644 --- a/s3api/controllers/bucket-put.go +++ b/s3api/controllers/bucket-put.go @@ -461,7 +461,7 @@ func (c S3ApiController) PutBucketAcl(ctx *fiber.Ctx) (*Response, error) { }, s3err.GetAPIError(s3err.ErrMissingSecurityHeader) } - updAcl, err := auth.UpdateACL(input, parsedAcl, c.iam, acct.Role == auth.RoleAdmin) + updAcl, err := auth.UpdateACL(input, parsedAcl, c.iam) if err != nil { return &Response{ MetaOpts: &MetaOptions{ @@ -487,13 +487,18 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { grantWrite := ctx.Get("X-Amz-Grant-Write") grantWriteACP := ctx.Get("X-Amz-Grant-Write-Acp") lockEnabled := strings.EqualFold(ctx.Get("X-Amz-Bucket-Object-Lock-Enabled"), "true") - acct := utils.ContextKeyAccount.Get(ctx).(auth.Account) grants := grantFullControl + grantRead + grantReadACP + grantWrite + grantWriteACP objectOwnership := types.ObjectOwnership( ctx.Get("X-Amz-Object-Ownership", string(types.ObjectOwnershipBucketOwnerEnforced)), ) - if acct.Role != auth.RoleAdmin && acct.Role != auth.RoleUserPlus { + creator := utils.ContextKeyAccount.Get(ctx).(auth.Account) + if !utils.ContextKeyBucketOwner.IsSet(ctx) { + utils.ContextKeyBucketOwner.Set(ctx, creator) + } + bucketOwner := utils.ContextKeyBucketOwner.Get(ctx).(auth.Account) + + if creator.Role != auth.RoleAdmin && creator.Role != auth.RoleUserPlus { return &Response{ MetaOpts: &MetaOptions{}, }, s3err.GetAPIError(s3err.ErrAccessDenied) @@ -503,7 +508,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { if ok := utils.IsValidBucketName(bucket); !ok { return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.GetAPIError(s3err.ErrInvalidBucketName) } @@ -513,7 +518,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { if err != nil { return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, err } @@ -522,7 +527,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { if ok := utils.IsValidOwnership(objectOwnership); !ok { return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.APIError{ Code: "InvalidArgument", @@ -535,7 +540,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { debuglogger.Logf("bucket acls are disabled for %v object ownership", objectOwnership) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.GetAPIError(s3err.ErrInvalidBucketAclWithObjectOwnership) } @@ -544,7 +549,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { debuglogger.Logf("invalid request: %q (grants) %q (acl)", grants, acl) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.GetAPIError(s3err.ErrBothCannedAndHeaderGrants) } @@ -557,7 +562,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { debuglogger.Logf("failed to parse the request body: %v", err) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.GetAPIError(s3err.ErrMalformedXML) } @@ -568,7 +573,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { debuglogger.Logf("invalid location constraint: %s", *body.LocationConstraint) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, s3err.GetAPIError(s3err.ErrInvalidLocationConstraint) } @@ -576,7 +581,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { } defACL := auth.ACL{ - Owner: acct.Access, + Owner: bucketOwner.Access, } updAcl, err := auth.UpdateACL(&auth.PutBucketAclInput{ @@ -587,15 +592,15 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { GrantWriteACP: &grantWriteACP, AccessControlPolicy: &auth.AccessControlPolicy{ Owner: &types.Owner{ - ID: &acct.Access, + ID: &bucketOwner.Access, }}, ACL: types.BucketCannedACL(acl), - }, defACL, c.iam, acct.Role == auth.RoleAdmin) + }, defACL, c.iam) if err != nil { debuglogger.Logf("failed to update bucket acl: %v", err) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, err } @@ -610,7 +615,7 @@ func (c S3ApiController) CreateBucket(ctx *fiber.Ctx) (*Response, error) { }, updAcl) return &Response{ MetaOpts: &MetaOptions{ - BucketOwner: acct.Access, + BucketOwner: bucketOwner.Access, }, }, err } diff --git a/s3api/router.go b/s3api/router.go index bde5e05a..30f2f6e7 100644 --- a/s3api/router.go +++ b/s3api/router.go @@ -28,16 +28,18 @@ import ( type S3ApiRouter struct { WithAdmSrv bool + Ctrl controllers.S3ApiController } func (sa *S3ApiRouter) Init(app *fiber.App, be backend.Backend, iam auth.IAMService, logger s3log.AuditLogger, aLogger s3log.AuditLogger, evs s3event.S3EventSender, mm metrics.Manager, readonly bool, region, virtualDomain string, root middlewares.RootUserConfig, corsAllowOrigin string) { ctrl := controllers.New(be, iam, logger, evs, mm, readonly, virtualDomain) + sa.Ctrl = ctrl adminServices := &controllers.Services{ Logger: aLogger, } if sa.WithAdmSrv { - adminController := controllers.NewAdminController(iam, be, aLogger) + adminController := controllers.NewAdminController(iam, be, aLogger, ctrl) // CreateUser admin api app.Patch("/create-user", @@ -110,6 +112,18 @@ func (sa *S3ApiRouter) Init(app *fiber.App, be backend.Backend, iam auth.IAMServ middlewares.ApplyDefaultCORSPreflight(corsAllowOrigin), middlewares.ApplyDefaultCORS(corsAllowOrigin), ) + + // CreateBucket admin api + app.Patch("/:bucket/create", + controllers.ProcessHandlers(adminController.CreateBucket, metrics.ActionAdminCreateBucket, adminServices, + middlewares.VerifyV4Signature(root, iam, region, false, true), + middlewares.IsAdmin(metrics.ActionAdminCreateBucket), + middlewares.ApplyDefaultCORS(corsAllowOrigin), + )) + app.Options("/:bucket/create", + middlewares.ApplyDefaultCORSPreflight(corsAllowOrigin), + middlewares.ApplyDefaultCORS(corsAllowOrigin), + ) } services := &controllers.Services{ diff --git a/s3api/server.go b/s3api/server.go index 628358c2..d32c506e 100644 --- a/s3api/server.go +++ b/s3api/server.go @@ -41,9 +41,9 @@ const ( ) type S3ApiServer struct { + Router *S3ApiRouter app *fiber.App backend backend.Backend - router *S3ApiRouter port string cert *tls.Certificate quiet bool @@ -67,7 +67,7 @@ func New( ) (*S3ApiServer, error) { server := &S3ApiServer{ backend: be, - router: new(S3ApiRouter), + Router: new(S3ApiRouter), port: port, } @@ -124,7 +124,7 @@ func New( app.Use(middlewares.DebugLogger()) } - server.router.Init(app, be, iam, l, adminLogger, evs, mm, server.readonly, region, server.virtualDomain, root, server.corsAllowOrigin) + server.Router.Init(app, be, iam, l, adminLogger, evs, mm, server.readonly, region, server.virtualDomain, root, server.corsAllowOrigin) return server, nil } @@ -139,7 +139,7 @@ func WithTLS(cert tls.Certificate) Option { // WithAdminServer runs admin endpoints with the gateway in the same network func WithAdminServer() Option { - return func(s *S3ApiServer) { s.router.WithAdmSrv = true } + return func(s *S3ApiServer) { s.Router.WithAdmSrv = true } } // WithQuiet silences default logging output diff --git a/s3api/server_test.go b/s3api/server_test.go index aede682b..f50eb268 100644 --- a/s3api/server_test.go +++ b/s3api/server_test.go @@ -35,7 +35,7 @@ func TestS3ApiServer_Serve(t *testing.T) { app: fiber.New(), backend: backend.BackendUnsupported{}, port: "Invalid address", - router: &S3ApiRouter{}, + Router: &S3ApiRouter{}, }, }, { @@ -45,7 +45,7 @@ func TestS3ApiServer_Serve(t *testing.T) { app: fiber.New(), backend: backend.BackendUnsupported{}, port: "Invalid address", - router: &S3ApiRouter{}, + Router: &S3ApiRouter{}, cert: &tls.Certificate{}, }, }, diff --git a/s3api/utils/context-keys.go b/s3api/utils/context-keys.go index fdb3cf18..76dfde0b 100644 --- a/s3api/utils/context-keys.go +++ b/s3api/utils/context-keys.go @@ -36,6 +36,7 @@ const ( ContextKeyBodyReader ContextKey = "body-reader" ContextKeySkip ContextKey = "__skip" ContextKeyStack ContextKey = "stack" + ContextKeyBucketOwner ContextKey = "bucket-owner" ) func (ck ContextKey) Values() []ContextKey { @@ -50,6 +51,7 @@ func (ck ContextKey) Values() []ContextKey { ContextKeyParsedAcl, ContextKeySkipResBodyLog, ContextKeyBodyReader, + ContextKeyBucketOwner, } } diff --git a/s3err/s3err.go b/s3err/s3err.go index eec7f5f3..6721d42d 100644 --- a/s3err/s3err.go +++ b/s3err/s3err.go @@ -201,6 +201,7 @@ const ( ErrAdminInvalidUserRole ErrAdminMissingUserAcess ErrAdminMethodNotSupported + ErrAdminEmptyBucketOwnerHeader ) var errorCodeResponse = map[ErrorCode]APIError{ @@ -898,6 +899,11 @@ var errorCodeResponse = map[ErrorCode]APIError{ Description: "The method is not supported in single root user mode.", HTTPStatusCode: http.StatusNotImplemented, }, + ErrAdminEmptyBucketOwnerHeader: { + Code: "XAdminInvalidRequest", + Description: "The x-vgw-owner header specifying the new bucket owner access key id is either missing or empty", + HTTPStatusCode: http.StatusBadRequest, + }, } // GetAPIError provides API Error for input API error code. diff --git a/tests/integration/group-tests.go b/tests/integration/group-tests.go index d6b2798a..92cef368 100644 --- a/tests/integration/group-tests.go +++ b/tests/integration/group-tests.go @@ -823,6 +823,7 @@ func TestFullFlow(ts *TestState) { if ts.conf.versioningEnabled { TestVersioning(ts) } + TestIAM(ts) } func TestPosix(ts *TestState) { @@ -945,7 +946,10 @@ func TestIAM(ts *TestState) { ts.Run(IAM_userplus_CreateBucket) ts.Run(IAM_admin_ChangeBucketOwner) ts.Run(IAM_ChangeBucketOwner_back_to_root) - ts.Run(IAM_ListBuckets) + ts.Sync(IAM_ListBuckets) + ts.Run(IAM_CreateBucket_empty_owner_header) + ts.Run(IAM_CreateBucket_non_existing_user) + ts.Run(IAM_CreateBucket_success) } func TestAccessControl(ts *TestState) { @@ -1672,6 +1676,10 @@ func GetIntTests() IntTests { "IAM_userplus_CreateBucket": IAM_userplus_CreateBucket, "IAM_admin_ChangeBucketOwner": IAM_admin_ChangeBucketOwner, "IAM_ChangeBucketOwner_back_to_root": IAM_ChangeBucketOwner_back_to_root, + "IAM_ListBuckets": IAM_ListBuckets, + "IAM_CreateBucket_empty_owner_header": IAM_CreateBucket_empty_owner_header, + "IAM_CreateBucket_non_existing_user": IAM_CreateBucket_non_existing_user, + "IAM_CreateBucket_success": IAM_CreateBucket_success, "AccessControl_default_ACL_user_access_denied": AccessControl_default_ACL_user_access_denied, "AccessControl_default_ACL_userplus_access_denied": AccessControl_default_ACL_userplus_access_denied, "AccessControl_default_ACL_admin_successful_access": AccessControl_default_ACL_admin_successful_access, diff --git a/tests/integration/iam.go b/tests/integration/iam.go index ebdf3fff..f1bfe695 100644 --- a/tests/integration/iam.go +++ b/tests/integration/iam.go @@ -16,11 +16,16 @@ package integration import ( "context" + "encoding/xml" "fmt" + "net/http" "strings" + "time" "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" "github.com/versity/versitygw/s3err" + "github.com/versity/versitygw/s3response" ) func IAM_user_access_denied(s *S3Conf) error { @@ -169,11 +174,191 @@ func IAM_ChangeBucketOwner_back_to_root(s *S3Conf) error { func IAM_ListBuckets(s *S3Conf) error { testName := "IAM_ListBuckets" return actionHandler(s, testName, func(s3client *s3.Client, bucket string) error { - err := listBuckets(s) + return listBuckets(s) + }) +} + +func IAM_CreateBucket_empty_owner_header(s *S3Conf) error { + testName := "IAM_CreateBucket_empty_owner_header" + return actionHandlerNoSetup(s, testName, func(s3client *s3.Client, bucket string) error { + req, err := createSignedReq( + http.MethodPatch, + s.endpoint, + fmt.Sprintf("%s/create", bucket), + s.awsID, + s.awsSecret, + "s3", + s.awsRegion, + nil, + time.Now(), + nil, + ) if err != nil { return err } - return nil + resp, err := s.httpClient.Do(req) + if err != nil { + return err + } + + return checkHTTPResponseApiErr(resp, s3err.GetAPIError(s3err.ErrAdminEmptyBucketOwnerHeader)) + }) +} + +func IAM_CreateBucket_non_existing_user(s *S3Conf) error { + testName := "IAM_CreateBucket_non_existing_user" + return actionHandlerNoSetup(s, testName, func(s3client *s3.Client, bucket string) error { + req, err := createSignedReq( + http.MethodPatch, + s.endpoint, + fmt.Sprintf("%s/create", bucket), + s.awsID, + s.awsSecret, + "s3", + s.awsRegion, + nil, + time.Now(), + map[string]string{ + "x-vgw-owner": "non-existing-user", + }, + ) + if err != nil { + return err + } + + resp, err := s.httpClient.Do(req) + if err != nil { + return err + } + + return checkHTTPResponseApiErr(resp, s3err.GetAPIError(s3err.ErrAdminUserNotFound)) + }) +} + +func IAM_CreateBucket_success(s *S3Conf) error { + testName := "IAM_CreateBucket_success" + return actionHandlerNoSetup(s, testName, func(s3client *s3.Client, bucket string) error { + tagSet := []types.Tag{ + {Key: getPtr("key"), Value: getPtr("value")}, + } + body, err := xml.Marshal(s3response.CreateBucketConfiguration{ + TagSet: tagSet, + }) + if err != nil { + return err + } + + testUser1, testUser2 := getUser("user"), getUser("user") + err = createUsers(s, []user{testUser1, testUser2}) + if err != nil { + return err + } + + req, err := createSignedReq( + http.MethodPatch, + s.endpoint, + fmt.Sprintf("%s/create", bucket), + s.awsID, + s.awsSecret, + "s3", + s.awsRegion, + body, + time.Now(), + map[string]string{ + "x-amz-bucket-object-lock-enabled": "true", + "x-amz-object-ownership": string(types.ObjectOwnershipBucketOwnerPreferred), + "x-amz-grant-read": testUser2.access, + "x-vgw-owner": testUser1.access, + }, + ) + if err != nil { + return err + } + + resp, err := s.httpClient.Do(req) + if err != nil { + return err + } + + if resp.StatusCode != http.StatusCreated { + return fmt.Errorf("expected the response status code to be %v, instead got %v", http.StatusCreated, resp.StatusCode) + } + + ctx, cancel := context.WithTimeout(context.Background(), shortTimeout) + tagging, err := s3client.GetBucketTagging(ctx, &s3.GetBucketTaggingInput{ + Bucket: &bucket, + }) + cancel() + if err != nil { + return err + } + + if !areTagsSame(tagSet, tagging.TagSet) { + return fmt.Errorf("expected the bucket tagging to be %v, instead got %v", tagSet, tagging.TagSet) + } + + ctx, cancel = context.WithTimeout(context.Background(), shortTimeout) + ownership, err := s3client.GetBucketOwnershipControls(ctx, &s3.GetBucketOwnershipControlsInput{ + Bucket: &bucket, + }) + cancel() + if err != nil { + return err + } + + var ownershipControls types.ObjectOwnership + if ownership.OwnershipControls != nil && len(ownership.OwnershipControls.Rules) == 1 { + ownershipControls = ownership.OwnershipControls.Rules[0].ObjectOwnership + } + + if ownershipControls != types.ObjectOwnershipBucketOwnerPreferred { + return fmt.Errorf("expected the bucket ownership controls to be %s, instaed got %s", types.ObjectOwnershipBucketOwnerPreferred, ownershipControls) + } + + ctx, cancel = context.WithTimeout(context.Background(), shortTimeout) + acl, err := s3client.GetBucketAcl(ctx, &s3.GetBucketAclInput{ + Bucket: &bucket, + }) + cancel() + if err != nil { + return err + } + + if len(acl.Grants) != 2 { + return fmt.Errorf("expected the length of acl grants to be 2, instead got %v", len(acl.Grants)) + } + + var granteeChecked bool + var ownerChecked bool + for _, grant := range acl.Grants { + // owner + if getString(grant.Grantee.ID) == testUser1.access { + ownerChecked = true + if grant.Permission != types.PermissionFullControl { + return fmt.Errorf("expected the owner '%s' to have %s permission, instead got %s", testUser1.access, types.PermissionFullControl, grant.Permission) + } + + continue + } + + if getString(grant.Grantee.ID) != testUser2.access { + return fmt.Errorf("expected the grantee ID to be %v, instaed got %v", testUser2.access, getString(grant.Grantee.ID)) + } + if grant.Permission != types.PermissionRead { + return fmt.Errorf("expected the %v user permission to be %s, instead got %s", testUser2.access, types.PermissionRead, grant.Permission) + } + + granteeChecked = true + } + + if !ownerChecked { + return fmt.Errorf("missing the owner '%s' full control acl", testUser1.access) + } + if !granteeChecked { + return fmt.Errorf("missing the user %s in read grantees acl", testUser2.access) + } + + return teardown(s, bucket) }) }