fix(submit): S3 后端按 8MiB 分片上传,大包不再撑爆 felis-api 内存,超出分片上限时报错并删掉截断的对象
This commit is contained in:
3 files changed
+194
-7
No files matched your search
@@ -155,21 +155,49 @@ func (s *S3ContextStore) keyFor(id string) (string, error) {
|
||||
return path.Join(s.prefix, id, contextBlobName), nil
|
||||
}
|
||||
|
||||
// s3PartSize is the size every upload is cut into, and so the one buffer an upload
|
||||
// holds in memory. Left unset, minio-go sizes the parts of a stream of unknown
|
||||
// length for a 5TiB object (528MiB each) and allocates one whole part before the
|
||||
// first byte arrives, twice felis-api's 256Mi memory limit: the first large modpack
|
||||
// would OOMKill the panel.
|
||||
const s3PartSize = 8 << 20
|
||||
|
||||
// s3MaxObject is what the 10,000 parts S3 allows hold at s3PartSize: about 78GiB,
|
||||
// far past the 1GiB default context cap.
|
||||
const s3MaxObject = s3PartSize * 10000
|
||||
|
||||
// Put streams r to the derived object key. Size is unknown (the Manager hands us a
|
||||
// size-capped reader), so it is uploaded with size -1 (multipart). PutObject is
|
||||
// atomic from a reader's perspective — a partial upload never becomes a readable
|
||||
// object — so a failed or oversize upload never replaces a good context. It
|
||||
// returns the number of bytes stored.
|
||||
// size-capped reader), so it is uploaded as s3PartSize parts. PutObject is atomic
|
||||
// from a reader's perspective — a partial upload never becomes a readable object —
|
||||
// so a failed or oversize upload never replaces a good context. It returns the
|
||||
// number of bytes stored.
|
||||
func (s *S3ContextStore) Put(ctx context.Context, id string, r io.Reader) (int64, error) {
|
||||
key, err := s.keyFor(id)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
info, err := s.client.PutObject(ctx, s.bucket, key, r, -1, minio.PutObjectOptions{ContentType: "application/gzip"})
|
||||
info, err := s.client.PutObject(ctx, s.bucket, key, r, -1, minio.PutObjectOptions{ContentType: "application/gzip", PartSize: s3PartSize})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("submit: put context blob: %w", err)
|
||||
}
|
||||
return info.Size, nil
|
||||
var probe [1]byte
|
||||
n, rerr := io.ReadFull(r, probe[:])
|
||||
if n == 0 && rerr == io.EOF {
|
||||
return info.Size, nil
|
||||
}
|
||||
// minio-go stops reading after the last part S3 allows and completes the object
|
||||
// with what it has, so a longer stream (a context cap configured past
|
||||
// s3MaxObject) is stored cut short and reported as a success. Anything but a
|
||||
// clean end here means the stored object is not the upload.
|
||||
if n > 0 {
|
||||
rerr = fmt.Errorf("submit: context blob exceeds the %d-byte multipart limit", int64(s3MaxObject))
|
||||
} else {
|
||||
rerr = fmt.Errorf("submit: put context blob: %w", rerr)
|
||||
}
|
||||
if err := s.client.RemoveObject(ctx, s.bucket, key, minio.RemoveObjectOptions{}); err != nil {
|
||||
return 0, fmt.Errorf("%w (removing the truncated object also failed: %v)", rerr, err)
|
||||
}
|
||||
return 0, rerr
|
||||
}
|
||||
|
||||
// Exists reports whether a context blob has been stored for id. Approve consults
|
||||
|
||||
@@ -1,11 +1,23 @@
|
||||
package submit
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"hash"
|
||||
"io"
|
||||
"math/rand"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"runtime"
|
||||
"slices"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/minio/minio-go/v7"
|
||||
@@ -19,12 +31,16 @@ type fakeS3 struct {
|
||||
putErr error
|
||||
statErr error // when set, StatObject returns it (e.g. auth rejected / bucket missing)
|
||||
removeErr error
|
||||
putLimit int64 // when set, PutObject stops reading after this many bytes, as minio-go does after the last part
|
||||
}
|
||||
|
||||
func (f *fakeS3) PutObject(_ context.Context, bucket, object string, r io.Reader, _ int64, _ minio.PutObjectOptions) (minio.UploadInfo, error) {
|
||||
if f.putErr != nil {
|
||||
return minio.UploadInfo{}, f.putErr
|
||||
}
|
||||
if f.putLimit > 0 {
|
||||
r = io.LimitReader(r, f.putLimit)
|
||||
}
|
||||
data, err := io.ReadAll(r)
|
||||
if err != nil {
|
||||
return minio.UploadInfo{}, err
|
||||
@@ -322,3 +338,143 @@ func s3KeysOf(m map[string][]byte) []string {
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// multipartS3 is an S3 endpoint speaking just enough of the multipart protocol for
|
||||
// minio-go to upload through it. It keeps no bytes: it records each part's size
|
||||
// and hashes the parts in order, so what it allocates stays out of the way of
|
||||
// what the uploader does.
|
||||
type multipartS3 struct {
|
||||
mu sync.Mutex
|
||||
partSizes []int64
|
||||
sum hash.Hash
|
||||
objects map[string]string // path -> sha256 of the completed object
|
||||
}
|
||||
|
||||
func (m *multipartS3) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
q := r.URL.Query()
|
||||
switch {
|
||||
case r.Method == http.MethodPost && q.Has("uploads"):
|
||||
m.partSizes, m.sum = nil, sha256.New()
|
||||
fmt.Fprint(w, `<InitiateMultipartUploadResult><UploadId>u1</UploadId></InitiateMultipartUploadResult>`)
|
||||
case r.Method == http.MethodPut && q.Has("partNumber"):
|
||||
if n, _ := strconv.Atoi(q.Get("partNumber")); n != len(m.partSizes)+1 {
|
||||
http.Error(w, "part out of order", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
size, err := s3Body(r, m.sum)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
m.partSizes = append(m.partSizes, size)
|
||||
w.Header().Set("ETag", fmt.Sprintf(`"part-%d"`, len(m.partSizes)))
|
||||
case r.Method == http.MethodPost && q.Has("uploadId"):
|
||||
m.objects[r.URL.Path] = hex.EncodeToString(m.sum.Sum(nil))
|
||||
fmt.Fprint(w, `<CompleteMultipartUploadResult><Bucket>felis-uploads</Bucket><ETag>"done"</ETag></CompleteMultipartUploadResult>`)
|
||||
case r.Method == http.MethodDelete && q.Has("uploadId"):
|
||||
m.partSizes = nil
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
default:
|
||||
http.Error(w, r.Method+" "+r.URL.String()+" is not part of this fake", http.StatusNotImplemented)
|
||||
}
|
||||
}
|
||||
|
||||
// s3Body copies a request's payload to w, taking off any aws-chunked framing (the
|
||||
// streaming signature minio-go uses over plain HTTP), and returns its size.
|
||||
func s3Body(r *http.Request, w io.Writer) (int64, error) {
|
||||
if !strings.HasPrefix(r.Header.Get("X-Amz-Content-Sha256"), "STREAMING-") {
|
||||
return io.Copy(w, r.Body)
|
||||
}
|
||||
br := bufio.NewReader(r.Body)
|
||||
var total int64
|
||||
for {
|
||||
line, err := br.ReadString('\n')
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
size, err := strconv.ParseInt(strings.TrimSpace(strings.SplitN(line, ";", 2)[0]), 16, 64)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if size == 0 {
|
||||
return total, nil
|
||||
}
|
||||
if _, err := io.CopyN(w, br, size); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if _, err := br.Discard(2); err != nil { // the chunk's CRLF
|
||||
return 0, err
|
||||
}
|
||||
total += size
|
||||
}
|
||||
}
|
||||
|
||||
// Put cuts an upload into s3PartSize parts, so felis-api holds one small part in
|
||||
// memory however large the modpack. Left to minio-go, a stream of unknown length
|
||||
// is cut into 528MiB parts and one is allocated up front: twice the api's 256Mi
|
||||
// memory limit, which OOMKilled the panel on the first large upload.
|
||||
func TestS3ContextStorePutStreamsInSmallParts(t *testing.T) {
|
||||
fake := &multipartS3{objects: map[string]string{}}
|
||||
ts := httptest.NewServer(fake)
|
||||
defer ts.Close()
|
||||
s, err := NewS3ContextStore(S3StoreConfig{Base: "s3://felis-uploads/builds", Endpoint: ts.URL, Region: "us-east-1", AccessKey: "a", SecretKey: "b"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
payload := make([]byte, 20<<20+123)
|
||||
rand.New(rand.NewSource(1)).Read(payload)
|
||||
want := sha256.Sum256(payload)
|
||||
|
||||
runtime.GC()
|
||||
var before, after runtime.MemStats
|
||||
runtime.ReadMemStats(&before)
|
||||
n, err := s.Put(context.Background(), "sub-big", bytes.NewReader(payload))
|
||||
runtime.ReadMemStats(&after)
|
||||
if err != nil {
|
||||
t.Fatalf("Put: %v", err)
|
||||
}
|
||||
|
||||
if n != int64(len(payload)) {
|
||||
t.Errorf("Put stored %d bytes, want %d", n, len(payload))
|
||||
}
|
||||
if got := fake.objects["/felis-uploads/builds/sub-big/"+contextBlobName]; got != hex.EncodeToString(want[:]) {
|
||||
t.Errorf("stored object differs from the upload (sha256 %q)", got)
|
||||
}
|
||||
if want := []int64{8 << 20, 8 << 20, 4<<20 + 123}; !slices.Equal(fake.partSizes, want) {
|
||||
t.Errorf("upload arrived as parts %v, want %v", fake.partSizes, want)
|
||||
}
|
||||
if alloc := after.TotalAlloc - before.TotalAlloc; alloc > 64<<20 {
|
||||
t.Errorf("uploading 20MiB allocated %d MiB, want one part's worth", alloc>>20)
|
||||
}
|
||||
}
|
||||
|
||||
// When minio-go stops reading before the stream ends (after the last part S3
|
||||
// allows), the object it completed is the upload cut short: Put reports the
|
||||
// failure and removes the object instead of returning a success.
|
||||
func TestS3ContextStorePutRefusesATruncatedObject(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
r io.Reader
|
||||
want error
|
||||
}{
|
||||
{"stream past the multipart limit", strings.NewReader("\x1f\x8b0123456789"), nil},
|
||||
{"past the context cap exactly at the limit", &cappedReader{r: strings.NewReader("\x1f\x8b0123456789"), left: 6, over: errContextTooLarge}, errContextTooLarge},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fake := &fakeS3{putLimit: 6}
|
||||
s := &S3ContextStore{client: fake, bucket: "felis-uploads", prefix: "builds"}
|
||||
n, err := s.Put(context.Background(), "sub-abc", tc.r)
|
||||
if err == nil {
|
||||
t.Fatalf("Put = (%d, nil), want an error", n)
|
||||
}
|
||||
if tc.want != nil && !errors.Is(err, tc.want) {
|
||||
t.Errorf("Put error = %v, want %v", err, tc.want)
|
||||
}
|
||||
if len(fake.objects) != 0 {
|
||||
t.Errorf("truncated object left behind: %v", s3KeysOf(fake.objects))
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user