From 60ba83b40fd6df4ed3fcbf3feba8d06888a6950f Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 8 Feb 2020 12:01:52 +0800 Subject: [PATCH] fix: progressively update image size while uploading --- pkg/image/models/images.go | 22 ++++++++++++++++++++-- pkg/multicloud/aws/shell/s3.go | 2 +- pkg/multicloud/objectstore/shell.go | 2 +- pkg/util/s3auth/v4.go | 2 +- pkg/util/streamutils/streamutils.go | 5 ++++- 5 files changed, 27 insertions(+), 6 deletions(-) diff --git a/pkg/image/models/images.go b/pkg/image/models/images.go index e8523a8227..4e04944ff3 100644 --- a/pkg/image/models/images.go +++ b/pkg/image/models/images.go @@ -244,7 +244,7 @@ func (self *SImage) CustomizedGetDetailsBody(ctx context.Context, userCred mccli } defer fp.Close() - _, err = streamutils.StreamPipe(fp, appParams.Response, false) + _, err = streamutils.StreamPipe(fp, appParams.Response, false, nil) if err != nil { return nil, httperrors.NewGeneralError(err) } @@ -448,7 +448,25 @@ func (self *SImage) saveImageFromStream(localPath string, reader io.Reader, calC return nil, err } defer fp.Close() - return streamutils.StreamPipe(reader, fp, calChecksum) + lastSaveTime := time.Now() + return streamutils.StreamPipe(reader, fp, calChecksum, func(saved int64) { + now := time.Now() + if now.Sub(lastSaveTime) > 5*time.Second { + self.saveSize(saved) + lastSaveTime = now + } + }) +} + +func (self *SImage) saveSize(newSize int64) error { + _, err := db.Update(self, func() error { + self.Size = newSize + return nil + }) + if err != nil { + return errors.Wrap(err, "Update size") + } + return nil } //Image always do probe and customize after save from stream diff --git a/pkg/multicloud/aws/shell/s3.go b/pkg/multicloud/aws/shell/s3.go index a74d0f3988..4fc026b485 100644 --- a/pkg/multicloud/aws/shell/s3.go +++ b/pkg/multicloud/aws/shell/s3.go @@ -121,7 +121,7 @@ func init() { } else { fio = os.Stdout } - prop, err := streamutils.StreamPipe(output.Body, fio, true) + prop, err := streamutils.StreamPipe(output.Body, fio, true, nil) if err != nil { return err } diff --git a/pkg/multicloud/objectstore/shell.go b/pkg/multicloud/objectstore/shell.go index de3eae56a3..7fde103908 100644 --- a/pkg/multicloud/objectstore/shell.go +++ b/pkg/multicloud/objectstore/shell.go @@ -416,7 +416,7 @@ func S3Shell() { defer fp.Close() target = fp } - prop, err := streamutils.StreamPipe(output, target, false) + prop, err := streamutils.StreamPipe(output, target, false, nil) if err != nil { return err } diff --git a/pkg/util/s3auth/v4.go b/pkg/util/s3auth/v4.go index ba8dfa768d..535afa3dcf 100644 --- a/pkg/util/s3auth/v4.go +++ b/pkg/util/s3auth/v4.go @@ -306,7 +306,7 @@ func SignV4(req http.Request, accessKey, secretAccessKey, location string, body h := sha256.New() if body != nil { - streamutils.StreamPipe(body, h, false) + streamutils.StreamPipe(body, h, false, nil) } req.Header.Set("X-Amz-Content-Sha256", hex.EncodeToString(h.Sum(nil))) diff --git a/pkg/util/streamutils/streamutils.go b/pkg/util/streamutils/streamutils.go index b044f56e73..38d792b15d 100644 --- a/pkg/util/streamutils/streamutils.go +++ b/pkg/util/streamutils/streamutils.go @@ -26,7 +26,7 @@ type SStreamProperty struct { Size int64 } -func StreamPipe(reader io.Reader, writer io.Writer, CalChecksum bool) (*SStreamProperty, error) { +func StreamPipe(reader io.Reader, writer io.Writer, CalChecksum bool, callback func(saved int64)) (*SStreamProperty, error) { sp := SStreamProperty{} var md5sum hash.Hash @@ -50,6 +50,9 @@ func StreamPipe(reader io.Reader, writer io.Writer, CalChecksum bool) (*SStreamP } offset += m } + if callback != nil { + callback(sp.Size) + } } if err != nil { if err == io.EOF {