|
|
@ -1,9 +1,12 @@ |
|
|
|
package s3api |
|
|
|
|
|
|
|
import ( |
|
|
|
"cmp" |
|
|
|
"encoding/hex" |
|
|
|
"encoding/xml" |
|
|
|
"fmt" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/stats" |
|
|
|
"golang.org/x/exp/slices" |
|
|
|
"math" |
|
|
|
"path/filepath" |
|
|
|
"sort" |
|
|
@ -11,18 +14,18 @@ import ( |
|
|
|
"strings" |
|
|
|
"time" |
|
|
|
|
|
|
|
"github.com/google/uuid" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err" |
|
|
|
"golang.org/x/exp/slices" |
|
|
|
|
|
|
|
"github.com/aws/aws-sdk-go/aws" |
|
|
|
"github.com/aws/aws-sdk-go/service/s3" |
|
|
|
"github.com/google/uuid" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err" |
|
|
|
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/filer" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog" |
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" |
|
|
|
) |
|
|
|
|
|
|
|
const multipartExt = ".part" |
|
|
|
|
|
|
|
type InitiateMultipartUploadResult struct { |
|
|
|
XMLName xml.Name `xml:"http://s3.amazonaws.com/doc/2006-03-01/ InitiateMultipartUploadResult"` |
|
|
|
s3.CreateMultipartUploadOutput |
|
|
@ -71,74 +74,97 @@ type CompleteMultipartUploadResult struct { |
|
|
|
func (s3a *S3ApiServer) completeMultipartUpload(input *s3.CompleteMultipartUploadInput, parts *CompleteMultipartUpload) (output *CompleteMultipartUploadResult, code s3err.ErrorCode) { |
|
|
|
|
|
|
|
glog.V(2).Infof("completeMultipartUpload input %v", input) |
|
|
|
|
|
|
|
completedParts := parts.Parts |
|
|
|
slices.SortFunc(completedParts, func(a, b CompletedPart) int { |
|
|
|
return a.PartNumber - b.PartNumber |
|
|
|
}) |
|
|
|
|
|
|
|
completedPartNumbers := []int{} |
|
|
|
completedPartMap := make(map[int][]string) |
|
|
|
for _, part := range parts.Parts { |
|
|
|
if _, ok := completedPartMap[part.PartNumber]; !ok { |
|
|
|
completedPartNumbers = append(completedPartNumbers, part.PartNumber) |
|
|
|
} |
|
|
|
completedPartMap[part.PartNumber] = append(completedPartMap[part.PartNumber], part.ETag) |
|
|
|
} |
|
|
|
uploadDirectory := s3a.genUploadsFolder(*input.Bucket) + "/" + *input.UploadId |
|
|
|
|
|
|
|
entries, _, err := s3a.list(uploadDirectory, "", "", false, maxPartsList) |
|
|
|
if err != nil || len(entries) == 0 { |
|
|
|
glog.Errorf("completeMultipartUpload %s %s error: %v, entries:%d", *input.Bucket, *input.UploadId, err, len(entries)) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedNoSuchUpload).Inc() |
|
|
|
return nil, s3err.ErrNoSuchUpload |
|
|
|
} |
|
|
|
|
|
|
|
pentry, err := s3a.getEntry(s3a.genUploadsFolder(*input.Bucket), *input.UploadId) |
|
|
|
if err != nil { |
|
|
|
glog.Errorf("completeMultipartUpload %s %s error: %v", *input.Bucket, *input.UploadId, err) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedNoSuchUpload).Inc() |
|
|
|
return nil, s3err.ErrNoSuchUpload |
|
|
|
} |
|
|
|
|
|
|
|
deleteEntries := []*filer_pb.Entry{} |
|
|
|
partEntries := make(map[int][]*filer_pb.Entry, len(entries)) |
|
|
|
for _, entry := range entries { |
|
|
|
foundEntry := false |
|
|
|
glog.V(4).Infof("completeMultipartUpload part entries %s", entry.Name) |
|
|
|
if strings.HasSuffix(entry.Name, ".part") && !entry.IsDirectory { |
|
|
|
var partNumberString string |
|
|
|
index := strings.Index(entry.Name, "_") |
|
|
|
if index != -1 { |
|
|
|
partNumberString = entry.Name[:index] |
|
|
|
} else { |
|
|
|
partNumberString = entry.Name[:len(entry.Name)-len(".part")] |
|
|
|
if entry.IsDirectory || !strings.HasSuffix(entry.Name, multipartExt) { |
|
|
|
continue |
|
|
|
} |
|
|
|
partNumber, err := strconv.Atoi(partNumberString) |
|
|
|
partNumber, err := parsePartNumber(entry.Name) |
|
|
|
if err != nil { |
|
|
|
glog.Errorf("completeMultipartUpload failed to pasre partNumber %s:%s", partNumberString, err) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedPartNumber).Inc() |
|
|
|
glog.Errorf("completeMultipartUpload failed to pasre partNumber %s:%s", entry.Name, err) |
|
|
|
continue |
|
|
|
} |
|
|
|
completedPartsByNumber, ok := completedPartMap[partNumber] |
|
|
|
if !ok { |
|
|
|
continue |
|
|
|
} |
|
|
|
for _, partETag := range completedPartsByNumber { |
|
|
|
partETag = strings.Trim(partETag, `"`) |
|
|
|
entryETag := hex.EncodeToString(entry.Attributes.GetMd5()) |
|
|
|
if partETag != "" && len(partETag) == 32 && entryETag != "" { |
|
|
|
if entryETag != partETag { |
|
|
|
glog.Errorf("completeMultipartUpload %s ETag mismatch chunk: %s part: %s", entry.Name, entryETag, partETag) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedEtagMismatch).Inc() |
|
|
|
continue |
|
|
|
} |
|
|
|
} else { |
|
|
|
glog.Warningf("invalid complete etag %s, partEtag %s", partETag, entryETag) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedEtagInvalid).Inc() |
|
|
|
} |
|
|
|
if len(entry.Chunks) == 0 { |
|
|
|
glog.Warningf("completeMultipartUpload %s empty chunks", entry.Name) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedPartEmpty).Inc() |
|
|
|
continue |
|
|
|
} |
|
|
|
//there maybe multi same part, because of client retry
|
|
|
|
partEntries[partNumber] = append(partEntries[partNumber], entry) |
|
|
|
foundEntry = true |
|
|
|
} |
|
|
|
if !foundEntry { |
|
|
|
deleteEntries = append(deleteEntries, entry) |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
mime := pentry.Attributes.Mime |
|
|
|
|
|
|
|
var finalParts []*filer_pb.FileChunk |
|
|
|
var offset int64 |
|
|
|
var deleteEntries []*filer_pb.Entry |
|
|
|
for _, part := range completedParts { |
|
|
|
entries := partEntries[part.PartNumber] |
|
|
|
// check whether completedParts is more than received parts
|
|
|
|
if len(entries) == 0 { |
|
|
|
glog.Errorf("part %d has no entry", part.PartNumber) |
|
|
|
sort.Ints(completedPartNumbers) |
|
|
|
for _, partNumber := range completedPartNumbers { |
|
|
|
partEntriesByNumber, ok := partEntries[partNumber] |
|
|
|
if !ok { |
|
|
|
glog.Errorf("part %d has no entry", partNumber) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedPartNotFound).Inc() |
|
|
|
return nil, s3err.ErrInvalidPart |
|
|
|
} |
|
|
|
|
|
|
|
found := false |
|
|
|
for _, entry := range entries { |
|
|
|
if found { |
|
|
|
deleteEntries = append(deleteEntries, entry) |
|
|
|
continue |
|
|
|
if len(partEntriesByNumber) > 1 { |
|
|
|
slices.SortFunc(partEntriesByNumber, func(a, b *filer_pb.Entry) int { |
|
|
|
return cmp.Compare(b.Chunks[0].ModifiedTsNs, a.Chunks[0].ModifiedTsNs) |
|
|
|
}) |
|
|
|
} |
|
|
|
|
|
|
|
partETag := strings.Trim(part.ETag, `"`) |
|
|
|
entryETag := hex.EncodeToString(entry.Attributes.GetMd5()) |
|
|
|
glog.Warningf("complete etag %s, partEtag %s", partETag, entryETag) |
|
|
|
if partETag != "" && len(partETag) == 32 && entryETag != "" && entryETag != partETag { |
|
|
|
err = fmt.Errorf("completeMultipartUpload %s ETag mismatch chunk: %s part: %s", entry.Name, entryETag, partETag) |
|
|
|
for _, entry := range partEntriesByNumber { |
|
|
|
if found { |
|
|
|
deleteEntries = append(deleteEntries, entry) |
|
|
|
stats.S3HandlerCounter.WithLabelValues(stats.ErrorCompletedPartEntryMismatch).Inc() |
|
|
|
continue |
|
|
|
} |
|
|
|
for _, chunk := range entry.GetChunks() { |
|
|
@ -154,13 +180,7 @@ func (s3a *S3ApiServer) completeMultipartUpload(input *s3.CompleteMultipartUploa |
|
|
|
offset += int64(chunk.Size) |
|
|
|
} |
|
|
|
found = true |
|
|
|
err = nil |
|
|
|
} |
|
|
|
if err != nil { |
|
|
|
glog.Errorf("%s", err) |
|
|
|
return nil, s3err.ErrInvalidPart |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
entryName := filepath.Base(*input.Key) |
|
|
@ -223,29 +243,15 @@ func (s3a *S3ApiServer) completeMultipartUpload(input *s3.CompleteMultipartUploa |
|
|
|
return |
|
|
|
} |
|
|
|
|
|
|
|
func findByPartNumber(fileName string, parts []CompletedPart) (etag string, found bool) { |
|
|
|
partNumber, formatErr := strconv.Atoi(fileName[:4]) |
|
|
|
if formatErr != nil { |
|
|
|
return |
|
|
|
} |
|
|
|
x := sort.Search(len(parts), func(i int) bool { |
|
|
|
return parts[i].PartNumber >= partNumber |
|
|
|
}) |
|
|
|
if x >= len(parts) { |
|
|
|
return |
|
|
|
} |
|
|
|
if parts[x].PartNumber != partNumber { |
|
|
|
return |
|
|
|
} |
|
|
|
y := 0 |
|
|
|
for i, part := range parts[x:] { |
|
|
|
if part.PartNumber == partNumber { |
|
|
|
y = i |
|
|
|
func parsePartNumber(fileName string) (int, error) { |
|
|
|
var partNumberString string |
|
|
|
index := strings.Index(fileName, "_") |
|
|
|
if index != -1 { |
|
|
|
partNumberString = fileName[:index] |
|
|
|
} else { |
|
|
|
break |
|
|
|
partNumberString = fileName[:len(fileName)-len(multipartExt)] |
|
|
|
} |
|
|
|
} |
|
|
|
return parts[x+y].ETag, true |
|
|
|
return strconv.Atoi(partNumberString) |
|
|
|
} |
|
|
|
|
|
|
|
func (s3a *S3ApiServer) abortMultipartUpload(input *s3.AbortMultipartUploadInput) (output *s3.AbortMultipartUploadOutput, code s3err.ErrorCode) { |
|
|
@ -361,7 +367,7 @@ func (s3a *S3ApiServer) listObjectParts(input *s3.ListPartsInput) (output *ListP |
|
|
|
StorageClass: aws.String("STANDARD"), |
|
|
|
} |
|
|
|
|
|
|
|
entries, isLast, err := s3a.list(s3a.genUploadsFolder(*input.Bucket)+"/"+*input.UploadId, "", fmt.Sprintf("%04d.part", *input.PartNumberMarker), false, uint32(*input.MaxParts)) |
|
|
|
entries, isLast, err := s3a.list(s3a.genUploadsFolder(*input.Bucket)+"/"+*input.UploadId, "", fmt.Sprintf("%04d%s", *input.PartNumberMarker, multipartExt), false, uint32(*input.MaxParts)) |
|
|
|
if err != nil { |
|
|
|
glog.Errorf("listObjectParts %s %s error: %v", *input.Bucket, *input.UploadId, err) |
|
|
|
return nil, s3err.ErrNoSuchUpload |
|
|
@ -373,15 +379,8 @@ func (s3a *S3ApiServer) listObjectParts(input *s3.ListPartsInput) (output *ListP |
|
|
|
output.IsTruncated = aws.Bool(!isLast) |
|
|
|
|
|
|
|
for _, entry := range entries { |
|
|
|
if strings.HasSuffix(entry.Name, ".part") && !entry.IsDirectory { |
|
|
|
var partNumberString string |
|
|
|
index := strings.Index(entry.Name, "_") |
|
|
|
if index != -1 { |
|
|
|
partNumberString = entry.Name[:index] |
|
|
|
} else { |
|
|
|
partNumberString = entry.Name[:len(entry.Name)-len(".part")] |
|
|
|
} |
|
|
|
partNumber, err := strconv.Atoi(partNumberString) |
|
|
|
if strings.HasSuffix(entry.Name, multipartExt) && !entry.IsDirectory { |
|
|
|
partNumber, err := parsePartNumber(entry.Name) |
|
|
|
if err != nil { |
|
|
|
glog.Errorf("listObjectParts %s %s parse %s: %v", *input.Bucket, *input.UploadId, entry.Name, err) |
|
|
|
continue |
|
|
|