From 4b2316208c92aea0b69f5ef2dd655feb8a989378 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 12 Aug 2019 10:34:51 +0800 Subject: [PATCH] feature: s3upload support multipart upload --- docs/bucket/uploadobject.yaml | 2 + docs/parameters/bucket.yaml | 18 ++- docs/schemas/bucket.yaml | 4 +- pkg/apis/compute/bucket.go | 1 + pkg/cloudprovider/objectstore.go | 86 +++++++++++++- pkg/compute/models/buckets.go | 28 ++++- pkg/image/models/images.go | 11 +- pkg/mcclient/modules/mod_images.go | 12 +- pkg/multicloud/aliyun/bucket.go | 124 +++++++++++++++++++- pkg/multicloud/aliyun/objects.go | 6 +- pkg/multicloud/aliyun/storagecache.go | 21 ++-- pkg/multicloud/aws/bucket.go | 127 ++++++++++++++++++--- pkg/multicloud/aws/s3object.go | 2 +- pkg/multicloud/aws/storagecache.go | 106 ++++++----------- pkg/multicloud/azure/blobobject.go | 2 +- pkg/multicloud/azure/storageaccount.go | 138 ++++++++++++++++++++++- pkg/multicloud/azure/storagecache.go | 6 +- pkg/multicloud/bucket_base.go | 8 ++ pkg/multicloud/huawei/bucket.go | 103 ++++++++++++++++- pkg/multicloud/huawei/object.go | 2 +- pkg/multicloud/huawei/storagecache.go | 50 ++++---- pkg/multicloud/objectstore/buckets.go | 72 +++++++++++- pkg/multicloud/objectstore/shell.go | 14 ++- pkg/multicloud/openstack/storagecache.go | 2 +- pkg/multicloud/qcloud/bucket.go | 91 ++++++++++++++- pkg/multicloud/qcloud/object.go | 8 +- pkg/multicloud/qcloud/storagecache.go | 6 +- pkg/multicloud/ucloud/storagecache.go | 4 +- pkg/multicloud/ucloud/ufile.go | 20 +++- pkg/multicloud/zstack/storagecache.go | 4 +- pkg/util/fileutils2/seeker.go | 81 +++++++++++++ pkg/util/fileutils2/seeker_test.go | 31 +++++ pkg/util/httputils/httputils.go | 6 +- 33 files changed, 1008 insertions(+), 188 deletions(-) create mode 100644 pkg/util/fileutils2/seeker.go create mode 100644 pkg/util/fileutils2/seeker_test.go diff --git a/docs/bucket/uploadobject.yaml b/docs/bucket/uploadobject.yaml index 17c58f4a9e..a078a49729 100644 --- a/docs/bucket/uploadobject.yaml +++ b/docs/bucket/uploadobject.yaml @@ -4,7 +4,9 @@ post: - $ref: "../parameters/bucket.yaml#/bucket_name" - $ref: "../parameters/bucket.yaml#/x-bucket-object-key" - $ref: "../parameters/bucket.yaml#/x-bucket-content-type" + - $ref: "../parameters/bucket.yaml#/x-bucket-content-length" - $ref: "../parameters/bucket.yaml#/x-bucket-storage-class" + - $ref: "../parameters/bucket.yaml#/x-bucket-acl" responses: 200: description: 上传成功 diff --git a/docs/parameters/bucket.yaml b/docs/parameters/bucket.yaml index aee480aa02..e7654fe77e 100644 --- a/docs/parameters/bucket.yaml +++ b/docs/parameters/bucket.yaml @@ -18,19 +18,31 @@ recursive: description: 是否展开对象列表,false则只显示当前目录层级下的对象,true则显示匹配前缀的所有对象 x-bucket-object-key: - name: x-bucket-object-key + name: X-Yunion-Bucket-Upload-Key type: string in: header description: 对象的key x-bucket-content-type: - name: x-bucket-content-type + name: Content-Type type: string in: header description: 对象的content-type +x-bucket-content-length: + name: Content-Length + type: string + in: header + description: 对象的content-length + x-bucket-storage-class: - name: x-bucket-storage-class + name: X-Yunion-Bucket-Upload-Storageclass type: string in: header description: 对象的storage_class + +x-bucket-acl: + name: X-Yunion-Bucket-Upload-Acl + type: string + in: header + description: 对象的acl,可能值为:private, public-read, public-read-write diff --git a/docs/schemas/bucket.yaml b/docs/schemas/bucket.yaml index e4e718fbc6..79ac6a7ee5 100644 --- a/docs/schemas/bucket.yaml +++ b/docs/schemas/bucket.yaml @@ -187,7 +187,7 @@ BucketGetACLResponse: properties: acl: type: string - description: bucket或者对象的ACL字串,可能为private, public-read, public-read-write, default(仅object支持) + description: bucket或者对象的ACL字串,可能为private, public-read, public-read-write, authenticated-read BucketMakedirInput: type: object @@ -202,7 +202,7 @@ BucketSetACLInput: acl: type: string required: true - description: bucket或者对象的ACL字串,可能为private, public-read, public-read-write, default(仅object支持) + description: bucket或者对象的ACL字串,可能为private, public-read, public-read-write key: type: string description: 如果设置对象的ACL,则此字段指定对象的key diff --git a/pkg/apis/compute/bucket.go b/pkg/apis/compute/bucket.go index 683ab7e820..4b1fd5ec81 100644 --- a/pkg/apis/compute/bucket.go +++ b/pkg/apis/compute/bucket.go @@ -27,5 +27,6 @@ const ( BUCKET_STATUS_DELETE_FAIL = "delete_fail" BUCKET_UPLOAD_OBJECT_KEY_HEADER = "X-Yunion-Bucket-Upload-Key" + BUCKET_UPLOAD_OBJECT_ACL_HEADER = "X-Yunion-Bucket-Upload-Acl" BUCKET_UPLOAD_OBJECT_STORAGECLASS_HEADER = "X-Yunion-Bucket-Upload-Storageclass" ) diff --git a/pkg/cloudprovider/objectstore.go b/pkg/cloudprovider/objectstore.go index 95f75a5547..221b44b08b 100644 --- a/pkg/cloudprovider/objectstore.go +++ b/pkg/cloudprovider/objectstore.go @@ -17,10 +17,10 @@ package cloudprovider import ( "context" "io" + "strings" "time" - "strings" - + "yunion.io/x/log" "yunion.io/x/pkg/errors" "yunion.io/x/s3cli" ) @@ -28,7 +28,10 @@ import ( type TBucketACLType string const ( - ACLDefault = TBucketACLType("default") + // 100 MB + MAX_PUT_OBJECT_SIZEBYTES = int64(1000 * 1000 * 100) + + // ACLDefault = TBucketACLType("default") ACLPrivate = TBucketACLType(s3cli.CANNED_ACL_PRIVATE) ACLAuthRead = TBucketACLType(s3cli.CANNED_ACL_AUTH_READ) @@ -75,6 +78,9 @@ type SListObjectResult struct { type ICloudBucket interface { IVirtualResource + MaxPartCount() int + MaxPartSizeBytes() int64 + //GetGlobalId() string //GetName() string GetAcl() TBucketACLType @@ -89,10 +95,15 @@ type ICloudBucket interface { ListObjects(prefix string, marker string, delimiter string, maxCount int) (SListObjectResult, error) GetIObjects(prefix string, isRecursive bool) ([]ICloudObject, error) - PutObject(ctx context.Context, key string, input io.Reader, contType string, storageClass string) error + DeleteObject(ctx context.Context, keys string) error GetTempUrl(method string, key string, expire time.Duration) (string, error) - // ObjectExist(key string) (bool, error) + + PutObject(ctx context.Context, key string, input io.Reader, sizeBytes int64, contType string, cannedAcl TBucketACLType, storageClassStr string) error + NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl TBucketACLType, storageClassStr string) (string, error) + UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) + CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error + AbortMultipartUpload(ctx context.Context, key string, uploadId string) error } type ICloudObject interface { @@ -247,9 +258,72 @@ func Makedir(ctx context.Context, bucket ICloudBucket, key string) error { } } path := strings.Join(segs, "/") + "/" - err := bucket.PutObject(ctx, path, strings.NewReader(""), "", "") + err := bucket.PutObject(ctx, path, strings.NewReader(""), 0, "", ACLPrivate, "") if err != nil { return errors.Wrap(err, "PutObject") } return nil } + +func UploadObject(ctx context.Context, bucket ICloudBucket, key string, blocksz int64, input io.Reader, sizeBytes int64, contType string, cannedAcl TBucketACLType, storageClass string, debug bool) error { + if blocksz <= 0 { + blocksz = MAX_PUT_OBJECT_SIZEBYTES + } + if sizeBytes < blocksz { + if debug { + log.Debugf("too small, put object in one shot") + } + return bucket.PutObject(ctx, key, input, sizeBytes, contType, cannedAcl, storageClass) + } + partSize := blocksz + partCount := sizeBytes / partSize + if partCount*partSize < sizeBytes { + partCount += 1 + } + if partCount > int64(bucket.MaxPartCount()) { + partCount = int64(bucket.MaxPartCount()) + partSize = sizeBytes / partCount + if partSize*partCount < sizeBytes { + partSize += 1 + } + if partSize > bucket.MaxPartSizeBytes() { + return errors.Error("too larget object") + } + } + if debug { + log.Debugf("multipart upload part count %d part size %d", partCount, partSize) + } + uploadId, err := bucket.NewMultipartUpload(ctx, key, contType, cannedAcl, storageClass) + if err != nil { + return errors.Wrap(err, "bucket.NewMultipartUpload") + } + etags := make([]string, partCount) + // offset := int64(0) + for i := 0; i < int(partCount); i += 1 { + if i == int(partCount)-1 { + partSize = sizeBytes - partSize*(partCount-1) + } + if debug { + log.Debugf("UploadPart %d %d", i+1, partSize) + } + etag, err := bucket.UploadPart(ctx, key, uploadId, i+1, io.LimitReader(input, partSize), partSize) + if err != nil { + err2 := bucket.AbortMultipartUpload(ctx, key, uploadId) + if err2 != nil { + log.Errorf("bucket.AbortMultipartUpload error %s", err2) + } + return errors.Wrap(err, "bucket.UploadPart") + } + // offset += partSize + etags[i] = etag + } + err = bucket.CompleteMultipartUpload(ctx, key, uploadId, etags) + if err != nil { + err2 := bucket.AbortMultipartUpload(ctx, key, uploadId) + if err2 != nil { + log.Errorf("bucket.AbortMultipartUpload error %s", err2) + } + return errors.Wrap(err, "CompleteMultipartUpload") + } + return nil +} diff --git a/pkg/compute/models/buckets.go b/pkg/compute/models/buckets.go index 085c09ae56..b24cf36762 100644 --- a/pkg/compute/models/buckets.go +++ b/pkg/compute/models/buckets.go @@ -29,6 +29,7 @@ import ( "yunion.io/x/pkg/util/compare" "yunion.io/x/sqlchemy" + "strconv" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -709,8 +710,27 @@ func (bucket *SBucket) PerformUpload( } contType := appParams.Request.Header.Get("Content-Type") + sizeStr := appParams.Request.Header.Get("Content-Length") + if len(sizeStr) == 0 { + return nil, httperrors.NewInputParameterError("missing Content-Length") + } + sizeBytes, err := strconv.ParseInt(sizeStr, 10, 64) + if err != nil { + return nil, httperrors.NewInputParameterError("Illegal Content-Length %s", sizeStr) + } + if sizeBytes <= 0 { + return nil, httperrors.NewInputParameterError("Content-Length not positive %d", sizeBytes) + } storageClass := appParams.Request.Header.Get(api.BUCKET_UPLOAD_OBJECT_STORAGECLASS_HEADER) - err = iBucket.PutObject(ctx, key, appParams.Request.Body, contType, storageClass) + aclStr := cloudprovider.TBucketACLType(appParams.Request.Header.Get(api.BUCKET_UPLOAD_OBJECT_ACL_HEADER)) + switch aclStr { + case cloudprovider.ACLPrivate, cloudprovider.ACLAuthRead, cloudprovider.ACLPublicRead, cloudprovider.ACLPublicReadWrite: + // do nothing + default: + return nil, httperrors.NewInputParameterError("invalid acl: %s", aclStr) + } + + err = cloudprovider.UploadObject(ctx, iBucket, key, 0, appParams.Request.Body, sizeBytes, contType, aclStr, storageClass, false) if err != nil { return nil, httperrors.NewInternalServerError("put object error %s", err) } @@ -739,12 +759,8 @@ func (bucket *SBucket) PerformAcl( switch cloudprovider.TBucketACLType(aclStr) { case cloudprovider.ACLPrivate, cloudprovider.ACLAuthRead, cloudprovider.ACLPublicRead, cloudprovider.ACLPublicReadWrite: // do nothing - case cloudprovider.ACLDefault: - if len(objKey) == 0 { - return nil, httperrors.NewInputParameterError("invalud acl") - } default: - return nil, httperrors.NewInputParameterError("invalud acl") + return nil, httperrors.NewInputParameterError("invalid acl: %s", aclStr) } iBucket, err := bucket.GetIBucket() diff --git a/pkg/image/models/images.go b/pkg/image/models/images.go index 0e95994e9f..8738e2395a 100644 --- a/pkg/image/models/images.go +++ b/pkg/image/models/images.go @@ -23,11 +23,13 @@ import ( "os" "os/exec" "path/filepath" + "strconv" "strings" "time" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/timeutils" "yunion.io/x/pkg/utils" @@ -211,9 +213,16 @@ func (self *SImage) CustomizedGetDetailsBody(ctx context.Context, userCred mccli appParams := appsrv.AppContextGetParams(ctx) + fstat, err := os.Stat(filePath) + if err != nil { + return nil, errors.Wrap(err, "os.Stat") + } + + appParams.Response.Header().Set("Content-Length", strconv.FormatInt(fstat.Size(), 10)) + fp, err := os.Open(filePath) if err != nil { - return nil, err + return nil, errors.Wrap(err, "os.Open") } defer fp.Close() diff --git a/pkg/mcclient/modules/mod_images.go b/pkg/mcclient/modules/mod_images.go index b4a58789c6..a72df037de 100644 --- a/pkg/mcclient/modules/mod_images.go +++ b/pkg/mcclient/modules/mod_images.go @@ -25,6 +25,7 @@ import ( "yunion.io/x/log" "yunion.io/x/pkg/utils" + "strconv" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/httputils" @@ -538,7 +539,7 @@ func (this *ImageManager) _update(s *mcclient.ClientSession, id string, params j return json.Get("image") } -func (this *ImageManager) Download(s *mcclient.ClientSession, id string, format string, torrent bool) (jsonutils.JSONObject, io.Reader, error) { +func (this *ImageManager) Download(s *mcclient.ClientSession, id string, format string, torrent bool) (jsonutils.JSONObject, io.Reader, int64, error) { query := jsonutils.NewDict() if len(format) > 0 { query.Add(jsonutils.NewString(format), "format") @@ -553,10 +554,15 @@ func (this *ImageManager) Download(s *mcclient.ClientSession, id string, format } resp, err := this.rawRequest(s, "GET", path, nil, nil) if err == nil && resp.StatusCode >= 200 && resp.StatusCode < 300 { - return FetchImageMeta(resp.Header), resp.Body, nil + sizeBytes, err := strconv.ParseInt(resp.Header.Get("Content-Length"), 10, 64) + if err != nil { + log.Errorf("Download image unknown size") + sizeBytes = -1 + } + return FetchImageMeta(resp.Header), resp.Body, sizeBytes, nil } else { _, _, err = s.ParseJSONResponse(resp, err) - return nil, nil, err + return nil, nil, -1, err } } diff --git a/pkg/multicloud/aliyun/bucket.go b/pkg/multicloud/aliyun/bucket.go index 0a38e48363..a6577c3b79 100644 --- a/pkg/multicloud/aliyun/bucket.go +++ b/pkg/multicloud/aliyun/bucket.go @@ -22,9 +22,9 @@ import ( "github.com/aliyun/aliyun-oss-go-sdk/oss" + "yunion.io/x/log" "yunion.io/x/pkg/errors" - "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/multicloud" ) @@ -183,7 +183,7 @@ func (b *SBucket) GetIObjects(prefix string, isRecursive bool) ([]cloudprovider. return cloudprovider.GetIObjects(b, prefix, isRecursive) } -func (b *SBucket) PutObject(ctx context.Context, key string, input io.Reader, contType string, storageClassStr string) error { +func (b *SBucket) PutObject(ctx context.Context, key string, input io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { osscli, err := b.region.GetOssClient() if err != nil { return errors.Wrap(err, "GetOssClient") @@ -193,9 +193,19 @@ func (b *SBucket) PutObject(ctx context.Context, key string, input io.Reader, co return errors.Wrap(err, "Bucket") } opts := make([]oss.Option, 0) + if sizeBytes > 0 { + opts = append(opts, oss.ContentLength(sizeBytes)) + } if len(contType) > 0 { opts = append(opts, oss.ContentType(contType)) } + if len(cannedAcl) > 0 { + acl, err := str2Acl(string(cannedAcl)) + if err != nil { + return errors.Wrap(err, "") + } + opts = append(opts, oss.ObjectACL(acl)) + } if len(storageClassStr) > 0 { storageClass, err := str2StorageClass(storageClassStr) if err != nil { @@ -206,6 +216,116 @@ func (b *SBucket) PutObject(ctx context.Context, key string, input io.Reader, co return bucket.PutObject(key, input, opts...) } +func (b *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + osscli, err := b.region.GetOssClient() + if err != nil { + return "", errors.Wrap(err, "GetOssClient") + } + bucket, err := osscli.Bucket(b.Name) + if err != nil { + return "", errors.Wrap(err, "Bucket") + } + opts := make([]oss.Option, 0) + if len(contType) > 0 { + opts = append(opts, oss.ContentType(contType)) + } + if len(cannedAcl) > 0 { + acl, err := str2Acl(string(cannedAcl)) + if err != nil { + return "", errors.Wrap(err, "") + } + opts = append(opts, oss.ObjectACL(acl)) + } + if len(storageClassStr) > 0 { + storageClass, err := str2StorageClass(storageClassStr) + if err != nil { + return "", errors.Wrap(err, "str2StorageClass") + } + opts = append(opts, oss.ObjectStorageClass(storageClass)) + } + result, err := bucket.InitiateMultipartUpload(key, opts...) + if err != nil { + return "", errors.Wrap(err, "bucket.InitiateMultipartUpload") + } + return result.UploadID, nil +} + +func (b *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) { + osscli, err := b.region.GetOssClient() + if err != nil { + return "", errors.Wrap(err, "GetOssClient") + } + bucket, err := osscli.Bucket(b.Name) + if err != nil { + return "", errors.Wrap(err, "Bucket") + } + imur := oss.InitiateMultipartUploadResult{ + Bucket: b.Name, + Key: key, + UploadID: uploadId, + } + part, err := bucket.UploadPart(imur, input, partSize, partIndex) + if err != nil { + return "", errors.Wrap(err, "bucket.UploadPart") + } + if b.region.client.Debug { + log.Debugf("upload part key:%s uploadId:%s partIndex:%d etag:%s", key, uploadId, partIndex, part.ETag) + } + return part.ETag, nil +} + +func (b *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + osscli, err := b.region.GetOssClient() + if err != nil { + return errors.Wrap(err, "GetOssClient") + } + bucket, err := osscli.Bucket(b.Name) + if err != nil { + return errors.Wrap(err, "Bucket") + } + imur := oss.InitiateMultipartUploadResult{ + Bucket: b.Name, + Key: key, + UploadID: uploadId, + } + parts := make([]oss.UploadPart, len(partEtags)) + for i := range partEtags { + parts[i] = oss.UploadPart{ + PartNumber: i + 1, + ETag: partEtags[i], + } + } + result, err := bucket.CompleteMultipartUpload(imur, parts) + if err != nil { + return errors.Wrap(err, "bucket.CompleteMultipartUpload") + } + if b.region.client.Debug { + log.Debugf("CompleteMultipartUpload bucket:%s key:%s etag:%s location:%s", result.Bucket, result.Key, result.ETag, result.Location) + } + return nil +} + +func (b *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + osscli, err := b.region.GetOssClient() + if err != nil { + return errors.Wrap(err, "GetOssClient") + } + bucket, err := osscli.Bucket(b.Name) + if err != nil { + return errors.Wrap(err, "Bucket") + } + imur := oss.InitiateMultipartUploadResult{ + Bucket: b.Name, + Key: key, + UploadID: uploadId, + } + err = bucket.AbortMultipartUpload(imur) + if err != nil { + return errors.Wrap(err, "AbortMultipartUpload") + } + return nil +} + func (b *SBucket) DeleteObject(ctx context.Context, key string) error { osscli, err := b.region.GetOssClient() if err != nil { diff --git a/pkg/multicloud/aliyun/objects.go b/pkg/multicloud/aliyun/objects.go index 330759bd6b..2488dbac1c 100644 --- a/pkg/multicloud/aliyun/objects.go +++ b/pkg/multicloud/aliyun/objects.go @@ -18,6 +18,7 @@ import ( "github.com/pkg/errors" "yunion.io/x/log" + "github.com/aliyun/aliyun-oss-go-sdk/oss" "yunion.io/x/onecloud/pkg/cloudprovider" ) @@ -32,7 +33,7 @@ func (o *SObject) GetIBucket() cloudprovider.ICloudBucket { } func (o *SObject) GetAcl() cloudprovider.TBucketACLType { - acl := cloudprovider.ACLDefault + acl := cloudprovider.ACLPrivate osscli, err := o.bucket.region.GetOssClient() if err != nil { log.Errorf("o.bucket.region.GetOssClient error %s", err) @@ -48,6 +49,9 @@ func (o *SObject) GetAcl() cloudprovider.TBucketACLType { log.Errorf("bucket.GetObjectACL error %s", err) return acl } + if result.ACL == string(oss.ACLDefault) { + return o.bucket.GetAcl() + } acl = cloudprovider.TBucketACLType(result.ACL) return acl } diff --git a/pkg/multicloud/aliyun/storagecache.go b/pkg/multicloud/aliyun/storagecache.go index 11dd017378..cfba5732ab 100644 --- a/pkg/multicloud/aliyun/storagecache.go +++ b/pkg/multicloud/aliyun/storagecache.go @@ -148,25 +148,21 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To // first upload image to oss s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.QCOW2), false) + meta, reader, sizeByte, err := modules.Images.Download(s, image.ImageId, string(qemuimg.QCOW2), false) if err != nil { return "", err } log.Infof("meta data %s", meta) - oss, err := self.region.GetOssClient() - if err != nil { - log.Errorf("GetOssClient err %s", err) - return "", err - } + bucketName := strings.ToLower(fmt.Sprintf("imgcache-%s-%s", self.region.GetId(), image.ImageId)) - exist, err := oss.IsBucketExist(bucketName) + exist, err := self.region.IBucketExist(bucketName) if err != nil { log.Errorf("IsBucketExist err %s", err) return "", err } if !exist { log.Debugf("Bucket %s not exists, to create ...", bucketName) - err = oss.CreateBucket(bucketName) + err = self.region.CreateIBucket(bucketName, "", "") if err != nil { log.Errorf("Create bucket error %s", err) return "", err @@ -175,21 +171,22 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To log.Debugf("Bucket %s exists", bucketName) } - defer oss.DeleteBucket(bucketName) // remove bucket + defer self.region.DeleteIBucket(bucketName) // remove bucket - bucket, err := oss.Bucket(bucketName) + bucket, err := self.region.GetIBucketByName(bucketName) if err != nil { log.Errorf("Bucket error %s %s", bucketName, err) return "", err } log.Debugf("To upload image to bucket %s ...", bucketName) - err = bucket.PutObject(image.ImageId, reader) + err = cloudprovider.UploadObject(context.Background(), bucket, image.ImageId, 0, reader, sizeByte, "", "", "", false) + // err = bucket.PutObject(image.ImageId, reader) if err != nil { log.Errorf("PutObject error %s %s", image.ImageId, err) return "", err } - defer bucket.DeleteObject(image.ImageId) // remove object + defer bucket.DeleteObject(context.Background(), image.ImageId) // remove object imageBaseName := image.ImageId if imageBaseName[0] >= '0' && imageBaseName[0] <= '9' { diff --git a/pkg/multicloud/aws/bucket.go b/pkg/multicloud/aws/bucket.go index 4c93bcb26d..a3486bb63f 100644 --- a/pkg/multicloud/aws/bucket.go +++ b/pkg/multicloud/aws/bucket.go @@ -20,10 +20,8 @@ import ( "io" "time" - "github.com/aws/aws-sdk-go/aws" "github.com/aws/aws-sdk-go/aws/request" "github.com/aws/aws-sdk-go/service/s3" - "github.com/aws/aws-sdk-go/service/s3/s3manager" "yunion.io/x/log" "yunion.io/x/pkg/errors" @@ -31,6 +29,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/multicloud" + "yunion.io/x/onecloud/pkg/util/fileutils2" ) type SBucket struct { @@ -210,30 +209,126 @@ func (b *SBucket) GetIObjects(prefix string, isRecursive bool) ([]cloudprovider. return cloudprovider.GetIObjects(b, prefix, isRecursive) } -func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, contType string, storageClassStr string) error { - sess, err := b.region.getAwsSession() +func (b *SBucket) PutObject(ctx context.Context, key string, body io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { + if sizeBytes <= 0 { + return errors.Error("content length expected") + } + s3cli, err := b.region.GetS3Client() if err != nil { - return errors.Wrap(err, "session.NewSession") + return errors.Wrap(err, "GetS3Client") } - - svc := s3manager.NewUploader(sess) - input := &s3manager.UploadInput{ - Bucket: aws.String(b.Name), - Key: aws.String(key), - Body: reader, + input := &s3.PutObjectInput{} + input.SetBucket(b.Name) + input.SetKey(key) + seeker, err := fileutils2.NewReadSeeker(body, sizeBytes) + if err != nil { + return errors.Wrap(err, "newFakeSeeker") } + defer seeker.Close() + input.SetBody(seeker) + input.SetContentLength(sizeBytes) if len(contType) > 0 { - input.ContentType = aws.String(contType) + input.SetContentType(contType) + } + if len(cannedAcl) > 0 { + input.SetACL(string(cannedAcl)) } if len(storageClassStr) > 0 { - input.StorageClass = aws.String(storageClassStr) + input.SetStorageClass(storageClassStr) } - - _, err = svc.Upload(input) + _, err = s3cli.PutObjectWithContext(ctx, input) if err != nil { - return errors.Wrap(err, "svc.Upload") + return errors.Wrap(err, "PutObjectWithContext") } + return nil +} +func (b *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + s3cli, err := b.region.GetS3Client() + if err != nil { + return "", errors.Wrap(err, "GetS3Client") + } + input := &s3.CreateMultipartUploadInput{} + input.SetBucket(b.Name) + input.SetKey(key) + if len(contType) > 0 { + input.SetContentType(contType) + } + if len(cannedAcl) > 0 { + input.SetACL(string(cannedAcl)) + } + if len(storageClassStr) > 0 { + input.SetStorageClass(storageClassStr) + } + output, err := s3cli.CreateMultipartUploadWithContext(ctx, input) + if err != nil { + return "", errors.Wrap(err, "CreateMultipartUpload") + } + return *output.UploadId, nil +} + +func (b *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, part io.Reader, partSize int64) (string, error) { + s3cli, err := b.region.GetS3Client() + if err != nil { + return "", errors.Wrap(err, "GetS3Client") + } + input := &s3.UploadPartInput{} + input.SetBucket(b.Name) + input.SetKey(key) + input.SetUploadId(uploadId) + input.SetPartNumber(int64(partIndex)) + seeker, err := fileutils2.NewReadSeeker(part, partSize) + if err != nil { + return "", errors.Wrap(err, "newFakeSeeker") + } + defer seeker.Close() + input.SetBody(seeker) + input.SetContentLength(partSize) + output, err := s3cli.UploadPartWithContext(ctx, input) + if err != nil { + return "", errors.Wrap(err, "UploadPartWithContext") + } + return *output.ETag, nil +} + +func (b *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + s3cli, err := b.region.GetS3Client() + if err != nil { + return errors.Wrap(err, "GetS3Client") + } + input := &s3.CompleteMultipartUploadInput{} + input.SetBucket(b.Name) + input.SetKey(key) + input.SetUploadId(uploadId) + uploads := &s3.CompletedMultipartUpload{} + parts := make([]*s3.CompletedPart, len(partEtags)) + for i := range partEtags { + parts[i] = &s3.CompletedPart{} + parts[i].SetPartNumber(int64(i + 1)) + parts[i].SetETag(partEtags[i]) + } + uploads.SetParts(parts) + input.SetMultipartUpload(uploads) + _, err = s3cli.CompleteMultipartUploadWithContext(ctx, input) + if err != nil { + return errors.Wrap(err, "CompleteMultipartUploadWithContext") + } + return nil +} + +func (b *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + s3cli, err := b.region.GetS3Client() + if err != nil { + return errors.Wrap(err, "GetS3Client") + } + input := &s3.AbortMultipartUploadInput{} + input.SetBucket(b.Name) + input.SetKey(key) + input.SetUploadId(uploadId) + _, err = s3cli.AbortMultipartUploadWithContext(ctx, input) + if err != nil { + return errors.Wrap(err, "AbortMultipartUploadWithContext") + } return nil } diff --git a/pkg/multicloud/aws/s3object.go b/pkg/multicloud/aws/s3object.go index 7c0c6876d0..29728037dc 100644 --- a/pkg/multicloud/aws/s3object.go +++ b/pkg/multicloud/aws/s3object.go @@ -34,7 +34,7 @@ func (o *SObject) GetIBucket() cloudprovider.ICloudBucket { } func (o *SObject) GetAcl() cloudprovider.TBucketACLType { - acl := cloudprovider.ACLDefault + acl := cloudprovider.ACLPrivate s3cli, err := o.bucket.region.GetS3Client() if err != nil { log.Errorf("o.bucket.region.GetS3Client error %s", err) diff --git a/pkg/multicloud/aws/storagecache.go b/pkg/multicloud/aws/storagecache.go index 256eaaa47d..5b61cdc9fe 100644 --- a/pkg/multicloud/aws/storagecache.go +++ b/pkg/multicloud/aws/storagecache.go @@ -23,10 +23,10 @@ import ( "github.com/aws/aws-sdk-go/service/ec2" "github.com/aws/aws-sdk-go/service/iam" "github.com/aws/aws-sdk-go/service/s3" - "github.com/aws/aws-sdk-go/service/s3/s3manager" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -153,65 +153,46 @@ func (self *SStoragecache) fetchImages() error { } func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.TokenCredential, image *cloudprovider.SImageCreateOption, isForce bool) (string, error) { + err := self.region.initVmimport() + if err != nil { + return "", errors.Wrap(err, "initVmimport") + } + bucketName := GetBucketName(self.region.GetId(), image.ImageId) - err := self.region.initVmimport(bucketName) + + exist, err := self.region.IBucketExist(bucketName) if err != nil { - return "", err + return "", errors.Wrap(err, "IBucketExist") } - // checking remote - s3client, err := self.region.GetS3Client() - if err != nil { - return "", err + if !exist { + err = self.region.CreateIBucket(bucketName, "", "") + if err != nil { + return "", errors.Wrap(err, "CreateIBucket") + } } - defer s3client.DeleteBucket(&s3.DeleteBucketInput{Bucket: &bucketName}) // remove bucket + defer self.region.DeleteIBucket(bucketName) - var diskFormat string s := auth.GetAdminSession(ctx, options.Options.Region, "") - _, err = s3client.GetObject(&s3.GetObjectInput{Bucket: &bucketName, Key: &image.ImageId}) + meta, reader, sizeBytes, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) if err != nil { - // first upload image to oss - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) - if err != nil { - return "", err - } - log.Debugf("Images meta data %s", meta) - - diskFormat, err = meta.GetString("disk_format") - if err != nil { - return "", err - } - - // uploader to aws s3 - input := &s3manager.UploadInput{ - Bucket: &bucketName, - Key: &image.ImageId, - Body: reader, - } - - s3Session, err := self.region.getAwsSession() - if err != nil { - return "", err - } - - uploader := s3manager.NewUploader(s3Session) - _, err = uploader.Upload(input) - if err != nil { - return "", err - } - defer s3client.DeleteObject(&s3.DeleteObjectInput{Bucket: &bucketName, Key: &image.ImageId}) // remove object - } else { - meta, _, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) - if err != nil { - return "", err - } - - diskFormat, err = meta.GetString("disk_format") - if err != nil { - return "", err - } + return "", errors.Wrap(err, "Images.Download") } + log.Debugf("Images meta data %s", meta) + + diskFormat, _ := meta.GetString("disk_format") + + bucket, err := self.region.GetIBucketByName(bucketName) + if err != nil { + return "", errors.Wrap(err, "GetIBucketByName") + } + err = cloudprovider.UploadObject(ctx, bucket, image.ImageId, 0, reader, sizeBytes, "", "", "", false) + if err != nil { + return "", errors.Wrap(err, "cloudprovider.UploadObject") + } + + defer bucket.DeleteObject(ctx, image.ImageId) imageBaseName := image.ImageId if imageBaseName[0] >= '0' && imageBaseName[0] <= '9' { @@ -477,26 +458,7 @@ func (self *SRegion) initVmimportRolePolicy() error { } } -func (self *SRegion) initVmimportBucket(bucketName string) error { - exists, err := self.IsBucketExist(bucketName) - if err != nil { - return err - } - - if exists { - return nil - } - - s3Client, err := self.GetS3Client() - if err != nil { - return err - } - - _, err = s3Client.CreateBucket(&s3.CreateBucketInput{Bucket: &bucketName}) - return err -} - -func (self *SRegion) initVmimport(bucketName string) error { +func (self *SRegion) initVmimport() error { if err := self.initVmimportRole(); err != nil { return err } @@ -505,10 +467,6 @@ func (self *SRegion) initVmimport(bucketName string) error { return err } - if err := self.initVmimportBucket(bucketName); err != nil { - return err - } - return nil } diff --git a/pkg/multicloud/azure/blobobject.go b/pkg/multicloud/azure/blobobject.go index e36807d6ba..f2ec965c13 100644 --- a/pkg/multicloud/azure/blobobject.go +++ b/pkg/multicloud/azure/blobobject.go @@ -29,7 +29,7 @@ func (o *SObject) GetIBucket() cloudprovider.ICloudBucket { } func (o *SObject) GetAcl() cloudprovider.TBucketACLType { - return cloudprovider.ACLDefault + return o.container.getAcl() } func (o *SObject) SetAcl(aclStr cloudprovider.TBucketACLType) error { diff --git a/pkg/multicloud/azure/storageaccount.go b/pkg/multicloud/azure/storageaccount.go index 61b4a5fec4..5de0a7dec1 100644 --- a/pkg/multicloud/azure/storageaccount.go +++ b/pkg/multicloud/azure/storageaccount.go @@ -32,6 +32,8 @@ import ( "yunion.io/x/pkg/errors" "yunion.io/x/pkg/utils" + "encoding/base64" + "strconv" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/multicloud" ) @@ -689,6 +691,14 @@ func (self *SStorageAccount) UploadStream(containerName string, key string, read return container.UploadStream(key, reader, contType) } +func (b *SStorageAccount) MaxPartSizeBytes() int64 { + return 100 * 1000 * 1000 +} + +func (b *SStorageAccount) MaxPartCount() int { + return 50000 +} + func (b *SStorageAccount) GetProjectId() string { return "" } @@ -927,7 +937,7 @@ func splitKey(key string) (string, string, error) { return containerName, key, nil } -func (b *SStorageAccount) PutObject(ctx context.Context, key string, reader io.Reader, contType string, storageClassStr string) error { +func (b *SStorageAccount) PutObject(ctx context.Context, key string, reader io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { containerName, blob, err := splitKey(key) if err != nil { return errors.Wrap(err, "splitKey") @@ -939,6 +949,132 @@ func (b *SStorageAccount) PutObject(ctx context.Context, key string, reader io.R return nil } +func (b *SStorageAccount) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + containerName, blob, err := splitKey(key) + if err != nil { + return "", errors.Wrap(err, "splitKey") + } + container, err := b.getOrCreateContainer(containerName, true) + if err != nil { + return "", errors.Wrap(err, "getOrCreateContainer") + } + containerRef, err := container.getContainerRef() + if err != nil { + return "", errors.Wrap(err, "getContainerRef") + } + blobRef := containerRef.GetBlobReference(blob) + + err = blobRef.CreateBlockBlob(&storage.PutBlobOptions{}) + if err != nil { + return "", errors.Wrap(err, "CreateBlockBlob") + } + + uploadId, err := blobRef.AcquireLease(-1, "", nil) + if err != nil { + return "", errors.Wrap(err, "blobRef.AcquireLease") + } + + return uploadId, nil +} + +func (b *SStorageAccount) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) { + containerName, blob, err := splitKey(key) + if err != nil { + return "", errors.Wrap(err, "splitKey") + } + container, err := b.getOrCreateContainer(containerName, true) + if err != nil { + return "", errors.Wrap(err, "getOrCreateContainer") + } + containerRef, err := container.getContainerRef() + if err != nil { + return "", errors.Wrap(err, "getContainerRef") + } + blobRef := containerRef.GetBlobReference(blob) + + opts := &storage.PutBlockOptions{} + opts.LeaseID = uploadId + + blockId := base64.URLEncoding.EncodeToString([]byte(strconv.FormatInt(int64(partIndex), 10))) + err = blobRef.PutBlockWithLength(blockId, uint64(partSize), input, opts) + if err != nil { + return "", errors.Wrap(err, "PutBlockWithLength") + } + + return blockId, nil +} + +func (b *SStorageAccount) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, blockIds []string) error { + containerName, blob, err := splitKey(key) + if err != nil { + return errors.Wrap(err, "splitKey") + } + container, err := b.getOrCreateContainer(containerName, true) + if err != nil { + return errors.Wrap(err, "getOrCreateContainer") + } + containerRef, err := container.getContainerRef() + if err != nil { + return errors.Wrap(err, "getContainerRef") + } + blobRef := containerRef.GetBlobReference(blob) + + blocks := make([]storage.Block, len(blockIds)) + for i := range blockIds { + blocks[i] = storage.Block{ + ID: blockIds[i], + Status: storage.BlockStatusLatest, + } + } + opts := &storage.PutBlockListOptions{} + opts.LeaseID = uploadId + err = blobRef.PutBlockList(blocks, opts) + if err != nil { + return errors.Wrap(err, "PutBlockList") + } + + err = blobRef.ReleaseLease(uploadId, nil) + if err != nil { + return errors.Wrap(err, "ReleaseLease") + } + + return nil +} + +func (b *SStorageAccount) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + containerName, blob, err := splitKey(key) + if err != nil { + return errors.Wrap(err, "splitKey") + } + container, err := b.getOrCreateContainer(containerName, false) + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil + } + return errors.Wrap(err, "getOrCreateContainer") + } + containerRef, err := container.getContainerRef() + if err != nil { + return errors.Wrap(err, "getContainerRef") + } + + blobRef := containerRef.GetBlobReference(blob) + err = blobRef.ReleaseLease(uploadId, nil) + if err != nil { + return errors.Wrap(err, "ReleaseLease") + } + + opts := &storage.DeleteBlobOptions{} + deleteSnapshots := true + opts.DeleteSnapshots = &deleteSnapshots + _, err = blobRef.DeleteIfExists(opts) + if err != nil { + return errors.Wrap(err, "DeleteIfExists") + } + + return nil +} + func (b *SStorageAccount) DeleteObject(ctx context.Context, key string) error { containerName, blob, err := splitKey(key) if err != nil { diff --git a/pkg/multicloud/azure/storagecache.go b/pkg/multicloud/azure/storagecache.go index b4d5fe5e24..c000a14e3b 100644 --- a/pkg/multicloud/azure/storagecache.go +++ b/pkg/multicloud/azure/storagecache.go @@ -160,7 +160,7 @@ func (self *SStoragecache) checkStorageAccount() (*SStorageAccount, error) { func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.TokenCredential, image *cloudprovider.SImageCreateOption, isForce bool, tmpPath string) (string, error) { s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VHD), false) + meta, reader, sizeBytes, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VHD), false) if err != nil { return "", err } @@ -214,9 +214,9 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To return "", err } - size, _ := meta.Int("size") + // size, _ := meta.Int("size") - img, err := self.region.CreateImageByBlob(image.ImageId, image.OsType, blobURI, int32(size>>30)) + img, err := self.region.CreateImageByBlob(image.ImageId, image.OsType, blobURI, int32(sizeBytes>>30)) if err != nil { return "", err } diff --git a/pkg/multicloud/bucket_base.go b/pkg/multicloud/bucket_base.go index 7c9b876c58..cf0185b375 100644 --- a/pkg/multicloud/bucket_base.go +++ b/pkg/multicloud/bucket_base.go @@ -4,6 +4,14 @@ import "yunion.io/x/jsonutils" type SBaseBucket struct{} +func (b *SBaseBucket) MaxPartCount() int { + return 10000 +} + +func (b *SBaseBucket) MaxPartSizeBytes() int64 { + return 5 * 1000 * 1000 * 1000 +} + func (b *SBaseBucket) GetId() string { return "" } diff --git a/pkg/multicloud/huawei/bucket.go b/pkg/multicloud/huawei/bucket.go index 0bc81d4809..7f47df8942 100644 --- a/pkg/multicloud/huawei/bucket.go +++ b/pkg/multicloud/huawei/bucket.go @@ -227,7 +227,7 @@ func (b *SBucket) ListObjects(prefix string, marker string, delimiter string, ma return result, nil } -func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, contType string, storageClassStr string) error { +func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { obscli, err := b.region.getOBSClient() if err != nil { return errors.Wrap(err, "GetOBSClient") @@ -236,12 +236,19 @@ func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, c input.Bucket = b.Name input.Key = key input.Body = reader + + if sizeBytes > 0 { + input.ContentLength = sizeBytes + } if len(storageClassStr) > 0 { input.StorageClass, err = str2StorageClass(storageClassStr) if err != nil { return err } } + if len(cannedAcl) > 0 { + input.ACL = obs.AclType(string(cannedAcl)) + } if len(contType) > 0 { input.ContentType = contType } @@ -252,6 +259,100 @@ func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, c return nil } +func (b *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + obscli, err := b.region.getOBSClient() + if err != nil { + return "", errors.Wrap(err, "GetOBSClient") + } + + input := &obs.InitiateMultipartUploadInput{} + input.Bucket = b.Name + input.Key = key + if len(contType) > 0 { + input.ContentType = contType + } + if len(cannedAcl) > 0 { + input.ACL = obs.AclType(string(cannedAcl)) + } + if len(storageClassStr) > 0 { + input.StorageClass, err = str2StorageClass(storageClassStr) + if err != nil { + return "", errors.Wrap(err, "str2StorageClass") + } + } + output, err := obscli.InitiateMultipartUpload(input) + if err != nil { + return "", errors.Wrap(err, "InitiateMultipartUpload") + } + + return output.UploadId, nil +} + +func (b *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, part io.Reader, partSize int64) (string, error) { + obscli, err := b.region.getOBSClient() + if err != nil { + return "", errors.Wrap(err, "GetOBSClient") + } + + input := &obs.UploadPartInput{} + input.Bucket = b.Name + input.Key = key + input.UploadId = uploadId + input.PartNumber = partIndex + input.PartSize = partSize + input.Body = part + output, err := obscli.UploadPart(input) + if err != nil { + return "", errors.Wrap(err, "UploadPart") + } + + return output.ETag, nil +} + +func (b *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + obscli, err := b.region.getOBSClient() + if err != nil { + return errors.Wrap(err, "GetOBSClient") + } + input := &obs.CompleteMultipartUploadInput{} + input.Bucket = b.Name + input.Key = key + input.UploadId = uploadId + parts := make([]obs.Part, len(partEtags)) + for i := range partEtags { + parts[i] = obs.Part{ + PartNumber: i + 1, + ETag: partEtags[i], + } + } + input.Parts = parts + _, err = obscli.CompleteMultipartUpload(input) + if err != nil { + return errors.Wrap(err, "CompleteMultipartUpload") + } + + return nil +} + +func (b *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + obscli, err := b.region.getOBSClient() + if err != nil { + return errors.Wrap(err, "GetOBSClient") + } + + input := &obs.AbortMultipartUploadInput{} + input.Bucket = b.Name + input.Key = key + input.UploadId = uploadId + + _, err = obscli.AbortMultipartUpload(input) + if err != nil { + return errors.Wrap(err, "AbortMultipartUpload") + } + + return nil +} + func (b *SBucket) DeleteObject(ctx context.Context, key string) error { obscli, err := b.region.getOBSClient() if err != nil { diff --git a/pkg/multicloud/huawei/object.go b/pkg/multicloud/huawei/object.go index bb5f705090..04c63aa901 100644 --- a/pkg/multicloud/huawei/object.go +++ b/pkg/multicloud/huawei/object.go @@ -33,7 +33,7 @@ func (o *SObject) GetIBucket() cloudprovider.ICloudBucket { } func (o *SObject) GetAcl() cloudprovider.TBucketACLType { - acl := cloudprovider.ACLDefault + acl := cloudprovider.ACLPrivate obscli, err := o.bucket.region.getOBSClient() if err != nil { log.Errorf("o.bucket.region.GetOssClient error %s", err) diff --git a/pkg/multicloud/huawei/storagecache.go b/pkg/multicloud/huawei/storagecache.go index eda57631aa..3867874dac 100644 --- a/pkg/multicloud/huawei/storagecache.go +++ b/pkg/multicloud/huawei/storagecache.go @@ -23,6 +23,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -30,7 +31,6 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/mcclient/modules" - "yunion.io/x/onecloud/pkg/multicloud/huawei/obs" "yunion.io/x/onecloud/pkg/util/qemuimg" ) @@ -162,34 +162,28 @@ func (self *SStoragecache) UploadImage(ctx context.Context, userCred mcclient.To func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.TokenCredential, image *cloudprovider.SImageCreateOption, isForce bool) (string, error) { bucketName := GetBucketName(self.region.GetId(), image.ImageId) - obsClient, err := self.region.getOBSClient() - if err != nil { - return "", err - } - // create bucket - input := &obs.CreateBucketInput{} - input.Bucket = bucketName - input.Location = self.region.GetId() - _, err = obsClient.CreateBucket(input) + exist, err := self.region.IBucketExist(bucketName) if err != nil { - return "", err + return "", errors.Wrap(err, "self.region.IBucketExist") } - defer obsClient.DeleteBucket(bucketName) + if !exist { + err = self.region.CreateIBucket(bucketName, "", "") + if err != nil { + return "", errors.Wrap(err, "CreateIBucket") + } + } + defer self.region.DeleteIBucket(bucketName) // upload to huawei cloud s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) + meta, reader, sizeByte, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) if err != nil { return "", err } log.Debugf("Images meta data %s", meta) - _image, err := modules.Images.Get(s, image.ImageId, nil) - if err != nil { - return "", err - } - minDiskMB, _ := _image.Int("min_disk") + minDiskMB, _ := meta.Int("min_disk") minDiskGB := int64(math.Ceil(float64(minDiskMB) / 1024)) // 在使用OBS桶的外部镜像文件制作镜像时生效且为必选字段。取值为40~1024GB。 if minDiskGB < 40 { @@ -198,21 +192,17 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To minDiskGB = 1024 } - // upload to huawei cloud - obj := &obs.PutObjectInput{} - obj.Bucket = bucketName - obj.Key = image.ImageId - obj.Body = reader - - _, err = obsClient.PutObject(obj) + bucket, err := self.region.GetIBucketByName(bucketName) if err != nil { - return "", err + return "", errors.Wrap(err, "GetIBucketByName") } - objDelete := &obs.DeleteObjectInput{} - objDelete.Bucket = bucketName - objDelete.Key = image.ImageId - defer obsClient.DeleteObject(objDelete) // remove object + err = cloudprovider.UploadObject(context.Background(), bucket, image.ImageId, 0, reader, sizeByte, "", "", "", false) + if err != nil { + return "", errors.Wrap(err, "cloudprovider.UploadObject") + } + + defer bucket.DeleteObject(context.Background(), image.ImageId) // check image name, avoid name conflict imageBaseName := image.ImageId diff --git a/pkg/multicloud/objectstore/buckets.go b/pkg/multicloud/objectstore/buckets.go index 7b664962d2..f6f2506683 100644 --- a/pkg/multicloud/objectstore/buckets.go +++ b/pkg/multicloud/objectstore/buckets.go @@ -26,9 +26,12 @@ import ( "yunion.io/x/s3cli" "yunion.io/x/onecloud/pkg/cloudprovider" + "yunion.io/x/onecloud/pkg/multicloud" ) type SBucket struct { + multicloud.SBaseBucket + client *SObjectStoreClient Name string @@ -150,16 +153,75 @@ func (bucket *SBucket) GetIObjects(prefix string, isRecursive bool) ([]cloudprov return ret, nil } -func (bucket *SBucket) PutObject(ctx context.Context, key string, input io.Reader, contType string, storageClass string) error { +func (bucket *SBucket) PutObject(ctx context.Context, key string, input io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { opts := s3cli.PutObjectOptions{} if len(contType) > 0 { opts.ContentType = contType } - if len(storageClass) > 0 { - opts.StorageClass = storageClass + if len(storageClassStr) > 0 { + opts.StorageClass = storageClassStr } - _, err := bucket.client.client.PutObjectWithContext(ctx, bucket.Name, key, input, -1, opts) - return err + opts.PartSize = uint64(cloudprovider.MAX_PUT_OBJECT_SIZEBYTES) + _, err := bucket.client.client.PutObjectDo(ctx, bucket.Name, key, input, "", "", sizeBytes, opts) + if err != nil { + return errors.Wrap(err, "PutObjectWithContext") + } + obj, err := cloudprovider.GetIObject(bucket, key) + if err != nil { + return errors.Wrap(err, "GetIObject") + } + err = obj.SetAcl(cannedAcl) + if err != nil { + return errors.Wrap(err, "obj.SetAcl") + } + return nil +} + +func (bucket *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + opts := s3cli.PutObjectOptions{} + if len(contType) > 0 { + opts.ContentType = contType + } + if len(storageClassStr) > 0 { + opts.StorageClass = storageClassStr + } + result, err := bucket.client.client.InitiateMultipartUpload(ctx, bucket.Name, key, opts) + if err != nil { + return "", errors.Wrap(err, "InitiateMultipartUpload") + } + return result.UploadID, nil +} + +func (bucket *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) { + part, err := bucket.client.client.UploadPart(ctx, bucket.Name, key, uploadId, input, partIndex, "", "", partSize, nil) + if err != nil { + return "", errors.Wrap(err, "UploadPart") + } + return part.ETag, nil +} + +func (bucket *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + complete := s3cli.CompleteMultipartUpload{} + complete.Parts = make([]s3cli.CompletePart, len(partEtags)) + for i := 0; i < len(partEtags); i += 1 { + complete.Parts[i] = s3cli.CompletePart{ + PartNumber: i + 1, + ETag: partEtags[i], + } + } + _, err := bucket.client.client.CompleteMultipartUpload(ctx, bucket.Name, key, uploadId, complete) + if err != nil { + return errors.Wrap(err, "CompleteMultipartUpload") + } + return nil +} + +func (bucket *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + err := bucket.client.client.AbortMultipartUpload(ctx, bucket.Name, key, uploadId) + if err != nil { + return errors.Wrap(err, "AbortMultipartUpload") + } + return nil } func (bucket *SBucket) DeleteObject(ctx context.Context, key string) error { diff --git a/pkg/multicloud/objectstore/shell.go b/pkg/multicloud/objectstore/shell.go index 01b29d62bb..560feb0d0e 100644 --- a/pkg/multicloud/objectstore/shell.go +++ b/pkg/multicloud/objectstore/shell.go @@ -21,6 +21,7 @@ import ( "os" "time" + "github.com/pkg/errors" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/util/printutils" "yunion.io/x/onecloud/pkg/util/shellutils" @@ -156,6 +157,9 @@ func S3Shell() { KEY string `help:"key of object"` Path string `help:"Path of file to upload"` + BlockSize int64 `help:"blocksz in MB" default:"100"` + + Acl string `help:"acl" choices:"private|public-read|public-read-write"` ContentType string `help:"content-type"` StorageClass string `help:"storage class"` } @@ -165,10 +169,16 @@ func S3Shell() { return err } var input io.ReadSeeker + var fSize int64 if len(args.Path) > 0 { + finfo, err := os.Stat(args.Path) + if err != nil { + return errors.Wrap(err, "os.Stat") + } + fSize = finfo.Size() file, err := os.Open(args.Path) if err != nil { - return err + return errors.Wrap(err, "os.Open") } defer file.Close() @@ -176,7 +186,7 @@ func S3Shell() { } else { input = os.Stdout } - err = bucket.PutObject(context.Background(), args.KEY, input, args.ContentType, args.StorageClass) + err = cloudprovider.UploadObject(context.Background(), bucket, args.KEY, args.BlockSize*1000*1000, input, fSize, args.ContentType, cloudprovider.TBucketACLType(args.Acl), args.StorageClass, true) if err != nil { return err } diff --git a/pkg/multicloud/openstack/storagecache.go b/pkg/multicloud/openstack/storagecache.go index 2527ecc13c..afdbe1010f 100644 --- a/pkg/multicloud/openstack/storagecache.go +++ b/pkg/multicloud/openstack/storagecache.go @@ -119,7 +119,7 @@ func (cache *SStoragecache) UploadImage(ctx context.Context, userCred mcclient.T func (cache *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.TokenCredential, image *cloudprovider.SImageCreateOption, isForce bool) (string, error) { s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) + meta, reader, _, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) if err != nil { return "", err } diff --git a/pkg/multicloud/qcloud/bucket.go b/pkg/multicloud/qcloud/bucket.go index 58bcaaa952..8369562934 100644 --- a/pkg/multicloud/qcloud/bucket.go +++ b/pkg/multicloud/qcloud/bucket.go @@ -216,15 +216,24 @@ func (b *SBucket) ListObjects(prefix string, marker string, delimiter string, ma return result, nil } -func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, contType string, storageClassStr string) error { +func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { coscli, err := b.region.GetCosClient(b) if err != nil { return errors.Wrap(err, "GetCosClient") } - opts := &cos.ObjectPutOptions{} + opts := &cos.ObjectPutOptions{ + ACLHeaderOptions: &cos.ACLHeaderOptions{}, + ObjectPutHeaderOptions: &cos.ObjectPutHeaderOptions{}, + } + if sizeBytes > 0 { + opts.ContentLength = int(sizeBytes) + } if len(contType) > 0 { opts.ContentType = contType } + if len(cannedAcl) > 0 { + opts.XCosACL = string(cannedAcl) + } if len(storageClassStr) > 0 { opts.XCosStorageClass = storageClassStr } @@ -235,6 +244,84 @@ func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, c return nil } +func (b *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + coscli, err := b.region.GetCosClient(b) + if err != nil { + return "", errors.Wrap(err, "GetCosClient") + } + opts := &cos.InitiateMultipartUploadOptions{ + ACLHeaderOptions: &cos.ACLHeaderOptions{}, + ObjectPutHeaderOptions: &cos.ObjectPutHeaderOptions{}, + } + if len(contType) > 0 { + opts.ContentType = contType + } + if len(cannedAcl) > 0 { + opts.XCosACL = string(cannedAcl) + } + if len(storageClassStr) > 0 { + opts.XCosStorageClass = storageClassStr + } + result, _, err := coscli.Object.InitiateMultipartUpload(ctx, key, opts) + if err != nil { + return "", errors.Wrap(err, "InitiateMultipartUpload") + } + + return result.UploadID, nil +} + +func (b *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) { + coscli, err := b.region.GetCosClient(b) + if err != nil { + return "", errors.Wrap(err, "GetCosClient") + } + opts := &cos.ObjectUploadPartOptions{} + opts.ContentLength = int(partSize) + resp, err := coscli.Object.UploadPart(ctx, key, uploadId, partIndex, input, opts) + if err != nil { + return "", errors.Wrap(err, "UploadPart") + } + + return resp.Header.Get("Etag"), nil +} + +func (b *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + coscli, err := b.region.GetCosClient(b) + if err != nil { + return errors.Wrap(err, "GetCosClient") + } + opts := &cos.CompleteMultipartUploadOptions{} + parts := make([]cos.Object, len(partEtags)) + for i := range partEtags { + parts[i] = cos.Object{ + PartNumber: i + 1, + ETag: partEtags[i], + } + } + opts.Parts = parts + _, _, err = coscli.Object.CompleteMultipartUpload(ctx, key, uploadId, opts) + + if err != nil { + return errors.Wrap(err, "CompleteMultipartUpload") + } + + return nil +} + +func (b *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { + coscli, err := b.region.GetCosClient(b) + if err != nil { + return errors.Wrap(err, "GetCosClient") + } + + _, err = coscli.Object.AbortMultipartUpload(ctx, key, uploadId) + if err != nil { + return errors.Wrap(err, "AbortMultipartUpload") + } + + return nil +} + func (b *SBucket) DeleteObject(ctx context.Context, key string) error { coscli, err := b.region.GetCosClient(b) if err != nil { diff --git a/pkg/multicloud/qcloud/object.go b/pkg/multicloud/qcloud/object.go index d9fa896d39..677694726e 100644 --- a/pkg/multicloud/qcloud/object.go +++ b/pkg/multicloud/qcloud/object.go @@ -15,11 +15,13 @@ package qcloud import ( + "context" + + "github.com/tencentyun/cos-go-sdk-v5" + "yunion.io/x/log" "yunion.io/x/pkg/errors" - "context" - "github.com/tencentyun/cos-go-sdk-v5" "yunion.io/x/onecloud/pkg/cloudprovider" ) @@ -34,7 +36,7 @@ func (o *SObject) GetIBucket() cloudprovider.ICloudBucket { } func (o *SObject) GetAcl() cloudprovider.TBucketACLType { - acl := cloudprovider.ACLDefault + acl := cloudprovider.ACLPrivate coscli, err := o.bucket.region.GetCosClient(o.bucket) if err != nil { log.Errorf("o.bucket.region.GetOssClient error %s", err) diff --git a/pkg/multicloud/qcloud/storagecache.go b/pkg/multicloud/qcloud/storagecache.go index 57de5605d3..bc0852329e 100644 --- a/pkg/multicloud/qcloud/storagecache.go +++ b/pkg/multicloud/qcloud/storagecache.go @@ -17,7 +17,6 @@ package qcloud import ( "context" "fmt" - "io" "strings" "time" @@ -163,7 +162,7 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To // first upload image to oss s := auth.GetAdminSession(ctx, options.Options.Region, "") - _, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) + _, reader, sizeBytes, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) if err != nil { return "", err } @@ -186,7 +185,8 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To if err != nil { return "", errors.Wrap(err, "GetIBucketByName") } - err = bucket.PutObject(context.Background(), image.ImageId, reader.(io.ReadSeeker), "", "") + err = cloudprovider.UploadObject(context.Background(), bucket, image.ImageId, 0, reader, sizeBytes, "", "", "", false) + // err = bucket.PutObject(context.Background(), image.ImageId, reader, sizeBytes, "", "", "") if err != nil { log.Errorf("UploadObject error %s %s", image.ImageId, err) return "", errors.Wrap(err, "bucket.PutObject") diff --git a/pkg/multicloud/ucloud/storagecache.go b/pkg/multicloud/ucloud/storagecache.go index 36ff41ed09..99b6ae70d8 100644 --- a/pkg/multicloud/ucloud/storagecache.go +++ b/pkg/multicloud/ucloud/storagecache.go @@ -162,7 +162,7 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To // upload to ucloud s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) + meta, reader, size, err := modules.Images.Download(s, image.ImageId, string(qemuimg.VMDK), false) if err != nil { return "", err } @@ -175,7 +175,7 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To } else if minDiskGB > 1024 { minDiskGB = 1024 } - size, _ := meta.Int("size") + // size, _ := meta.Int("size") md5, _ := meta.GetString("checksum") diskFormat, _ := meta.GetString("disk_format") // upload to ucloud diff --git a/pkg/multicloud/ucloud/ufile.go b/pkg/multicloud/ucloud/ufile.go index 4288702f05..9a1249805b 100644 --- a/pkg/multicloud/ucloud/ufile.go +++ b/pkg/multicloud/ucloud/ufile.go @@ -175,7 +175,7 @@ func (self *SFile) GetContentType() string { } func (self *SFile) GetAcl() cloudprovider.TBucketACLType { - return cloudprovider.ACLDefault + return self.bucket.GetAcl() } func (self *SFile) SetAcl(cloudprovider.TBucketACLType) error { @@ -343,7 +343,23 @@ func (b *SBucket) ListObjects(prefix string, marker string, delimiter string, ma return result, nil } -func (b *SBucket) PutObject(ctx context.Context, key string, reader io.Reader, contType string, storageClassStr string) error { +func (b *SBucket) PutObject(ctx context.Context, key string, input io.Reader, sizeBytes int64, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) error { + return cloudprovider.ErrNotSupported +} + +func (b *SBucket) NewMultipartUpload(ctx context.Context, key string, contType string, cannedAcl cloudprovider.TBucketACLType, storageClassStr string) (string, error) { + return "", cloudprovider.ErrNotSupported +} + +func (b *SBucket) UploadPart(ctx context.Context, key string, uploadId string, partIndex int, input io.Reader, partSize int64) (string, error) { + return "", cloudprovider.ErrNotSupported +} + +func (b *SBucket) CompleteMultipartUpload(ctx context.Context, key string, uploadId string, partEtags []string) error { + return cloudprovider.ErrNotSupported +} + +func (b *SBucket) AbortMultipartUpload(ctx context.Context, key string, uploadId string) error { return cloudprovider.ErrNotSupported } diff --git a/pkg/multicloud/zstack/storagecache.go b/pkg/multicloud/zstack/storagecache.go index bf0c656821..8c50ff361b 100644 --- a/pkg/multicloud/zstack/storagecache.go +++ b/pkg/multicloud/zstack/storagecache.go @@ -118,13 +118,13 @@ func (scache *SStoragecache) UploadImage(ctx context.Context, userCred mcclient. func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.TokenCredential, image *cloudprovider.SImageCreateOption, isForce bool) (string, error) { s := auth.GetAdminSession(ctx, options.Options.Region, "") - meta, reader, err := modules.Images.Download(s, image.ImageId, string(qemuimg.QCOW2), false) + meta, reader, size, err := modules.Images.Download(s, image.ImageId, string(qemuimg.QCOW2), false) if err != nil { return "", err } log.Infof("meta data %s", meta) - size, _ := meta.Int("size") + // size, _ := meta.Int("size") img, err := self.region.CreateImage(self.ZoneId, image.ImageName, string(qemuimg.QCOW2), image.OsType, "", reader, size) if err != nil { return "", err diff --git a/pkg/util/fileutils2/seeker.go b/pkg/util/fileutils2/seeker.go new file mode 100644 index 0000000000..5c3920884b --- /dev/null +++ b/pkg/util/fileutils2/seeker.go @@ -0,0 +1,81 @@ +package fileutils2 + +import ( + "io" + "io/ioutil" + "os" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" +) + +type SReadSeeker struct { + reader io.Reader + + offset int64 + + readerOffset int64 + readerSize int64 + + tmpFile *os.File +} + +func NewReadSeeker(reader io.Reader, size int64) (*SReadSeeker, error) { + tmpfile, err := ioutil.TempFile("", "fakeseeker") + if err != nil { + return nil, errors.Wrap(err, "TempFile") + } + log.Debugf("use tmpfile %s", tmpfile.Name()) + return &SReadSeeker{ + reader: reader, + readerOffset: 0, + readerSize: size, + + tmpFile: tmpfile, + }, nil +} + +func (s *SReadSeeker) Read(p []byte) (int, error) { + if s.offset == s.readerOffset && s.offset < s.readerSize { + n, err := s.reader.Read(p) + if n > 0 { + wn, werr := s.tmpFile.Write(p[:n]) + if werr != nil { + return n, werr + } + if wn < n { + return n, errors.Error("sFakeSeeker write less bytes") + } + s.offset += int64(n) + s.readerOffset += int64(n) + } + return n, err + } else { + n, err := s.tmpFile.ReadAt(p, s.offset) + if n > 0 { + s.offset += int64(n) + } + return n, err + } +} + +func (s *SReadSeeker) Seek(offset int64, whence int) (int64, error) { + switch whence { + case io.SeekStart: + // offset = offset + case io.SeekCurrent: + offset = s.offset + offset + case io.SeekEnd: + offset = s.readerSize + offset + } + if offset < 0 || offset > s.readerSize { + log.Debugf("offset out of range: %d", offset) + return -1, io.ErrUnexpectedEOF + } + s.offset = offset + return offset, nil +} + +func (s *SReadSeeker) Close() error { + return os.Remove(s.tmpFile.Name()) +} diff --git a/pkg/util/fileutils2/seeker_test.go b/pkg/util/fileutils2/seeker_test.go new file mode 100644 index 0000000000..88f50694c8 --- /dev/null +++ b/pkg/util/fileutils2/seeker_test.go @@ -0,0 +1,31 @@ +package fileutils2 + +import ( + "io" + "strings" + "testing" +) + +func TestNewReadSeeker(t *testing.T) { + testStr := "This is a test reader string" + seeker, err := NewReadSeeker(strings.NewReader(testStr), int64(len(testStr))) + if err != nil { + t.Fatalf("NewReadSeeker error %s", err) + } + defer seeker.Close() + buf1 := make([]byte, 1024) + n, err := seeker.Read(buf1) + if n != len(testStr) { + t.Fatalf("read buf1 error %s", err) + } + buf2 := make([]byte, 1024) + n, err = seeker.Read(buf2) + if n != 0 { + t.Fatalf("read buf2 should fail") + } + seeker.Seek(0, io.SeekStart) + n, err = seeker.Read(buf2) + if n != len(testStr) { + t.Fatalf("read buf2 error %s", err) + } +} diff --git a/pkg/util/httputils/httputils.go b/pkg/util/httputils/httputils.go index 9641a29f35..f7438c46d2 100644 --- a/pkg/util/httputils/httputils.go +++ b/pkg/util/httputils/httputils.go @@ -130,12 +130,14 @@ func GetAddrPort(urlStr string) (string, int, error) { func GetTransport(insecure bool, timeout time.Duration) *http.Transport { return &http.Transport{ DialContext: (&net.Dialer{ - Timeout: 5 * time.Second, + Timeout: timeout, + KeepAlive: timeout, }).DialContext, - IdleConnTimeout: 5 * time.Second, + IdleConnTimeout: 90 * time.Second, TLSHandshakeTimeout: 10 * time.Second, ExpectContinueTimeout: 1 * time.Second, TLSClientConfig: &tls.Config{InsecureSkipVerify: insecure}, + DisableCompression: true, } }