Refactor: param parse in object storage interface handlers

Signed-off-by: Mervin <mofei2816@gmail.com>
This commit is contained in:
Mervin 2020-03-10 00:43:57 +08:00
parent cec09e50fa
commit fbbe1e0a7e
No known key found for this signature in database
GPG Key ID: 7C19725C197AE672
25 changed files with 1056 additions and 662 deletions

View File

@ -1784,7 +1784,7 @@ func (m *Server) createUserWithPolicy(userID, volName string) (err error) {
return
}
policy := &proto.UserPolicy{
OwnVol: []string{volName},
OwnVols: []string{volName},
}
if _, err = m.user.addPolicy(akPolicy.AccessKey, policy); err != nil {
return

View File

@ -58,7 +58,7 @@ func (u *User) createKey(owner string) (akPolicy *proto.AKPolicy, err error) {
accessKey = util.RandomString(accessKeyLength, util.Numeric|util.LowerLetter|util.UpperLetter)
_, exit = u.akStore.Load(accessKey)
}
userPolicy = &proto.UserPolicy{OwnVol: make([]string, 0), NoneOwnVol: make(map[string][]string)}
userPolicy = &proto.UserPolicy{OwnVols: make([]string, 0), NoneOwnVol: make(map[string][]string)}
akPolicy = &proto.AKPolicy{AccessKey: accessKey, SecretKey: secretKey, Policy: userPolicy, UserID: owner}
userAK = &proto.UserAK{UserID: owner, AccessKey: accessKey}
if err = u.syncAddAKPolicy(akPolicy); err != nil {
@ -96,7 +96,7 @@ func (u *User) createUserWithKey(owner, accessKey, secretKey string) (akPolicy *
err = proto.ErrDuplicateAccessKey
goto errHandler
}
userPolicy = &proto.UserPolicy{OwnVol: make([]string, 0), NoneOwnVol: make(map[string][]string)}
userPolicy = &proto.UserPolicy{OwnVols: make([]string, 0), NoneOwnVol: make(map[string][]string)}
akPolicy = &proto.AKPolicy{AccessKey: accessKey, SecretKey: secretKey, Policy: userPolicy, UserID: owner}
userAK = &proto.UserAK{UserID: owner, AccessKey: accessKey}
if err = u.syncAddAKPolicy(akPolicy); err != nil {
@ -129,7 +129,7 @@ func (u *User) deleteKey(owner string) (err error) {
if akPolicy, err = u.getAKInfo(userAK.AccessKey); err != nil {
goto errHandler
}
if len(akPolicy.Policy.OwnVol) > 0 {
if len(akPolicy.Policy.OwnVols) > 0 {
err = proto.ErrOwnVolExits
goto errHandler
}
@ -244,7 +244,7 @@ func (u *User) deleteVolPolicy(vol string) (err error) {
}
var userPolicy *proto.UserPolicy
if action == ALL {
userPolicy = &proto.UserPolicy{OwnVol: []string{vol}}
userPolicy = &proto.UserPolicy{OwnVols: []string{vol}}
} else {
userPolicy = &proto.UserPolicy{NoneOwnVol: map[string][]string{vol: {action}}}
}
@ -269,11 +269,11 @@ errHandler:
func (u *User) transferVol(vol, ak, targetKey string) (targetAKPolicy *proto.AKPolicy, err error) {
var akPolicy *proto.AKPolicy
userPolicy := &proto.UserPolicy{OwnVol: []string{vol}}
userPolicy := &proto.UserPolicy{OwnVols: []string{vol}}
if akPolicy, err = u.getAKInfo(ak); err != nil {
goto errHandler
}
if !contains(akPolicy.Policy.OwnVol, vol) {
if !contains(akPolicy.Policy.OwnVols, vol) {
err = proto.ErrHaveNoPolicy
goto errHandler
}
@ -303,7 +303,7 @@ func (u *User) getAKInfo(ak string) (akPolicy *proto.AKPolicy, err error) {
func (u *User) addVolAKs(ak string, policy *proto.UserPolicy) (err error) {
u.volAKsMutex.Lock()
defer u.volAKsMutex.Unlock()
for _, vol := range policy.OwnVol {
for _, vol := range policy.OwnVols {
if err = u.addAKToVol(ak+separator+ALL, vol); err != nil {
return
}
@ -339,7 +339,7 @@ func (u *User) addAKToVol(akAndAction string, vol string) (err error) {
}
func (u *User) deleteVolAKs(ak string, policy *proto.UserPolicy) (err error) {
for _, vol := range policy.OwnVol {
for _, vol := range policy.OwnVols {
if err = u.deleteAKFromVol(ak+separator+ALL, vol); err != nil {
return
}

View File

@ -203,8 +203,8 @@ func (acp *AccessControlPolicy) SetBucketStandardACL(param *RequestParam, acl st
if uri, ok := aclRoleURIMap[role]; ok {
grantee.URI = uri
} else {
grantee.Id = param.account
grantee.DisplayName = param.account
grantee.Id = param.accessKey
grantee.DisplayName = param.accessKey
}
for _, p := range permissions {
grant := Grant{
@ -218,8 +218,8 @@ func (acp *AccessControlPolicy) SetBucketStandardACL(param *RequestParam, acl st
func (acp *AccessControlPolicy) SetBucketGrantACL(param *RequestParam, permission Permission) {
grantee := Grantee{
Id: param.account,
DisplayName: param.account,
Id: param.accessKey,
DisplayName: param.accessKey,
}
grant := Grant{
Grantee: grantee,
@ -228,8 +228,8 @@ func (acp *AccessControlPolicy) SetBucketGrantACL(param *RequestParam, permissio
acp.Acl.Grants = append(acp.Acl.Grants, grant)
}
func (acl *AccessControlPolicy) Marshal() ([]byte, error) {
data, err := xml.Marshal(acl)
func (acp *AccessControlPolicy) Marshal() ([]byte, error) {
data, err := xml.Marshal(acp)
if err != nil {
return nil, err
}
@ -254,7 +254,7 @@ func ParseACL(bytes []byte, bucket string) (*AccessControlPolicy, error) {
return acl, nil
}
func storeBucketACL(bytes []byte, vol *volume) (*AccessControlPolicy, error) {
func storeBucketACL(bytes []byte, vol *Volume) (*AccessControlPolicy, error) {
store, err1 := vol.vm.GetStore()
if err1 != nil {
return nil, err1
@ -280,9 +280,9 @@ func (g Grant) Validate() bool {
}
func (g *Grant) IsAllowed(param *RequestParam) bool {
if param.account != g.Grantee.Id {
if param.accessKey != g.Grantee.Id {
return false
}
actions := aclBucketPermissionActions[g.Permission]
return IsIntersectionActions(actions, param.actions)
return IsIntersectionActions(actions, param.Action())
}

View File

@ -18,7 +18,6 @@ package objectnode
import (
"encoding/xml"
"errors"
"io"
"io/ioutil"
"net/http"
@ -54,13 +53,18 @@ func (o *ObjectNode) getBucketACLHandler(w http.ResponseWriter, r *http.Request)
)
defer o.errorResponse(w, r, err, ec)
_, bucket, _, vol, err := o.parseRequestParams(r)
if bucket == "" {
var param *RequestParam
param = ParseRequestParam(r)
if param.Bucket() == "" {
ec = &NoSuchBucket
return
}
//om := vol.OSSMeta()
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
ec = &NoSuchBucket
return
}
acl := vol.loadACL()
var aclData []byte
if acl != nil {
@ -98,63 +102,53 @@ func (o *ObjectNode) putBucketACLHandler(w http.ResponseWriter, r *http.Request)
defer o.errorResponse(w, r, err, ec)
log.LogInfof("Put bucket acl")
_, bucket, _, vol, err1 := o.parseRequestParams(r)
if err1 != nil {
err = err1
var param *RequestParam
param = ParseRequestParam(r)
if param.Bucket() == "" {
ec = &NoSuchBucket
return
}
if bucket == "" {
err = errors.New("")
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
ec = &NoSuchBucket
return
}
bytes, err2 := ioutil.ReadAll(r.Body)
if err2 != nil && err2 != io.EOF {
err = err2
var bytes []byte
if bytes, err = ioutil.ReadAll(r.Body); err != nil && err != io.EOF {
return
}
acl, err3 := ParseACL(bytes, vol.name)
if err3 != nil {
err = err3
var acp *AccessControlPolicy
if acp, err = ParseACL(bytes, param.Bucket()); err != nil {
return
}
if acl == nil {
err = errors.New("")
if acp == nil {
return
}
//add standard acl request header
// https://docs.aws.amazon.com/zh_cn/AmazonS3/latest/dev/acl-overview.html
p, err5 := o.parseRequestParam(r)
if err5 != nil {
err = err5
return
}
if standardAcls, found := r.Header["x-amz-acl"]; found {
acl.SetBucketStandardACL(p, standardAcls[0])
acp.SetBucketStandardACL(param, standardAcls[0])
} else {
for grant, permission := range aclGrantKeyPermissionMap {
if _, found2 := r.Header[grant]; found2 {
acl.SetBucketGrantACL(p, permission)
acp.SetBucketGrantACL(param, permission)
}
}
}
newBytes, err6 := acl.Marshal()
if err6 != nil {
err = err6
var newBytes []byte
if newBytes, err = acp.Marshal(); err != nil {
return
}
// store bucket acl
_, err4 := storeBucketACL(newBytes, vol)
if err4 != nil {
err = err4
if _, err = storeBucketACL(newBytes, vol); err != nil {
return
}
return
}

View File

@ -25,31 +25,47 @@ import (
)
type RequestParam struct {
account string
resource string
bucket string
object string
actions Action
sourceIP string
vol *volume
condVals map[string][]string
isOwner bool
vars map[string]string
accessKey string
resource string
bucket string
object string
action Action
sourceIP string
conditionVars map[string][]string
vars map[string]string
accessKey string
}
func (p *RequestParam) Bucket() string {
return p.bucket
}
func (p *RequestParam) Object() string {
return p.object
}
func (p *RequestParam) Action() Action {
return p.action
}
func (p *RequestParam) GetVar(name string) string {
return p.vars[name]
}
func (o *ObjectNode) parseRequestParam(r *http.Request) (*RequestParam, error) {
func (p *RequestParam) GetConditionVar(name string) []string {
return p.conditionVars[name]
}
func (p *RequestParam) AccessKey() string {
return p.accessKey
}
func ParseRequestParam(r *http.Request) *RequestParam {
p := new(RequestParam)
p.vars = mux.Vars(r)
p.bucket = p.vars["bucket"]
p.object = p.vars["object"]
p.vol, _ = o.getVol(p.bucket)
p.sourceIP = getRequestIP(r)
p.condVals = getCondtionValues(r)
p.conditionVars = getCondtionValues(r)
if len(p.bucket) > 0 {
p.resource = p.bucket
if len(p.object) > 0 {
@ -61,50 +77,44 @@ func (o *ObjectNode) parseRequestParam(r *http.Request) (*RequestParam, error) {
}
}
auth := parseRequestAuthInfo(r)
if auth != nil && p.vol != nil {
accessKey, _ := p.vol.OSSSecure()
p.account = accessKey
if auth.accessKey == accessKey {
p.isOwner = true
}
if auth != nil {
p.accessKey = auth.accessKey
}
p.action = GetActionFromContext(r)
if !p.action.IsKnown() {
p.action = ActionFromRouteName(mux.CurrentRoute(r).GetName())
}
return p, nil
return p
}
//Deprecated:
func (o *ObjectNode) parseRequestParams(r *http.Request) (vars map[string]string, bucket, object string, vl *volume, err error) {
vars = mux.Vars(r)
bucket = vars["bucket"]
object = vars["object"]
if bucket != "" {
if vm, ok := o.vm.(*volumeManager); ok {
vl, err = vm.loadVolume(bucket)
if err != nil {
log.LogErrorf("parseRequestParams: load volume fail, requestId(%v) bucket(%v) err(%v)",
GetRequestID(r), bucket, err)
}
} else {
log.LogErrorf("parseRequestParams: load volume fail, requestId(%v) bucket(%v) err(%v)",
GetRequestID(r), bucket, err)
}
}
return
}
//func (o *ObjectNode) parseRequestParams(r *http.Request) (vars map[string]string, bucket, object string, vl *Volume, err error) {
// vars = mux.Vars(r)
// bucket = vars["bucket"]
// object = vars["object"]
// if bucket != "" {
// if vm, ok := o.vm.(*VolumeManager); ok {
// vl, err = vm.loadVolume(bucket)
// if err != nil {
// log.LogErrorf("parseRequestParams: load Volume fail, requestId(%v) bucket(%v) err(%v)",
// GetRequestID(r), bucket, err)
// }
// } else {
// log.LogErrorf("parseRequestParams: load Volume fail, requestId(%v) bucket(%v) err(%v)",
// GetRequestID(r), bucket, err)
// }
// }
// return
//}
func (o *ObjectNode) getVol(bucket string) (vol *volume, err error) {
func (o *ObjectNode) getVol(bucket string) (vol *Volume, err error) {
if bucket == "" {
return nil, errors.New("bucket name is empty")
}
vm, ok := o.vm.(*volumeManager)
if !ok {
return nil, errors.New("volumeManger is invalid")
}
vol, err = vm.loadVolume(bucket)
vol, err = o.vm.loadVolume(bucket)
if err != nil {
log.LogErrorf("parseRequestParams: load volume fail, bucket(%v) err(%v)", bucket, err)
log.LogErrorf("parseRequestParams: load Volume fail, bucket(%v) err(%v)", bucket, err)
return nil, err
}

View File

@ -107,7 +107,7 @@ func (o *ObjectNode) deleteBucketHandler(w http.ResponseWriter, r *http.Request)
_ = InternalError.ServeResponse(w, r)
return
}
// get volume use state
// get Volume use state
if mc, err = o.vm.GetMasterClient(); err != nil {
log.LogErrorf("get master client error: err(%v)", err)
_ = InternalError.ServeResponse(w, r)
@ -127,7 +127,7 @@ func (o *ObjectNode) deleteBucketHandler(w http.ResponseWriter, r *http.Request)
_ = InternalError.ServeResponse(w, r)
return
}
// delete volume from master
// delete Volume from master
if authKey, err = calculateAuthKey(akPolicy.UserID); err != nil {
log.LogErrorf("delete bucket[%v] error: calculate authKey(%v) err(%v)", bucket, akPolicy.UserID, err)
_ = InternalError.ServeResponse(w, r)
@ -139,7 +139,7 @@ func (o *ObjectNode) deleteBucketHandler(w http.ResponseWriter, r *http.Request)
return
}
// release volume from volume manager
// release Volume from Volume manager
o.vm.Release(bucket)
return
}
@ -158,7 +158,7 @@ func (o *ObjectNode) listBucketsHandler(w http.ResponseWriter, r *http.Request)
return
}
ownVols := akPolicy.Policy.OwnVol
ownVols := akPolicy.Policy.OwnVols
var buckets = make([]*Bucket, 0)
for _, ownVol := range ownVols {
var bucket = &Bucket{Name: ownVol, CreationDate: time.Now()} //todo time
@ -208,20 +208,20 @@ func (o *ObjectNode) getBucketLocation(w http.ResponseWriter, r *http.Request) {
// API reference: https://docs.aws.amazon.com/AmazonS3/latest/API/API_GetBucketTagging.html
func (o *ObjectNode) getBucketTaggingHandler(w http.ResponseWriter, r *http.Request) {
var err error
var param *RequestParam
if param, err = o.parseRequestParam(r); err != nil {
log.LogErrorf("getBucketTaggingHandler: parse request param fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InvalidArgument.ServeResponse(w, r)
var param = ParseRequestParam(r)
if len(param.Bucket()) == 0 {
_ = InvalidBucketName.ServeResponse(w, r)
return
}
if param.vol == nil {
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
_ = NoSuchBucket.ServeResponse(w, r)
return
}
var xattrInfo *proto.XAttrInfo
if xattrInfo, err = param.vol.GetXAttr("/", XAttrKeyOSSTagging); err != nil {
log.LogErrorf("getBucketTaggingHandler: volume get XAttr fail: requestID(%v) err(%v)", GetRequestID(r), err)
if xattrInfo, err = vol.GetXAttr("/", XAttrKeyOSSTagging); err != nil {
log.LogErrorf("getBucketTaggingHandler: Volume get XAttr fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
return
}
@ -250,14 +250,22 @@ func (o *ObjectNode) getBucketTaggingHandler(w http.ResponseWriter, r *http.Requ
// API reference: https://docs.aws.amazon.com/AmazonS3/latest/API/API_PutBucketTagging.html
func (o *ObjectNode) putBucketTaggingHandler(w http.ResponseWriter, r *http.Request) {
var err error
var param *RequestParam
if param, err = o.parseRequestParam(r); err != nil {
log.LogErrorf("putBucketTaggingHandler: parse request param fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InvalidArgument.ServeResponse(w, r)
var errorCode *ErrorCode
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
}
}()
var param = ParseRequestParam(r)
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
if param.vol == nil {
_ = NoSuchBucket.ServeResponse(w, r)
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
errorCode = &NoSuchBucket
return
}
@ -282,7 +290,7 @@ func (o *ObjectNode) putBucketTaggingHandler(w http.ResponseWriter, r *http.Requ
return
}
if err = param.vol.SetXAttr("/", XAttrKeyOSSTagging, encoded); err != nil {
if err = vol.SetXAttr("/", XAttrKeyOSSTagging, encoded); err != nil {
_ = InternalError.ServeResponse(w, r)
return
}
@ -294,27 +302,31 @@ func (o *ObjectNode) putBucketTaggingHandler(w http.ResponseWriter, r *http.Requ
// API reference: https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteBucketTagging.html
func (o *ObjectNode) deleteBucketTaggingHandler(w http.ResponseWriter, r *http.Request) {
var (
param *RequestParam
err error
err error
errorCode *ErrorCode
)
if param, err = o.parseRequestParam(r); err != nil {
log.LogErrorf("deleteBucketTaggingHandler: parse request param fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InvalidArgument.ServeResponse(w, r)
return
}
var volume Volume
if len(param.bucket) == 0 {
_ = NoSuchBucket.ServeResponse(w, r)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
var param = ParseRequestParam(r)
if len(param.Bucket()) == 0 {
errorCode = &InvalidBucketName
return
}
if volume, err = o.vm.Volume(param.bucket); err != nil {
log.LogErrorf("deleteBucketTaggingHandler: load volume fail: requestID(%v) volume(%v) err(%v)", GetRequestID(r), param.bucket, err)
_ = NoSuchBucket.ServeResponse(w, r)
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("deleteBucketTaggingHandler: load Volume fail: requestID(%v) Volume(%v) err(%v)", GetRequestID(r), param.bucket, err)
errorCode = &NoSuchBucket
return
}
if err = volume.DeleteXAttr("/", XAttrKeyOSSTagging); err != nil {
log.LogErrorf("deleteBucketTaggingHandler: volume delete tagging xattr fail: requestID(%v) err(%v)", GetRequestID(r), err)
if err = vol.DeleteXAttr("/", XAttrKeyOSSTagging); err != nil {
log.LogErrorf("deleteBucketTaggingHandler: Volume delete tagging xattr fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
return
}

View File

@ -27,25 +27,46 @@ func (o *ObjectNode) createMultipleUploadHandler(w http.ResponseWriter, r *http.
log.LogInfof("createMultipleUploadHandler: init multiple upload, requestID(%v) remote(%v)",
GetRequestID(r), r.RemoteAddr)
_, bucket, object, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("createMultipleUploadHandler: parse request parameters fail, requestID(%v) err(%v)",
GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
var err error
var errorCode *ErrorCode
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
var param = ParseRequestParam(r)
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
uploadId, initErr := vl.InitMultipart(object)
if initErr != nil {
if param.Object() == "" {
errorCode = &InvalidKey
return
}
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("createMultipleUploadHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
var uploadID string
if uploadID, err = vol.InitMultipart(param.Object()); err != nil {
log.LogErrorf("createMultipleUploadHandler: init multipart fail, requestID(%v) err(%v)",
GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
errorCode = &InternalError
return
}
initResult := InitMultipartResult{
Bucket: bucket,
Key: object,
UploadId: uploadId,
Bucket: param.Bucket(),
Key: param.Object(),
UploadId: uploadID,
}
var bytes []byte
@ -53,7 +74,7 @@ func (o *ObjectNode) createMultipleUploadHandler(w http.ResponseWriter, r *http.
if bytes, marshalError = MarshalXMLEntity(initResult); marshalError != nil {
log.LogErrorf("createMultipleUploadHandler: marshal result fail, requestID(%v) err(%v)",
GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}
@ -73,20 +94,28 @@ func (o *ObjectNode) createMultipleUploadHandler(w http.ResponseWriter, r *http.
func (o *ObjectNode) uploadPartHandler(w http.ResponseWriter, r *http.Request) {
log.LogInfof("uploadPartHandler: upload part, requestID(%v) remote(%v)",
GetRequestID(r), r.RemoteAddr)
var (
err error
errorCode *ErrorCode
)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
// check args
params, _, object, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("uploadPartHandler: parse request parameters fail, requestID(%v) err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
return
}
var param = ParseRequestParam(r)
//// get upload id and part number
uploadId := params[ParamUploadId]
partNumber := params[ParamPartNumber]
uploadId := param.GetVar(ParamUploadId)
partNumber := param.GetVar(ParamPartNumber)
if uploadId == "" || partNumber == "" {
log.LogErrorf("uploadPartHandler: illegal uploadID or partNumber, requestID(%v)", GetRequestID(r))
_ = InvalidArgument.ServeResponse(w, r)
errorCode = &InvalidArgument
return
}
@ -94,16 +123,32 @@ func (o *ObjectNode) uploadPartHandler(w http.ResponseWriter, r *http.Request) {
if partNumberInt, err = strconv.ParseUint(partNumber, 10, 64); err != nil {
log.LogErrorf("uploadPartHandler: parse part number fail, requestID(%v) raw(%v) err(%v)",
GetRequestID(r), partNumber, err)
_ = InvalidArgument.ServeResponse(w, r)
errorCode = &InvalidArgument
return
}
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
if param.Object() == "" {
errorCode = &InvalidKey
return
}
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("uploadPartHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
// handle exception
var fsFileInfo *FSFileInfo
if fsFileInfo, err = vl.WritePart(object, uploadId, uint16(partNumberInt), r.Body); err != nil {
if fsFileInfo, err = vol.WritePart(param.Object(), uploadId, uint16(partNumberInt), r.Body); err != nil {
log.LogErrorf("uploadPartHandler: write part fail, requestID(%v) err(%v)", GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}
log.LogDebugf("uploadPartHandler: write part, requestID(%v) fsFileInfo(%v)", GetRequestID(r), fsFileInfo)
@ -119,17 +164,24 @@ func (o *ObjectNode) uploadPartHandler(w http.ResponseWriter, r *http.Request) {
func (o *ObjectNode) listPartsHandler(w http.ResponseWriter, r *http.Request) {
log.LogInfof("listPartsHandler: list parts, requestID(%v) remote(%v)", GetRequestID(r), r.RemoteAddr)
// check args
params, bucket, object, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("listPartsHandler: parse request parameters fail, requestID(%v) err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
return
}
//// get upload id and part number
uploadId := params[ParamUploadId]
maxParts := params[ParamMaxParts]
partNoMarker := params[ParamPartNoMarker]
var (
err error
errorCode *ErrorCode
)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
var param = ParseRequestParam(r)
// get upload id and part number
uploadId := param.GetVar(ParamUploadId)
maxParts := param.GetVar(ParamMaxParts)
partNoMarker := param.GetVar(ParamPartNoMarker)
var maxPartsInt uint64
var partNoMarkerInt uint64
@ -163,27 +215,43 @@ func (o *ObjectNode) listPartsHandler(w http.ResponseWriter, r *http.Request) {
partNoMarkerInt = res
}
fsParts, nextMarker, isTruncated, err := vl.ListParts(object, uploadId, maxPartsInt, partNoMarkerInt)
if err != nil {
log.LogErrorf("listPartsHandler: volume list parts fail, requestID(%v) uploadID(%v) maxParts(%v) partNoMarker(%v) err(%v)",
GetRequestID(r), uploadId, maxPartsInt, partNoMarkerInt, err)
_ = InternalError.ServeResponse(w, r)
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
log.LogDebugf("listPartsHandler: volume list parts, "+
if param.Object() == "" {
errorCode = &InvalidKey
return
}
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("listPartsHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
fsParts, nextMarker, isTruncated, err := vol.ListParts(param.Object(), uploadId, maxPartsInt, partNoMarkerInt)
if err != nil {
log.LogErrorf("listPartsHandler: Volume list parts fail, requestID(%v) uploadID(%v) maxParts(%v) partNoMarker(%v) err(%v)",
GetRequestID(r), uploadId, maxPartsInt, partNoMarkerInt, err)
errorCode = &InternalError
return
}
log.LogDebugf("listPartsHandler: Volume list parts, "+
"requestID(%v) uploadID(%v) maxParts(%v) partNoMarker(%v) numFSParts(%v) nextMarker(%v) isTruncated(%v)",
GetRequestID(r), uploadId, maxPartsInt, partNoMarkerInt, len(fsParts), nextMarker, isTruncated)
// get owner
accessKey, _ := vl.OSSSecure()
bucketOwner := NewBucketOwner(accessKey)
bucketOwner := NewBucketOwner(param.accessKey)
// get parts
parts := NewParts(fsParts)
listPartsResult := ListPartsResult{
Bucket: bucket,
Key: object,
Bucket: param.Bucket(),
Key: param.Object(),
UploadId: uploadId,
StorageClass: StorageClassStandard,
NextMarker: int(nextMarker),
@ -198,7 +266,7 @@ func (o *ObjectNode) listPartsHandler(w http.ResponseWriter, r *http.Request) {
if bytes, marshalError = MarshalXMLEntity(listPartsResult); marshalError != nil {
log.LogErrorf("listPartsHandler: marshal result fail, requestID(%v) err(%v)",
GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}
@ -217,35 +285,59 @@ func (o *ObjectNode) listPartsHandler(w http.ResponseWriter, r *http.Request) {
func (o *ObjectNode) completeMultipartUploadHandler(w http.ResponseWriter, r *http.Request) {
log.LogInfof("completeMultipartUploadHandler: complete multiple upload, requestID(%v) remote(%v)", GetRequestID(r), r.RemoteAddr)
params, bucket, object, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("completeMultipartUploadHandler: parse request params fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
return
}
var (
err error
errorCode *ErrorCode
)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
var param = ParseRequestParam(r)
// get upload id and part number
uploadId := params[ParamUploadId]
uploadId := param.GetVar(ParamUploadId)
if uploadId == "" {
log.LogErrorf("completeMultipartUploadHandler: non upload ID specified: requestID(%v)", GetRequestID(r))
_ = InvalidArgument.ServeResponse(w, r)
errorCode = &InvalidArgument
return
}
fsFileInfo, err := vl.CompleteMultipart(object, uploadId)
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
if param.Object() == "" {
errorCode = &InvalidKey
return
}
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("completeMultipartUploadHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
fsFileInfo, err := vol.CompleteMultipart(param.Object(), uploadId)
if err != nil {
log.LogErrorf("completeMultipartUploadHandler: complete multipart fail, requestID(%v) uploadID(%v) err(%v)",
GetRequestID(r), uploadId, err)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}
log.LogDebugf("completeMultipartUploadHandler: complete multipart, requestID(%v) uploadID(%v) path(%v)",
GetRequestID(r), uploadId, object)
GetRequestID(r), uploadId, param.Object())
// write response
completeResult := CompleteMultipartResult{
Bucket: bucket,
Key: object,
Bucket: param.Bucket(),
Key: param.Object(),
ETag: fsFileInfo.ETag,
}
@ -253,7 +345,7 @@ func (o *ObjectNode) completeMultipartUploadHandler(w http.ResponseWriter, r *ht
var marshalError error
if bytes, marshalError = MarshalXMLEntity(completeResult); marshalError != nil {
log.LogErrorf("completeMultipartUploadHandler: marshal result fail, requestID(%v) err(%v)", GetRequestID(r), marshalError)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}
@ -272,23 +364,50 @@ func (o *ObjectNode) completeMultipartUploadHandler(w http.ResponseWriter, r *ht
func (o *ObjectNode) abortMultipartUploadHandler(w http.ResponseWriter, r *http.Request) {
log.LogInfof("abortMultipartUploadHandler: abort multiple upload, requestID(%v) remote(%v)", GetRequestID(r), r.RemoteAddr)
var (
err error
errorCode *ErrorCode
)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
// check args
params, _, object, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("abortMultipartUploadHandler: parse request parameters fail, requestID(%v) err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
var param = ParseRequestParam(r)
uploadId := param.GetVar(ParamUploadId)
if uploadId == "" {
errorCode = &InvalidArgument
return
}
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
if param.Object() == "" {
errorCode = &InvalidKey
return
}
uploadId := params["uploadId"]
//// Abort multipart upload
err = vl.AbortMultipart(object, uploadId)
if err != nil {
log.LogErrorf("abortMultipartUploadHandler: volume abort multipart fail, requestID(%v) uploadID(%v) err(%v)", GetRequestID(r), uploadId, err)
_ = InternalError.ServeResponse(w, r)
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("abortMultipartUploadHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
log.LogDebugf("abortMultipartUploadHandler: volume abort multipart, requestID(%v) uploadID(%v) path(%v)", GetRequestID(r), uploadId, object)
// Abort multipart upload
if err = vol.AbortMultipart(param.Object(), uploadId); err != nil {
log.LogErrorf("abortMultipartUploadHandler: Volume abort multipart fail, requestID(%v) uploadID(%v) err(%v)", GetRequestID(r), uploadId, err)
errorCode = &InternalError
return
}
log.LogDebugf("abortMultipartUploadHandler: Volume abort multipart, requestID(%v) uploadID(%v) path(%v)", GetRequestID(r), uploadId, param.Object())
return
}
@ -296,19 +415,27 @@ func (o *ObjectNode) abortMultipartUploadHandler(w http.ResponseWriter, r *http.
// API reference: https://docs.aws.amazon.com/AmazonS3/latest/API/API_ListMultipartUploads.html
func (o *ObjectNode) listMultipartUploadsHandler(w http.ResponseWriter, r *http.Request) {
log.LogInfof("abortMultipartUploadHandler: list multipart uploads, requestID(%v) remote(%v)", GetRequestID(r), r.RemoteAddr)
// check args
params, bucket, _, vl, err := o.parseRequestParams(r)
if err != nil {
log.LogErrorf("listMultipartUploadsHandler: parse request parameters fail, requestID(%v) err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
}
var (
err error
errorCode *ErrorCode
)
defer func() {
if errorCode != nil {
_ = errorCode.ServeResponse(w, r)
return
}
}()
var param = ParseRequestParam(r)
// get list uploads parameter
prefix := params[ParamPrefix]
keyMarker := params[ParamKeyMarker]
delimiter := params[ParamPartDelimiter]
maxUploads := params[ParamPartMaxUploads]
uploadIdMarker := params[ParamUploadIdMarker]
prefix := param.GetVar(ParamPrefix)
keyMarker := param.GetVar(ParamKeyMarker)
delimiter := param.GetVar(ParamPartDelimiter)
maxUploads := param.GetVar(ParamPartMaxUploads)
uploadIdMarker := param.GetVar(ParamUploadIdMarker)
var maxUploadsInt uint64
if maxUploads == "" {
@ -325,15 +452,27 @@ func (o *ObjectNode) listMultipartUploadsHandler(w http.ResponseWriter, r *http.
}
}
fsUploads, nextKeyMarker, nextUploadIdMarker, IsTruncated, prefixes, err := vl.ListMultipartUploads(prefix, delimiter, keyMarker, uploadIdMarker, maxUploadsInt)
if err != nil {
log.LogErrorf("listMultipartUploadsHandler: volume list multipart uploads fail: requestID(%v), err(%v)", GetRequestID(r), err)
_ = NoSuchBucket.ServeResponse(w, r)
if param.Bucket() == "" {
errorCode = &InvalidBucketName
return
}
accessKey, _ := vl.OSSSecure()
uploads := NewUploads(fsUploads, accessKey)
var vol *Volume
if vol, err = o.vm.Volume(param.Bucket()); err != nil {
log.LogErrorf("listMultipartUploadsHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
fsUploads, nextKeyMarker, nextUploadIdMarker, IsTruncated, prefixes, err := vol.ListMultipartUploads(prefix, delimiter, keyMarker, uploadIdMarker, maxUploadsInt)
if err != nil {
log.LogErrorf("listMultipartUploadsHandler: Volume list multipart uploads fail: requestID(%v), err(%v)", GetRequestID(r), err)
errorCode = &NoSuchBucket
return
}
uploads := NewUploads(fsUploads, param.AccessKey())
var commonPrefixes = make([]*CommonPrefix, 0)
for _, prefix := range prefixes {
@ -344,7 +483,7 @@ func (o *ObjectNode) listMultipartUploadsHandler(w http.ResponseWriter, r *http.
}
listUploadsResult := ListUploadsResult{
Bucket: bucket,
Bucket: param.Bucket(),
KeyMarker: keyMarker,
UploadIdMarker: uploadIdMarker,
NextKeyMarker: nextKeyMarker,
@ -361,7 +500,7 @@ func (o *ObjectNode) listMultipartUploadsHandler(w http.ResponseWriter, r *http.
var marshalError error
if bytes, marshalError = MarshalXMLEntity(listUploadsResult); marshalError != nil {
log.LogErrorf("listMultipartUploadsHandler: marshal xml entity fail: requestID(%v) err(%v)", GetRequestID(r), err)
_ = InternalError.ServeResponse(w, r)
errorCode = &InternalError
return
}

File diff suppressed because it is too large Load Diff

View File

@ -59,6 +59,9 @@ func (o *ObjectNode) traceMiddleware(next http.Handler) http.Handler {
SetRequestID(r, requestID)
w.Header().Set(HeaderNameRequestId, requestID)
var action = ActionFromRouteName(mux.CurrentRoute(r).GetName())
SetRequestAction(r, action)
var startTime = time.Now()
next.ServeHTTP(w, r)
@ -79,7 +82,7 @@ func (o *ObjectNode) traceMiddleware(next http.Handler) http.Handler {
" requestID(%v) host(%v) method(%v) url(%v)\n"+
" header(%v)\n"+
" remote(%v) cost(%v)",
ActionFromRouteName(mux.CurrentRoute(r).GetName()).String(),
action,
requestID, r.Host, r.Method, r.URL.String(),
headerToString(r.Header),
getRequestIP(r), time.Since(startTime))
@ -155,7 +158,7 @@ func (o *ObjectNode) policyCheckMiddleware(next http.Handler) http.Handler {
next.ServeHTTP(w, r)
return
}
wrappedNext := o.policyCheck(next.ServeHTTP, action)
wrappedNext := o.policyCheck(next.ServeHTTP)
wrappedNext.ServeHTTP(w, r)
return
})

View File

@ -257,16 +257,13 @@ func (o *ObjectNode) validateUrlBySignatureAlgorithmV2(r *http.Request) (bool, e
return false, nil
}
params, _, _, vl, err := o.parseRequestParams(r)
if err != nil || vl == nil {
log.LogInfof("check PresignedSignatureV2 error: %v %v", err, vl)
return false, err
}
accessKey := params["accessKey"]
signature := params["signature"]
expires := params["expires"]
var param = ParseRequestParam(r)
accessKey := param.GetVar("accessKey")
signature := param.GetVar("signature")
expires := param.GetVar("expires")
if accessKey == "" || signature == "" || expires == "" {
log.LogInfof("validateUrlBySignatureAlgorithmV2 params not valid: %v", params)
log.LogInfof("validateUrlBySignatureAlgorithmV2: incomplete authentication information: requestID(%v)",
GetRequestID(r))
return false, nil
}

View File

@ -64,6 +64,7 @@ const (
ParamFetchOwner = "fetch-owner"
ParamMaxKeys = "max-keys"
ParamStartAfter = "start-after"
ParamKey = "key"
ParamMaxParts = "max-parts"
ParamUploadIdMarker = "upload-id-marker"

View File

@ -15,25 +15,11 @@
package objectnode
import (
"io"
"os"
"sort"
"time"
"github.com/chubaofs/chubaofs/proto"
"github.com/chubaofs/chubaofs/sdk/master"
)
type VolumeManager interface {
Volume(volName string) (Volume, error)
Release(volName string)
GetStore() (Store, error)
InitStore(s Store)
InitMasterClient(masters []string, useSSL bool)
GetMasterClient() (*master.MasterClient, error)
Close()
}
type FSFileInfo struct {
Path string
Size int64
@ -73,41 +59,3 @@ type FSPart struct {
ETag string
Size int
}
type Volume interface {
OSSSecure() (accessKey, secretKey string) //todo delete
OSSMeta() *OSSMeta
// ListFiles return an FileInfo slice of specified volume, like read dir for hole volume.
// The result will be ordered by full path.
ListFilesV1(request *ListBucketRequestV1) ([]*FSFileInfo, string, bool, []string, error)
ListFilesV2(request *ListBucketRequestV2) ([]*FSFileInfo, uint64, string, bool, []string, error)
// PutObject create file in specified volume with specified path.
WriteFile(path string, reader io.Reader) (*FSFileInfo, error)
// DeleteFile delete specified file from specified volume. If target is not exists then returns error.
DeleteFile(path string) error
FileInfo(path string) (*FSFileInfo, error)
// operation about multipart uploads
InitMultipart(path string) (multipartID string, err error)
WritePart(path, multipartID string, partId uint16, reader io.Reader) (*FSFileInfo, error)
ListParts(path, multipartID string, maxParts, partNumberMarker uint64) ([]*FSPart, uint64, bool, error)
CompleteMultipart(path, multipartID string) (*FSFileInfo, error)
AbortMultipart(path, multipartID string) error
ListMultipartUploads(prefix, delimiter, keyMarker, uploadIdMarker string, maxUploads uint64) ([]*FSUpload, string, string, bool, []string, error)
ReadFile(path string, writer io.Writer, offset, size uint64) error
CopyFile(path, sourcePath string) (*FSFileInfo, error)
SetXAttr(path string, key string, data []byte) error
GetXAttr(path string, key string) (*proto.XAttrInfo, error)
DeleteXAttr(path string, key string) error
ListXAttrs(path string) (info *proto.XAttrInfo, err error)
Close() error
}

View File

@ -22,32 +22,32 @@ import (
"github.com/chubaofs/chubaofs/util/log"
)
type volumeManager struct {
type VolumeManager struct {
masters []string
mc *master.MasterClient
volumes map[string]*volume // volume key -> vol
volumes map[string]*Volume // Volume key -> vol
volMu sync.RWMutex
store Store
closeOnce sync.Once
}
func (m *volumeManager) Release(volName string) {
func (m *VolumeManager) Release(volName string) {
m.volMu.Lock()
defer m.volMu.Unlock()
delete(m.volumes, volName)
}
func (m *volumeManager) ReleaseAll() {
func (m *VolumeManager) ReleaseAll() {
panic("implement me")
}
func (m *volumeManager) Volume(volName string) (Volume, error) {
func (m *VolumeManager) Volume(volName string) (*Volume, error) {
return m.loadVolume(volName)
}
func (m *volumeManager) loadVolume(volName string) (*volume, error) {
func (m *VolumeManager) loadVolume(volName string) (*Volume, error) {
var err error
var volume *volume
var volume *Volume
var exist bool
m.volMu.RLock()
volume, exist = m.volumes[volName]
@ -64,7 +64,7 @@ func (m *volumeManager) loadVolume(volName string) (*volume, error) {
return nil, err
}
ak, sk := volume.OSSSecure()
log.LogDebugf("[loadVolume] load volume: Name[%v] AccessKey[%v] SecretKey[%v]", volName, ak, sk)
log.LogDebugf("[loadVolume] load Volume: Name[%v] AccessKey[%v] SecretKey[%v]", volName, ak, sk)
m.volumes[volName] = volume
volume.vm = m
m.volMu.Unlock()
@ -76,42 +76,42 @@ func (m *volumeManager) loadVolume(volName string) (*volume, error) {
}
// Release all
func (m *volumeManager) Close() {
func (m *VolumeManager) Close() {
m.volMu.Lock()
defer m.volMu.Unlock()
for volKey, vol := range m.volumes {
_ = vol.Close()
log.LogDebugf("release volume %v", volKey)
log.LogDebugf("release Volume %v", volKey)
}
m.volumes = make(map[string]*volume)
m.volumes = make(map[string]*Volume)
}
func (m *volumeManager) InitStore(s Store) {
func (m *VolumeManager) InitStore(s Store) {
s.Init(m)
m.store = s
}
func (m *volumeManager) GetStore() (Store, error) {
func (m *VolumeManager) GetStore() (Store, error) {
if m.store == nil {
return nil, errors.New("store not init")
}
return m.store, nil
}
func (m *volumeManager) InitMasterClient(masters []string, useSSL bool) {
func (m *VolumeManager) InitMasterClient(masters []string, useSSL bool) {
m.mc = master.NewMasterClient(masters, useSSL)
}
func (m *volumeManager) GetMasterClient() (*master.MasterClient, error) {
func (m *VolumeManager) GetMasterClient() (*master.MasterClient, error) {
if m.mc == nil {
return nil, errors.New("master client not init")
}
return m.mc, nil
}
func NewVolumeManager(masters []string) VolumeManager {
vc := &volumeManager{
volumes: make(map[string]*volume),
func NewVolumeManager(masters []string) *VolumeManager {
vc := &VolumeManager{
volumes: make(map[string]*Volume),
masters: masters,
}
return vc

View File

@ -19,7 +19,7 @@ type MetaStore interface {
// MetaStore
type Store interface {
Init(vm *volumeManager)
Init(vm *VolumeManager)
Put(ns, obj, key string, data []byte) error
Get(ns, obj, key string) (data []byte, err error)
List(ns, obj string) (data [][]byte, err error)

View File

@ -23,10 +23,10 @@ const (
)
type objectStore struct {
vm *volumeManager
vm *VolumeManager
}
func (s *objectStore) Init(vm *volumeManager) {
func (s *objectStore) Init(vm *VolumeManager) {
s.vm = vm
//TODO: init meta dir
}

View File

@ -27,14 +27,14 @@ const (
)
type xattrStore struct {
vm *volumeManager //vol *volume
vm *VolumeManager //vol *Volume
}
func (s *xattrStore) Init(vm *volumeManager) {
func (s *xattrStore) Init(vm *VolumeManager) {
s.vm = vm
}
func (s *xattrStore) getInode(vol, path string) (*volume, uint64, error) {
func (s *xattrStore) getInode(vol, path string) (*Volume, uint64, error) {
v, err := s.vm.loadVolume(vol)
if err != nil {
return nil, 0, err
@ -69,7 +69,7 @@ func (s *xattrStore) Put(vol, path, key string, data []byte) (err error) {
}
func (s *xattrStore) Get(vol, path, key string) (val []byte, err error) {
var v *volume
var v *Volume
v, err = s.vm.loadVolume(vol)
if err != nil {
return

View File

@ -29,9 +29,6 @@ const (
OSSMetaUpdateDuration = time.Duration(time.Second * 30)
)
// Used to validate interface implementation
var _ Volume = &volume{}
type OSSMeta struct {
policy *Policy
acl *AccessControlPolicy
@ -39,28 +36,28 @@ type OSSMeta struct {
aclLock sync.RWMutex
}
func (v *volume) loadPolicy() (p *Policy) {
func (v *Volume) loadPolicy() (p *Policy) {
v.om.policyLock.RLock()
p = v.om.policy
v.om.policyLock.RUnlock()
return
}
func (v *volume) storePolicy(p *Policy) {
func (v *Volume) storePolicy(p *Policy) {
v.om.policyLock.Lock()
v.om.policy = p
v.om.policyLock.Unlock()
return
}
func (v *volume) loadACL() (p *AccessControlPolicy) {
func (v *Volume) loadACL() (p *AccessControlPolicy) {
v.om.aclLock.RLock()
p = v.om.acl
v.om.aclLock.RUnlock()
return
}
func (v *volume) storeACL(p *AccessControlPolicy) {
func (v *Volume) storeACL(p *AccessControlPolicy) {
v.om.aclLock.Lock()
v.om.acl = p
v.om.aclLock.Unlock()
@ -68,10 +65,10 @@ func (v *volume) storeACL(p *AccessControlPolicy) {
}
// This struct is implementation of Volume interface
type volume struct {
type Volume struct {
mw *meta.MetaWrapper
ec *stream.ExtentClient
vm *volumeManager
vm *VolumeManager
name string
om *OSSMeta
ticker *time.Ticker
@ -82,7 +79,7 @@ type volume struct {
closedCh chan struct{}
}
func (v *volume) syncOSSMeta() {
func (v *Volume) syncOSSMeta() {
defer v.ticker.Stop()
v.ticker = time.NewTicker(OSSMetaUpdateDuration)
for {
@ -96,7 +93,7 @@ func (v *volume) syncOSSMeta() {
}
}
func (v *volume) stopOSSMetaSync() {
func (v *Volume) stopOSSMetaSync() {
v.closingCh <- struct{}{}
select {
@ -105,8 +102,8 @@ func (v *volume) stopOSSMetaSync() {
}
}
// update volume meta info
func (v *volume) loadOSSMeta() {
// update Volume meta info
func (v *Volume) loadOSSMeta() {
policy, _ := v.loadBucketPolicy()
if policy != nil {
v.storePolicy(policy)
@ -118,8 +115,12 @@ func (v *volume) loadOSSMeta() {
}
}
func (v *Volume) Name() string {
return v.name
}
// load bucket policy from vm
func (v *volume) loadBucketPolicy() (policy *Policy, err error) {
func (v *Volume) loadBucketPolicy() (policy *Policy, err error) {
var store Store
store, err = v.vm.GetStore()
if err != nil {
@ -128,7 +129,7 @@ func (v *volume) loadBucketPolicy() (policy *Policy, err error) {
var data []byte
data, err = store.Get(v.name, bucketRootPath, XAttrKeyOSSPolicy)
if err != nil {
log.LogErrorf("loadBucketPolicy: load bucket policy fail: volume(%v) err(%v)", v.name, err)
log.LogErrorf("loadBucketPolicy: load bucket policy fail: Volume(%v) err(%v)", v.name, err)
return
}
policy = &Policy{}
@ -138,7 +139,7 @@ func (v *volume) loadBucketPolicy() (policy *Policy, err error) {
return
}
func (v *volume) loadBucketACL() (*AccessControlPolicy, error) {
func (v *Volume) loadBucketACL() (*AccessControlPolicy, error) {
store, err1 := v.vm.GetStore()
if err1 != nil {
return nil, err1
@ -155,11 +156,11 @@ func (v *volume) loadBucketACL() (*AccessControlPolicy, error) {
return acl, nil
}
func (v *volume) OSSMeta() *OSSMeta {
func (v *Volume) OSSMeta() *OSSMeta {
return v.om
}
func (v *volume) getInodeFromPath(path string) (inode uint64, err error) {
func (v *Volume) getInodeFromPath(path string) (inode uint64, err error) {
if path == "/" {
return volumeRootInode, nil
}
@ -191,7 +192,7 @@ func (v *volume) getInodeFromPath(path string) (inode uint64, err error) {
return
}
func (v *volume) SetXAttr(path string, key string, data []byte) error {
func (v *Volume) SetXAttr(path string, key string, data []byte) error {
var err error
var inode uint64
if inode, err = v.getInodeFromPath(path); err != nil && err != syscall.ENOENT {
@ -212,7 +213,7 @@ func (v *volume) SetXAttr(path string, key string, data []byte) error {
return v.mw.XAttrSet_ll(inode, []byte(key), data)
}
func (v *volume) GetXAttr(path string, key string) (info *proto.XAttrInfo, err error) {
func (v *Volume) GetXAttr(path string, key string) (info *proto.XAttrInfo, err error) {
var inode uint64
inode, err = v.getInodeFromPath(path)
if err != nil {
@ -225,7 +226,7 @@ func (v *volume) GetXAttr(path string, key string) (info *proto.XAttrInfo, err e
return
}
func (v *volume) DeleteXAttr(path string, key string) (err error) {
func (v *Volume) DeleteXAttr(path string, key string) (err error) {
inode, err1 := v.getInodeFromPath(path)
if err1 != nil {
err = err1
@ -238,7 +239,7 @@ func (v *volume) DeleteXAttr(path string, key string) (err error) {
return
}
func (v *volume) ListXAttrs(path string) (info *proto.XAttrInfo, err error) {
func (v *Volume) ListXAttrs(path string) (info *proto.XAttrInfo, err error) {
var inode uint64
inode, err = v.getInodeFromPath(path)
if err != nil {
@ -251,11 +252,11 @@ func (v *volume) ListXAttrs(path string) (info *proto.XAttrInfo, err error) {
return
}
func (v *volume) OSSSecure() (accessKey, secretKey string) {
func (v *Volume) OSSSecure() (accessKey, secretKey string) {
return v.mw.OSSSecure()
}
func (v *volume) ListFilesV1(request *ListBucketRequestV1) ([]*FSFileInfo, string, bool, []string, error) {
func (v *Volume) ListFilesV1(request *ListBucketRequestV1) ([]*FSFileInfo, string, bool, []string, error) {
//prefix, delimiter, marker string, maxKeys uint64
marker := request.marker
@ -287,7 +288,7 @@ func (v *volume) ListFilesV1(request *ListBucketRequestV1) ([]*FSFileInfo, strin
return infos, nextMarker, isTruncated, prefixes, nil
}
func (v *volume) ListFilesV2(request *ListBucketRequestV2) ([]*FSFileInfo, uint64, string, bool, []string, error) {
func (v *Volume) ListFilesV2(request *ListBucketRequestV2) ([]*FSFileInfo, uint64, string, bool, []string, error) {
delimiter := request.delimiter
maxKeys := request.maxKeys
@ -322,7 +323,7 @@ func (v *volume) ListFilesV2(request *ListBucketRequestV2) ([]*FSFileInfo, uint6
return infos, keyCount, nextToken, isTruncated, prefixes, nil
}
func (v *volume) WriteFile(path string, reader io.Reader) (*FSFileInfo, error) {
func (v *Volume) WriteFile(path string, reader io.Reader) (*FSFileInfo, error) {
var err error
var fInfo *FSFileInfo
@ -430,7 +431,7 @@ func (v *volume) WriteFile(path string, reader io.Reader) (*FSFileInfo, error) {
return fInfo, nil
}
func (v *volume) DeleteFile(path string) error {
func (v *Volume) DeleteFile(path string) error {
var err error
dirs, filename := splitPath(path)
@ -465,7 +466,7 @@ func (v *volume) DeleteFile(path string) error {
return nil
}
func (v *volume) InitMultipart(path string) (multipartID string, err error) {
func (v *Volume) InitMultipart(path string) (multipartID string, err error) {
// Invoke meta service to get a session id
// Create parent path
@ -486,7 +487,7 @@ func (v *volume) InitMultipart(path string) (multipartID string, err error) {
return multipartID, nil
}
func (v *volume) WritePart(path string, multipartId string, partId uint16, reader io.Reader) (*FSFileInfo, error) {
func (v *Volume) WritePart(path string, multipartId string, partId uint16, reader io.Reader) (*FSFileInfo, error) {
var parentId uint64
var err error
var fInfo *FSFileInfo
@ -584,7 +585,7 @@ func (v *volume) WritePart(path string, multipartId string, partId uint16, reade
return fInfo, nil
}
func (v *volume) AbortMultipart(path string, multipartID string) (err error) {
func (v *Volume) AbortMultipart(path string, multipartID string) (err error) {
// TODO: cleanup data while abort multipart
var parentId uint64
@ -627,7 +628,7 @@ func (v *volume) AbortMultipart(path string, multipartID string) (err error) {
return nil
}
func (v *volume) CompleteMultipart(path string, multipartID string) (fsFileInfo *FSFileInfo, err error) {
func (v *Volume) CompleteMultipart(path string, multipartID string) (fsFileInfo *FSFileInfo, err error) {
const mode = 0600
@ -776,7 +777,7 @@ func (v *volume) CompleteMultipart(path string, multipartID string) (fsFileInfo
return fInfo, nil
}
func (v *volume) appendInodeHash(h hash.Hash, inode uint64, total uint64, preAllocatedBuf []byte) (err error) {
func (v *Volume) appendInodeHash(h hash.Hash, inode uint64, total uint64, preAllocatedBuf []byte) (err error) {
if err = v.ec.OpenStream(inode); err != nil {
log.LogErrorf("appendInodeHash: data open stream fail: inode(%v) err(%v)",
inode, err)
@ -825,7 +826,7 @@ func (v *volume) appendInodeHash(h hash.Hash, inode uint64, total uint64, preAll
return
}
func (v *volume) applyInodeToNewDentry(parentID uint64, name string, inode uint64) (err error) {
func (v *Volume) applyInodeToNewDentry(parentID uint64, name string, inode uint64) (err error) {
const mode = 0600
if err = v.mw.DentryCreate_ll(parentID, name, inode, mode); err != nil {
log.LogErrorf("applyInodeToNewDentry: meta dentry create fail: parentID(%v) name(%v) inode(%v) mode(%v) err(%v)",
@ -835,7 +836,7 @@ func (v *volume) applyInodeToNewDentry(parentID uint64, name string, inode uint6
return
}
func (v *volume) applyInodeToExistDentry(parentID uint64, name string, inode uint64) (err error) {
func (v *Volume) applyInodeToExistDentry(parentID uint64, name string, inode uint64) (err error) {
var oldInode uint64
oldInode, err = v.mw.DentryUpdate_ll(parentID, name, inode)
if err != nil {
@ -879,7 +880,7 @@ func (v *volume) applyInodeToExistDentry(parentID uint64, name string, inode uin
return
}
func (v *volume) ReadFile(path string, writer io.Writer, offset, size uint64) error {
func (v *Volume) ReadFile(path string, writer io.Writer, offset, size uint64) error {
var err error
dirs, filename := splitPath(path)
@ -950,7 +951,7 @@ func (v *volume) ReadFile(path string, writer io.Writer, offset, size uint64) er
return nil
}
func (v *volume) FileInfo(path string) (info *FSFileInfo, err error) {
func (v *Volume) FileInfo(path string) (info *FSFileInfo, err error) {
dirs, filename := splitPath(path)
// process path
@ -991,7 +992,7 @@ func (v *volume) FileInfo(path string) (info *FSFileInfo, err error) {
return
}
func (v *volume) Close() error {
func (v *Volume) Close() error {
v.closeOnce.Do(func() {
v.mw.Close()
v.stopOSSMetaSync()
@ -999,7 +1000,7 @@ func (v *volume) Close() error {
return nil
}
func (v *volume) lookupDirectories(dirs []string, autoCreate bool) (inode uint64, err error) {
func (v *Volume) lookupDirectories(dirs []string, autoCreate bool) (inode uint64, err error) {
var parentId = rootIno
// check and create dirs
for _, dir := range dirs {
@ -1049,7 +1050,7 @@ func (v *volume) lookupDirectories(dirs []string, autoCreate bool) (inode uint64
return
}
func (v *volume) listFilesV1(prefix, marker, delimiter string, maxKeys uint64) (infos []*FSFileInfo, prefixes Prefixes, err error) {
func (v *Volume) listFilesV1(prefix, marker, delimiter string, maxKeys uint64) (infos []*FSFileInfo, prefixes Prefixes, err error) {
var prefixMap = PrefixMap(make(map[string]struct{}))
parentId, dirs, err := v.findParentId(prefix)
@ -1061,7 +1062,7 @@ func (v *volume) listFilesV1(prefix, marker, delimiter string, maxKeys uint64) (
// recursion call listDir method
infos, prefixMap, err = v.listDir(infos, prefixMap, parentId, maxKeys, dirs, prefix, marker, delimiter)
if err != nil {
log.LogErrorf("listFilesV1: volume list dir fail: volume(%v) err(%v)", v.name, err)
log.LogErrorf("listFilesV1: Volume list dir fail: Volume(%v) err(%v)", v.name, err)
return
}
@ -1073,13 +1074,13 @@ func (v *volume) listFilesV1(prefix, marker, delimiter string, maxKeys uint64) (
prefixes = prefixMap.Prefixes()
log.LogDebugf("listFilesV1: volume list dir: volume(%v) prefix(%v) marker(%v) delimiter(%v) maxKeys(%v) infos(%v) prefixes(%v)",
log.LogDebugf("listFilesV1: Volume list dir: Volume(%v) prefix(%v) marker(%v) delimiter(%v) maxKeys(%v) infos(%v) prefixes(%v)",
v.name, prefix, marker, delimiter, maxKeys, len(infos), len(prefixes))
return
}
func (v *volume) listFilesV2(prefix, startAfter, contToken, delimiter string, maxKeys uint64) (infos []*FSFileInfo, prefixes Prefixes, err error) {
func (v *Volume) listFilesV2(prefix, startAfter, contToken, delimiter string, maxKeys uint64) (infos []*FSFileInfo, prefixes Prefixes, err error) {
var prefixMap = PrefixMap(make(map[string]struct{}))
var marker string
@ -1099,7 +1100,7 @@ func (v *volume) listFilesV2(prefix, startAfter, contToken, delimiter string, ma
// recursion call listDir method
infos, prefixMap, err = v.listDir(infos, prefixMap, parentId, maxKeys, dirs, prefix, marker, delimiter)
if err != nil {
log.LogErrorf("listFilesV2: volume list dir fail, volume(%v) err(%v)", v.name, err)
log.LogErrorf("listFilesV2: Volume list dir fail, Volume(%v) err(%v)", v.name, err)
return
}
@ -1112,13 +1113,13 @@ func (v *volume) listFilesV2(prefix, startAfter, contToken, delimiter string, ma
prefixes = prefixMap.Prefixes()
log.LogDebugf("listFilesV2: volume list dir: volume(%v) prefix(%v) marker(%v) delimiter(%v) maxKeys(%v) infos(%v) prefixes(%v)",
log.LogDebugf("listFilesV2: Volume list dir: Volume(%v) prefix(%v) marker(%v) delimiter(%v) maxKeys(%v) infos(%v) prefixes(%v)",
v.name, prefix, marker, delimiter, maxKeys, len(infos), len(prefixes))
return
}
func (v *volume) findParentId(prefix string) (inode uint64, prefixDirs []string, err error) {
func (v *Volume) findParentId(prefix string) (inode uint64, prefixDirs []string, err error) {
prefixDirs = make([]string, 0)
// if prefix and marker are both not empty, use marker
@ -1153,7 +1154,7 @@ func (v *volume) findParentId(prefix string) (inode uint64, prefixDirs []string,
return
}
func (v *volume) listDir(fileInfos []*FSFileInfo, prefixMap PrefixMap, parentId, maxKeys uint64, dirs []string, prefix, marker, delimiter string) ([]*FSFileInfo, PrefixMap, error) {
func (v *Volume) listDir(fileInfos []*FSFileInfo, prefixMap PrefixMap, parentId, maxKeys uint64, dirs []string, prefix, marker, delimiter string) ([]*FSFileInfo, PrefixMap, error) {
children, err := v.mw.ReadDir_ll(parentId)
if err != nil {
@ -1202,7 +1203,7 @@ func (v *volume) listDir(fileInfos []*FSFileInfo, prefixMap PrefixMap, parentId,
return fileInfos, prefixMap, nil
}
func (v *volume) supplyListFileInfo(fileInfos []*FSFileInfo) (err error) {
func (v *Volume) supplyListFileInfo(fileInfos []*FSFileInfo) (err error) {
// get all file info indoes
inoInfoMap := make(map[uint64]*proto.InodeInfo)
var inodes []uint64
@ -1241,7 +1242,7 @@ func (v *volume) supplyListFileInfo(fileInfos []*FSFileInfo) (err error) {
return
}
func (v *volume) ListMultipartUploads(prefix, delimiter, keyMarker string, multipartIdMarker string, maxUploads uint64) ([]*FSUpload, string, string, bool, []string, error) {
func (v *Volume) ListMultipartUploads(prefix, delimiter, keyMarker string, multipartIdMarker string, maxUploads uint64) ([]*FSUpload, string, string, bool, []string, error) {
sessions, err := v.mw.ListMultipart_ll(prefix, delimiter, keyMarker, multipartIdMarker, maxUploads)
if err != nil {
return nil, "", "", false, nil, err
@ -1294,7 +1295,7 @@ func (v *volume) ListMultipartUploads(prefix, delimiter, keyMarker string, multi
return uploads, NextMarker, NextSessionIdMarker, IsTruncated, prefixes, nil
}
func (v *volume) ListParts(path, sessionId string, maxParts, partNumberMarker uint64) (parts []*FSPart, nextMarker uint64, isTruncated bool, err error) {
func (v *Volume) ListParts(path, sessionId string, maxParts, partNumberMarker uint64) (parts []*FSPart, nextMarker uint64, isTruncated bool, err error) {
var parentId uint64
dirs, _ := splitPath(path)
// process path
@ -1332,7 +1333,7 @@ func (v *volume) ListParts(path, sessionId string, maxParts, partNumberMarker ui
return parts, nextMarker, isTruncated, nil
}
func (v *volume) CopyFile(targetPath, sourcePath string) (info *FSFileInfo, err error) {
func (v *Volume) CopyFile(targetPath, sourcePath string) (info *FSFileInfo, err error) {
sourceDirs, sourceFilename := splitPath(sourcePath)
// process source targetPath
@ -1385,7 +1386,7 @@ func (v *volume) CopyFile(targetPath, sourcePath string) (info *FSFileInfo, err
return
}
func (v *volume) copyFile(parentID uint64, newFileName string, sourceFileInode uint64, mode uint32) (info *proto.InodeInfo, err error) {
func (v *Volume) copyFile(parentID uint64, newFileName string, sourceFileInode uint64, mode uint32) (info *proto.InodeInfo, err error) {
if err = v.mw.DentryCreate_ll(parentID, newFileName, sourceFileInode, mode); err != nil {
return
@ -1396,7 +1397,7 @@ func (v *volume) copyFile(parentID uint64, newFileName string, sourceFileInode u
return
}
func newVolume(masters []string, vol string) (*volume, error) {
func newVolume(masters []string, vol string) (*Volume, error) {
var err error
opt := &proto.MountOptions{
Volname: vol,
@ -1420,7 +1421,7 @@ func newVolume(masters []string, vol string) (*volume, error) {
return nil, err
}
v := &volume{mw: mw, ec: ec, name: vol, om: new(OSSMeta)}
v := &Volume{mw: mw, ec: ec, name: vol, om: new(OSSMeta)}
go v.syncOSSMeta()
return v, nil
}

View File

@ -86,7 +86,7 @@ func parseArn(str string) (*Arn, error) {
}
// write bucket policy into store and update vol policy meta
func storeBucketPolicy(bytes []byte, vol *volume) (*Policy, error) {
func storeBucketPolicy(bytes []byte, vol *Volume) (*Policy, error) {
store, err1 := vol.vm.GetStore()
if err1 != nil {
return nil, err1
@ -171,11 +171,6 @@ func (p *Policy) IsAllowed(params *RequestParam) bool {
}
}
}
if params.isOwner {
return true
}
for _, s := range p.Statements {
if s.Effect == Allow {
if s.IsAllowed(params) {
@ -187,7 +182,7 @@ func (p *Policy) IsAllowed(params *RequestParam) bool {
return false
}
func (o *ObjectNode) policyCheck(f http.HandlerFunc, action Action) http.HandlerFunc {
func (o *ObjectNode) policyCheck(f http.HandlerFunc) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
var (
err error
@ -205,59 +200,75 @@ func (o *ObjectNode) policyCheck(f http.HandlerFunc, action Action) http.Handler
}
}()
param, err1 := o.parseRequestParam(r)
if err1 != nil {
err = err1
log.LogInfof("parse Request Param err %v", err)
return
}
if param.vol == nil {
log.LogInfof("vol is null")
param := ParseRequestParam(r)
if param.Bucket() == "" {
log.LogDebugf("policyCheck: no bucket specified: requestID(%v)", GetRequestID(r))
allowed = true
return
}
param.actions = action
//check policy and acl
acl := param.vol.loadACL()
policy := param.vol.loadPolicy()
//check ip policy
if policy != nil && !policy.IsEmpty() {
allowed = policy.IsAllowed(param)
if !allowed {
log.LogWarnf("policy not allowed %v", param)
var vol *Volume
var acl *AccessControlPolicy
var policy *Policy
var loadBucketMeta = func(bucket string) (err error) {
if vol, err = o.getVol(bucket); err != nil {
return
}
acl = vol.loadACL()
policy = vol.loadPolicy()
return
}
switch param.action {
case CreateBucketAction:
default:
if err = loadBucketMeta(param.Bucket()); err != nil {
log.LogErrorf("policyCheck: load bucket metadata fail: requestID(%v) err(%v)", GetRequestID(r), err)
allowed = false
ec = &NoSuchBucket
return
}
}
if acl != nil && !acl.IsAclEmpty() {
if vol != nil && policy != nil && !policy.IsEmpty() {
allowed = policy.IsAllowed(param)
if !allowed {
log.LogWarnf("policyCheck: bucket policy not allowed: requestID(%v) volume(%v)", GetRequestID(r), vol.Name())
return
}
}
if vol != nil && acl != nil && !acl.IsAclEmpty() {
allowed = acl.IsAllowed(param)
if !allowed {
log.LogWarnf("acl not allowed %v", param)
log.LogWarnf("policyCheck: bucket ACL not allowed: requestID(%v) volume(%v)", GetRequestID(r), vol.Name())
return
}
}
//check user policy
var akPolicy *proto.AKPolicy
if akPolicy, err = o.getAkInfo(param.accessKey); err != nil {
log.LogInfof("get user policy from master error: accessKey(%v), err(%v)", param.accessKey, err)
log.LogErrorf("policyCheck: load user policy from master fail: requestID(%v) accessKey(%v) err(%v)",
GetRequestID(r), param.AccessKey(), err)
allowed = false
return
}
if contains(akPolicy.Policy.OwnVol, param.bucket) {
if contains(akPolicy.Policy.OwnVols, param.bucket) {
allowed = true
return
}
if apis, exit := akPolicy.Policy.NoneOwnVol[param.bucket]; exit {
if !contains(apis, action.String()) {
if !contains(apis, param.Action().String()) {
allowed = false
log.LogWarnf("user policy not allowed %v", param)
log.LogWarnf("policyCheck: user policy not allowed: requestID(%v) accessKey(%v) action(%v)",
GetRequestID(r), param.AccessKey(), param.Action())
return
}
allowed = true
} else {
allowed = false
log.LogWarnf("user policy not allowed %v", param)
log.LogWarnf("policyCheck: user policy not allowed: requestID(%v) accessKey(%v) action(%v)",
GetRequestID(r), param.AccessKey(), param.Action())
return
}
}

View File

@ -167,12 +167,9 @@ func (s Statement) checkActions(p *RequestParam) bool {
if s.Actions.Empty() {
return true
}
for _, pa := range p.actions {
if s.Actions.ContainsWithAny(string(pa)) {
return true
}
if s.Actions.ContainsWithAny(p.Action().String()) {
return true
}
return false
}
@ -180,12 +177,9 @@ func (s Statement) checkNotActions(p *RequestParam) bool {
if s.NotActions.Empty() {
return true
}
for _, pa := range p.actions {
if s.NotActions.ContainsWithAny(string(pa)) {
return false
}
if s.NotActions.ContainsWithAny(p.Action().String()) {
return false
}
return true
}

View File

@ -226,7 +226,7 @@ func StringLikeFunc(reqParam *RequestParam, storeCondVals ConditionValues) bool
for k, storeVals := range storeCondVals {
key := TrimAwsPrefixKey(k)
canonicalKey := http.CanonicalHeaderKey(key)
if reqVals, ok := reqParam.condVals[canonicalKey]; ok {
if reqVals, ok := reqParam.conditionVars[canonicalKey]; ok {
for _, rv := range reqVals {
for sv, _ := range storeVals.values {
if match := patternMatch(rv, sv); match {
@ -248,13 +248,13 @@ func StringEqualsFunc(reqParam *RequestParam, storeCondVals ConditionValues) boo
for k, storeVals := range storeCondVals {
key := TrimAwsPrefixKey(k)
canonicalKey := http.CanonicalHeaderKey(key)
if reqVals, ok := reqParam.condVals[canonicalKey]; ok {
if reqVals, ok := reqParam.conditionVars[canonicalKey]; ok {
for _, rv := range reqVals {
if storeVals.Contains(rv) {
return true
}
}
} else if reqVals, ok := reqParam.condVals[key]; ok {
} else if reqVals, ok := reqParam.conditionVars[key]; ok {
for _, rv := range reqVals {
if storeVals.Contains(rv) {
return true
@ -295,7 +295,7 @@ func BoolFunc(p *RequestParam, policyCondtion ConditionValues) bool {
for condKey, condVal := range policyCondtion {
for vals, _ := range condVal.values {
val1, _ := strconv.ParseBool(vals)
if cond, ok := p.condVals[condKey]; ok {
if cond, ok := p.conditionVars[condKey]; ok {
for _, c := range cond {
val2, _ := strconv.ParseBool(c)
return val1 == val2
@ -314,7 +314,7 @@ func DateEqualsFunc(p *RequestParam, policyVals ConditionValues) bool {
if err != nil {
return false
}
if reqVals, ok := p.condVals[k]; ok {
if reqVals, ok := p.conditionVars[k]; ok {
for _, reqVal := range reqVals {
reqDate, err := time.Parse(AMZTimeFormat, reqVal)
if err != nil {
@ -340,7 +340,7 @@ func DateLessThanFunc(p *RequestParam, value ConditionValues) bool {
if err != nil {
return false
}
if reqVals, ok := p.condVals[k]; ok {
if reqVals, ok := p.conditionVars[k]; ok {
for _, reqVal := range reqVals {
reqDate, err := time.Parse(AMZTimeFormat, reqVal)
if err != nil {
@ -361,7 +361,7 @@ func DateLessThanEqualsFunc(p *RequestParam, value ConditionValues) bool {
if err != nil {
return false
}
if reqVals, ok := p.condVals[k]; ok {
if reqVals, ok := p.conditionVars[k]; ok {
for _, reqVal := range reqVals {
reqDate, err := time.Parse(AMZTimeFormat, reqVal)
if err != nil {

View File

@ -33,12 +33,18 @@ func (o *ObjectNode) getBucketPolicyHandler(w http.ResponseWriter, r *http.Reque
)
defer o.errorResponse(w, r, err, ec)
_, bucket, _, vol, err := o.parseRequestParams(r)
if bucket == "" {
var param = ParseRequestParam(r)
if param.Bucket() == "" {
ec = &InvalidBucketName
return
}
var vol *Volume
if vol, err = o.getVol(param.Bucket()); err != nil {
log.LogErrorf("getBucketPolicyHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
ec = &NoSuchBucket
return
}
ossMeta := vol.OSSMeta()
if ossMeta == nil {
ec = &InternalError
@ -52,7 +58,7 @@ func (o *ObjectNode) getBucketPolicyHandler(w http.ResponseWriter, r *http.Reque
return
}
w.Write(policyData)
_, _ = w.Write(policyData)
return
}
@ -64,8 +70,15 @@ func (o *ObjectNode) putBucketPolicyHandler(w http.ResponseWriter, r *http.Reque
)
defer o.errorResponse(w, r, err, ec)
_, bucket, _, vol, err := o.parseRequestParams(r)
if bucket == "" {
var param = ParseRequestParam(r)
if param.Bucket() == "" {
ec = &InvalidBucketName
return
}
var vol *Volume
if vol, err = o.getVol(param.Bucket()); err != nil {
log.LogErrorf("putBucketPolicyHandler: load volume fail: requestID(%v) err(%v)",
GetRequestID(r), err)
ec = &NoSuchBucket
return
}
@ -75,21 +88,28 @@ func (o *ObjectNode) putBucketPolicyHandler(w http.ResponseWriter, r *http.Reque
return
}
bytes, err2 := ioutil.ReadAll(r.Body)
if err2 != nil && err2 != io.EOF {
err = err2
log.LogInfof("read body err, %v", err)
var bytes []byte
bytes, err = ioutil.ReadAll(r.Body)
if err != nil && err != io.EOF {
log.LogErrorf("putBucketPolicyHandler: read request body fail: requestID(%v) err(%v)", GetRequestID(r), err)
ec = &ErrorCode{
ErrorCode: http.StatusText(http.StatusBadRequest),
ErrorMessage: err.Error(),
StatusCode: http.StatusBadRequest,
}
return
}
policy, err3 := storeBucketPolicy(bytes, vol)
if err3 != nil {
err = err3
log.LogErrorf("store policy err, %v", err)
var policy *Policy
policy, err = storeBucketPolicy(bytes, vol)
if err != nil {
log.LogErrorf("putBucketPolicyHandler: store policy fail: requestID(%v) err(%v)", GetRequestID(r), err)
ec = &InternalError
return
}
log.LogInfof("put bucket policy %v %v", bucket, policy)
log.LogInfof("putBucketPolicyHandler: put bucket policy: requestID(%v) volume(%v) policy(%v)",
GetRequestID(r), param.Bucket(), policy)
return
}

View File

@ -85,7 +85,7 @@ func (s Statement) checkPrincipal(p *RequestParam) bool {
return true
}
for _, principal := range s.Principal {
if principal.ContainsWild(p.account) {
if principal.ContainsWild(p.AccessKey()) {
return true
}
}

View File

@ -53,7 +53,7 @@ type ObjectNode struct {
listen string
region string
httpServer *http.Server
vm VolumeManager
vm *VolumeManager
mc *master.MasterClient
state uint32
wg sync.WaitGroup

View File

@ -15,9 +15,9 @@ type AKPolicy struct {
}
type UserPolicy struct {
OwnVol []string
OwnVols []string
NoneOwnVol map[string][]string // k: vol, v: apis
sync.RWMutex
mu sync.RWMutex
}
type VolAK struct {
@ -26,10 +26,36 @@ type VolAK struct {
sync.RWMutex
}
func (policy *UserPolicy) AddOwnVol(volume string) {
policy.mu.Lock()
defer policy.mu.Unlock()
for _, ownVol := range policy.OwnVols {
if ownVol == volume {
return
}
}
policy.OwnVols = append(policy.OwnVols, volume)
}
func (policy *UserPolicy) RemoveOwnVol(volume string) {
policy.mu.Lock()
defer policy.mu.Unlock()
for i, ownVol := range policy.OwnVols {
if ownVol == volume {
if i == len(policy.OwnVols)-1 {
policy.OwnVols = policy.OwnVols[:i]
return
}
policy.OwnVols = append(policy.OwnVols[:i], policy.OwnVols[i+1:]...)
return
}
}
}
func (policy *UserPolicy) Add(addPolicy *UserPolicy) {
policy.Lock()
defer policy.Unlock()
policy.OwnVol = append(policy.OwnVol, addPolicy.OwnVol...)
policy.mu.Lock()
defer policy.mu.Unlock()
policy.OwnVols = append(policy.OwnVols, addPolicy.OwnVols...)
for k, v := range addPolicy.NoneOwnVol {
if apis, ok := policy.NoneOwnVol[k]; ok {
policy.NoneOwnVol[k] = append(apis, addPolicy.NoneOwnVol[k]...)
@ -40,9 +66,9 @@ func (policy *UserPolicy) Add(addPolicy *UserPolicy) {
}
func (policy *UserPolicy) Delete(deletePolicy *UserPolicy) {
policy.Lock()
defer policy.Unlock()
policy.OwnVol = removeSlice(policy.OwnVol, deletePolicy.OwnVol)
policy.mu.Lock()
defer policy.mu.Unlock()
policy.OwnVols = removeSlice(policy.OwnVols, deletePolicy.OwnVols)
for k, v := range deletePolicy.NoneOwnVol {
if apis, ok := policy.NoneOwnVol[k]; ok {
policy.NoneOwnVol[k] = removeSlice(apis, v)
@ -67,11 +93,11 @@ func removeSlice(s []string, removeSlice []string) []string {
func CleanPolicy(policy *UserPolicy) (newUserPolicy *UserPolicy) {
m := make(map[string]bool)
newUserPolicy = &UserPolicy{OwnVol: make([]string, 0), NoneOwnVol: make(map[string][]string)}
for _, vol := range policy.OwnVol {
newUserPolicy = &UserPolicy{OwnVols: make([]string, 0), NoneOwnVol: make(map[string][]string)}
for _, vol := range policy.OwnVols {
if _, exit := m[vol]; !exit {
m[vol] = true
newUserPolicy.OwnVol = append(newUserPolicy.OwnVol, vol)
newUserPolicy.OwnVols = append(newUserPolicy.OwnVols, vol)
}
}
for vol, apis := range policy.NoneOwnVol {

View File

@ -99,6 +99,7 @@ type Ticket struct {
}
func NewMetaWrapper(opt *proto.MountOptions, validateOwner bool) (*MetaWrapper, error) {
var err error
mw := new(MetaWrapper)
mw.closeCh = make(chan struct{}, 1)
if opt.Authenticate {
@ -124,8 +125,12 @@ func NewMetaWrapper(opt *proto.MountOptions, validateOwner bool) (*MetaWrapper,
mw.partitions = make(map[uint64]*MetaPartition)
mw.ranges = btree.New(32)
mw.rwPartitions = make([]*MetaPartition, 0)
_ = mw.updateClusterInfo()
_ = mw.updateVolStatInfo()
if err = mw.updateClusterInfo(); err != nil {
return nil, err
}
if err = mw.updateVolStatInfo(); err != nil {
return nil, err
}
limit := MaxMountRetryLimit
retry: