forked from TrueCloudLab/frostfs-s3-gw
[#346] api: Do not use io.Pipe
in CompleteMultipartUpload
Replace `layer.objectWritePayload` method with `initObjectPayloadReader` which returns `io.Reader` of the object payload. Copy payload data to the parameterized `io.Writer` in `layer.GetObject`. Remove `io.Pipe` from `CompleteMultipartUpload` implementation and build analogue of `io.MultiReader` for the part list. Signed-off-by: Leonard Lyubich <leonard@nspcc.ru>
This commit is contained in:
parent
eac4c4d849
commit
8fb3835250
3 changed files with 79 additions and 63 deletions
|
@ -567,7 +567,6 @@ func (n *layer) ListBuckets(ctx context.Context) ([]*data.BucketInfo, error) {
|
|||
func (n *layer) GetObject(ctx context.Context, p *GetObjectParams) error {
|
||||
var params getParams
|
||||
|
||||
params.w = p.Writer
|
||||
params.oid = p.ObjectInfo.ID
|
||||
params.cid = p.ObjectInfo.CID
|
||||
|
||||
|
@ -580,10 +579,22 @@ func (n *layer) GetObject(ctx context.Context, p *GetObjectParams) error {
|
|||
params.ln = p.Range.End - p.Range.Start + 1
|
||||
}
|
||||
|
||||
err := n.objectWritePayload(ctx, params)
|
||||
payload, err := n.initObjectPayloadReader(ctx, params)
|
||||
if err != nil {
|
||||
n.objCache.Delete(p.ObjectInfo.Address())
|
||||
return fmt.Errorf("couldn't get object, cid: %s : %w", p.ObjectInfo.CID, err)
|
||||
return fmt.Errorf("init object payload reader: %w", err)
|
||||
}
|
||||
|
||||
if params.ln == 0 {
|
||||
params.ln = 4096 // configure?
|
||||
}
|
||||
|
||||
// alloc buffer for copying
|
||||
buf := make([]byte, params.ln) // sync-pool it?
|
||||
|
||||
// copy full payload
|
||||
_, err = io.CopyBuffer(p.Writer, payload, buf)
|
||||
if err != nil {
|
||||
return fmt.Errorf("copy object payload: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
|
|
@ -2,6 +2,8 @@ package layer
|
|||
|
||||
import (
|
||||
"context"
|
||||
stderrors "errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"sort"
|
||||
"strconv"
|
||||
|
@ -169,6 +171,45 @@ func (n *layer) UploadPartCopy(ctx context.Context, p *UploadCopyParams) (*data.
|
|||
return n.CopyObject(ctx, c)
|
||||
}
|
||||
|
||||
// implements io.Reader of payloads of the object list stored in the NeoFS network.
|
||||
type multiObjectReader struct {
|
||||
ctx context.Context
|
||||
|
||||
layer *layer
|
||||
|
||||
prm getParams
|
||||
|
||||
curReader io.Reader
|
||||
|
||||
parts []*data.ObjectInfo
|
||||
}
|
||||
|
||||
func (x *multiObjectReader) Read(p []byte) (n int, err error) {
|
||||
if x.curReader != nil {
|
||||
n, err = x.curReader.Read(p)
|
||||
if !stderrors.Is(err, io.EOF) {
|
||||
return n, err
|
||||
}
|
||||
}
|
||||
|
||||
if len(x.parts) == 0 {
|
||||
return n, io.EOF
|
||||
}
|
||||
|
||||
x.prm.oid = x.parts[0].ID
|
||||
|
||||
x.curReader, err = x.layer.initObjectPayloadReader(x.ctx, x.prm)
|
||||
if err != nil {
|
||||
return n, fmt.Errorf("init payload reader for the next part: %w", err)
|
||||
}
|
||||
|
||||
x.parts = x.parts[1:]
|
||||
|
||||
next, err := x.Read(p[n:])
|
||||
|
||||
return n + next, err
|
||||
}
|
||||
|
||||
func (n *layer) CompleteMultipartUpload(ctx context.Context, p *CompleteMultipartParams) (*data.ObjectInfo, error) {
|
||||
var obj *data.ObjectInfo
|
||||
|
||||
|
@ -231,10 +272,6 @@ func (n *layer) CompleteMultipartUpload(ctx context.Context, p *CompleteMultipar
|
|||
initMetadata[api.ContentType] = objects[0].ContentType
|
||||
}
|
||||
|
||||
pr, pw := io.Pipe()
|
||||
done := make(chan bool)
|
||||
uploadCompleted := false
|
||||
|
||||
/* We will keep "S3-Upload-Id" attribute in completed object to determine is it "common" object or completed object.
|
||||
We will need to differ these objects if something goes wrong during completing multipart upload.
|
||||
I.e. we had completed the object but didn't put tagging/acl for some reason */
|
||||
|
@ -242,11 +279,18 @@ func (n *layer) CompleteMultipartUpload(ctx context.Context, p *CompleteMultipar
|
|||
delete(initMetadata, UploadKeyAttributeName)
|
||||
delete(initMetadata, attrVersionsIgnore)
|
||||
|
||||
go func(done chan bool) {
|
||||
r := &multiObjectReader{
|
||||
ctx: ctx,
|
||||
layer: n,
|
||||
parts: parts,
|
||||
}
|
||||
|
||||
r.prm.cid = p.Info.Bkt.CID
|
||||
|
||||
obj, err = n.objectPut(ctx, p.Info.Bkt, &PutObjectParams{
|
||||
Bucket: p.Info.Bkt.Name,
|
||||
Object: p.Info.Key,
|
||||
Reader: pr,
|
||||
Reader: r,
|
||||
Header: initMetadata,
|
||||
})
|
||||
if err != nil {
|
||||
|
@ -254,34 +298,7 @@ func (n *layer) CompleteMultipartUpload(ctx context.Context, p *CompleteMultipar
|
|||
zap.String("uploadID", p.Info.UploadID),
|
||||
zap.String("uploadKey", p.Info.Key),
|
||||
zap.Error(err))
|
||||
done <- true
|
||||
return
|
||||
}
|
||||
uploadCompleted = true
|
||||
done <- true
|
||||
}(done)
|
||||
|
||||
var prmGet getParams
|
||||
prmGet.w = pw
|
||||
prmGet.cid = p.Info.Bkt.CID
|
||||
|
||||
for _, part := range parts {
|
||||
prmGet.oid = part.ID
|
||||
|
||||
err = n.objectWritePayload(ctx, prmGet)
|
||||
if err != nil {
|
||||
_ = pw.Close()
|
||||
n.log.Error("could not download a part of multipart upload",
|
||||
zap.String("uploadID", p.Info.UploadID),
|
||||
zap.String("part number", part.Headers[UploadPartNumberAttributeName]),
|
||||
zap.Error(err))
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
_ = pw.Close()
|
||||
<-done
|
||||
|
||||
if !uploadCompleted {
|
||||
return nil, errors.GetAPIError(errors.ErrInternalError)
|
||||
}
|
||||
|
||||
|
|
|
@ -27,8 +27,6 @@ type (
|
|||
}
|
||||
|
||||
getParams struct {
|
||||
w io.Writer
|
||||
|
||||
// payload range
|
||||
off, ln uint64
|
||||
|
||||
|
@ -114,9 +112,9 @@ func (n *layer) objectHead(ctx context.Context, idCnr *cid.ID, idObj *oid.ID) (*
|
|||
return res.Head, nil
|
||||
}
|
||||
|
||||
// writes payload part of the NeoFS object to the provided io.Writer.
|
||||
// initializes payload reader of the NeoFS object.
|
||||
// Zero range corresponds to full payload (panics if only offset is set).
|
||||
func (n *layer) objectWritePayload(ctx context.Context, p getParams) error {
|
||||
func (n *layer) initObjectPayloadReader(ctx context.Context, p getParams) (io.Reader, error) {
|
||||
prm := PrmObjectRead{
|
||||
Container: *p.cid,
|
||||
Object: *p.oid,
|
||||
|
@ -127,21 +125,11 @@ func (n *layer) objectWritePayload(ctx context.Context, p getParams) error {
|
|||
n.prepareAuthParameters(ctx, &prm.PrmAuth)
|
||||
|
||||
res, err := n.neoFS.ReadObject(ctx, prm)
|
||||
if err == nil {
|
||||
defer res.Payload.Close()
|
||||
|
||||
if p.ln == 0 {
|
||||
p.ln = 4096 // configure?
|
||||
if err != nil {
|
||||
return nil, n.transformNeofsError(ctx, err)
|
||||
}
|
||||
|
||||
// alloc buffer for copying
|
||||
buf := make([]byte, p.ln) // sync-pool it?
|
||||
|
||||
// copy full payload
|
||||
_, err = io.CopyBuffer(p.w, res.Payload, buf)
|
||||
}
|
||||
|
||||
return n.transformNeofsError(ctx, err)
|
||||
return res.Payload, nil
|
||||
}
|
||||
|
||||
// objectGet returns an object with payload in the object.
|
||||
|
|
Loading…
Reference in a new issue