fix(submit): 分块上传的预算检查复用一分钟内读过的 blob 大小,存储读不到时回 503 让面板自动重试
This commit is contained in:
13 files changed
+335
-25
No files matched your search
@@ -0,0 +1,175 @@
|
||||
package submit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// countingBlobs counts the Size reads per id.
|
||||
type countingBlobs struct {
|
||||
*fakeBlobs
|
||||
sizes map[string]int
|
||||
}
|
||||
|
||||
func (c *countingBlobs) Size(ctx context.Context, id string) (int64, bool, error) {
|
||||
c.sizes[id]++
|
||||
return c.fakeBlobs.Size(ctx, id)
|
||||
}
|
||||
|
||||
// putThenFail stores the bytes and still fails, like an object-store upload that
|
||||
// completed but whose answer was lost.
|
||||
type putThenFail struct{ *fakeBlobs }
|
||||
|
||||
func (p putThenFail) Put(ctx context.Context, id string, r io.Reader) (int64, error) {
|
||||
if _, err := p.fakeBlobs.Put(ctx, id, r); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return 0, errors.New("put: connection reset")
|
||||
}
|
||||
|
||||
// The budget check ahead of each part reads every other blob's size once per
|
||||
// blobSizeTTL; the one on completion reads them all from the store.
|
||||
func TestPartBudgetRemembersBlobSizes(t *testing.T) {
|
||||
m, _, fb, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
clock := testNow
|
||||
m.Now = func() time.Time { return clock }
|
||||
b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-2"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fb.stored[b.ID] = []byte("xyz")
|
||||
cb := &countingBlobs{fakeBlobs: fb, sizes: map[string]int{}}
|
||||
m.Blobs = cb
|
||||
|
||||
sendPart(t, m, a, 0, "\x1f\x8b\x08\x00")
|
||||
sendPart(t, m, a, 4, "abcd")
|
||||
if n := cb.sizes[b.ID]; n != 1 {
|
||||
t.Fatalf("two parts read the other blob's size %d times, want once", n)
|
||||
}
|
||||
clock = clock.Add(blobSizeTTL)
|
||||
sendPart(t, m, a, 8, "ef")
|
||||
if n := cb.sizes[b.ID]; n != 2 {
|
||||
t.Fatalf("a part once the TTL ran out: %d reads in all, want 2", n)
|
||||
}
|
||||
if _, err := m.CompleteUpload(ctx, a, "user-1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n := cb.sizes[b.ID]; n != 3 {
|
||||
t.Fatalf("completion within the TTL: %d reads in all, want 3 (it reads the store)", n)
|
||||
}
|
||||
}
|
||||
|
||||
// What this Manager stores or deletes counts in the next part's budget at once,
|
||||
// without waiting out the TTL.
|
||||
func TestPartBudgetSeesThisManagersOwnWritesAtOnce(t *testing.T) {
|
||||
m, _, _, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
m.MaxStoredBytesPerUser = 10
|
||||
m.PartMaxBytes = 16
|
||||
b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-1"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The first part remembers B as holding nothing.
|
||||
sendPart(t, m, a, 0, "\x1f\x8b")
|
||||
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader("\x1f\x8b\x08\x00ab")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := m.UploadPart(ctx, a, "user-1", 2, strings.NewReader("cde")); !errors.Is(err, ErrQuotaExceeded) {
|
||||
t.Fatalf("5 bytes staged beside B's 6 under a 10-byte budget = %v, want ErrQuotaExceeded", err)
|
||||
}
|
||||
if _, err := m.Withdraw(ctx, b.ID, "user-1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Its row is gone, so nothing counts it; nor is its size kept, or the
|
||||
// remembered sizes would grow with every submission ever deleted.
|
||||
if _, kept := m.sizes[b.ID]; kept {
|
||||
t.Fatal("a withdrawn submission's blob size is still remembered")
|
||||
}
|
||||
if got := sendPart(t, m, a, 2, "cdefgh"); got != 8 {
|
||||
t.Fatalf("staged %d after B was withdrawn, want 8", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPartBudgetSeesAReapedBlobAtOnce(t *testing.T) {
|
||||
m, st, _, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
m.MaxStoredBytesPerUser = 10
|
||||
m.PartMaxBytes = 16
|
||||
b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-1"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader("\x1f\x8b\x08\x00ab")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
reviewed := testNow.Add(-48 * time.Hour)
|
||||
st.subs[b.ID].Status = StatusRejected
|
||||
st.subs[b.ID].ReviewedAt = &reviewed
|
||||
if n, err := m.ReapRejected(ctx, 24*time.Hour); n != 1 || err != nil {
|
||||
t.Fatalf("ReapRejected = %d, %v; want 1, nil", n, err)
|
||||
}
|
||||
if got := sendPart(t, m, a, 0, "\x1f\x8b\x08\x00abcd"); got != 8 {
|
||||
t.Fatalf("staged %d after B's blob was reaped, want 8", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A Put that failed may still have replaced the blob, so its size is read again.
|
||||
func TestPartBudgetRereadsABlobAfterAFailedPut(t *testing.T) {
|
||||
m, _, fb, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
m.MaxStoredBytesPerUser = 10
|
||||
m.PartMaxBytes = 16
|
||||
b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-1"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader("\x1f\x8b\x08\x00ab")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
m.Blobs = putThenFail{fb}
|
||||
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader("\x1f\x8b")); err == nil {
|
||||
t.Fatal("setup: the failing Put succeeded")
|
||||
}
|
||||
m.Blobs = fb
|
||||
if got := sendPart(t, m, a, 0, "\x1f\x8b\x08\x00abcd"); got != 8 {
|
||||
t.Fatalf("staged %d beside B's 2 stored bytes, want 8", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A size that cannot be read fails the check as ErrStoreUnavailable, which the
|
||||
// API answers 503 for the client to send again, never as a bare error (500).
|
||||
func TestBudgetReadFailuresAreRetryable(t *testing.T) {
|
||||
m, _, fb, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
if _, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-2"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
fb.sizeErr = errors.New("dial tcp: i/o timeout")
|
||||
if _, err := m.UploadPart(ctx, a, "user-1", 0, strings.NewReader("\x1f\x8b\x08\x00")); !errors.Is(err, ErrStoreUnavailable) {
|
||||
t.Fatalf("a part with the blob store down = %v, want ErrStoreUnavailable", err)
|
||||
}
|
||||
if _, err := m.UploadContext(ctx, a, "user-1", strings.NewReader("\x1f\x8b\x08\x00")); !errors.Is(err, ErrStoreUnavailable) {
|
||||
t.Fatalf("an upload with the blob store down = %v, want ErrStoreUnavailable", err)
|
||||
}
|
||||
if _, ok := fb.stored[a]; ok {
|
||||
t.Fatal("an upload whose budget could not be checked was stored")
|
||||
}
|
||||
|
||||
fb.sizeErr = nil
|
||||
notDir := filepath.Join(t.TempDir(), "parts")
|
||||
if err := os.WriteFile(notDir, nil, 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
m.Parts = &PartStore{Dir: notDir}
|
||||
if _, err := m.UploadContext(ctx, a, "user-1", strings.NewReader("\x1f\x8b\x08\x00")); !errors.Is(err, ErrStoreUnavailable) {
|
||||
t.Fatalf("an upload with the staging directory unreadable = %v, want ErrStoreUnavailable", err)
|
||||
}
|
||||
}
|
||||
+89
-13
@@ -67,6 +67,7 @@ import (
|
||||
"log"
|
||||
"regexp"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/build"
|
||||
@@ -109,6 +110,11 @@ var (
|
||||
// was never uploaded). The internal context-fetch route maps it to 404, the
|
||||
// same distinction Exists draws for Approve.
|
||||
ErrBlobNotFound = errors.New("submit: context blob not found")
|
||||
// ErrStoreUnavailable means the sizes the storage budgets are checked against
|
||||
// could not be read: the blob store (an object store, most likely) or the
|
||||
// staging directory failed to answer. Nothing was written and the request can
|
||||
// simply be sent again, so the API answers 503 and the panel retries.
|
||||
ErrStoreUnavailable = errors.New("submit: the uploads store did not answer")
|
||||
)
|
||||
|
||||
// invalidf wraps ErrInvalid so every malformed-request case maps to one 400.
|
||||
@@ -335,8 +341,9 @@ type Blobs interface {
|
||||
Open(ctx context.Context, id string) (io.ReadCloser, error)
|
||||
}
|
||||
|
||||
// Manager orchestrates the approval lane. It holds no mutable state; the clock
|
||||
// and id generator are injectable for hermetic tests.
|
||||
// Manager orchestrates the approval lane. Its only mutable state is the blob
|
||||
// sizes it remembers for the budget checks (blobSize); the clock and id
|
||||
// generator are injectable for hermetic tests. Use it by pointer.
|
||||
type Manager struct {
|
||||
Store Store
|
||||
Builds Builds
|
||||
@@ -390,8 +397,27 @@ type Manager struct {
|
||||
|
||||
Now func() time.Time
|
||||
IDGen func() string
|
||||
|
||||
// sizes remembers each blob's size as last read or written here, so the
|
||||
// budget check ahead of every chunked-upload part does not stat every blob
|
||||
// in the store again (blobSize).
|
||||
sizesMu sync.Mutex
|
||||
sizes map[string]knownSize
|
||||
}
|
||||
|
||||
// knownSize is a blob's size and when this Manager last read or wrote it.
|
||||
type knownSize struct {
|
||||
n int64
|
||||
at time.Time
|
||||
}
|
||||
|
||||
// blobSizeTTL is how long a blob size read from the store stands in for it in
|
||||
// the budget check ahead of a chunked-upload part. The blobs this Manager writes
|
||||
// or deletes update it at once; only another api replica's changes wait out the
|
||||
// TTL, and the check that decides what is stored (UploadContext, which
|
||||
// CompleteUpload goes through) always reads the store.
|
||||
const blobSizeTTL = time.Minute
|
||||
|
||||
// ContextLimit is the effective cap on one uploaded context.
|
||||
func (m *Manager) ContextLimit() int64 { return m.maxContextBytes() }
|
||||
|
||||
@@ -602,7 +628,7 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i
|
||||
return nil, invalidf("build context must be a gzip-compressed tarball (.tar.gz)")
|
||||
}
|
||||
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id)
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id, true)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -614,9 +640,14 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i
|
||||
}
|
||||
}
|
||||
h := sha256.New()
|
||||
if _, err := m.Blobs.Put(ctx, id, io.TeeReader(&cappedReader{r: br, left: limit, over: over}, h)); err != nil {
|
||||
n, err := m.Blobs.Put(ctx, id, io.TeeReader(&cappedReader{r: br, left: limit, over: over}, h))
|
||||
if err != nil {
|
||||
// A failed Put may or may not have replaced the blob; the next read
|
||||
// finds out.
|
||||
m.forgetSize(id)
|
||||
return nil, err
|
||||
}
|
||||
m.noteSize(id, n)
|
||||
digest := hex.EncodeToString(h.Sum(nil))
|
||||
won, err := m.Store.SetContextDigest(ctx, id, digest)
|
||||
if err != nil {
|
||||
@@ -658,9 +689,10 @@ func (m *Manager) ownPending(ctx context.Context, id, submittedBy string) (*Subm
|
||||
// package: a burst that reaches two api replicas (or any direct caller of the
|
||||
// Manager) can overshoot by up to one blob per interleaved upload — each write
|
||||
// still bounded by the single-blob cap — while the API's per-user upload
|
||||
// reservation collapses the single-replica case.
|
||||
func (m *Manager) uploadLimit(ctx context.Context, submittedBy, id string) (int64, error, error) {
|
||||
used, total, err := m.storedBytes(ctx, submittedBy, id)
|
||||
// reservation collapses the single-replica case. fresh reads every blob's size
|
||||
// from the store; otherwise a size read within blobSizeTTL stands in for it.
|
||||
func (m *Manager) uploadLimit(ctx context.Context, submittedBy, id string, fresh bool) (int64, error, error) {
|
||||
used, total, err := m.storedBytes(ctx, submittedBy, id, fresh)
|
||||
if err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
@@ -735,7 +767,9 @@ func (m *Manager) UploadPart(ctx context.Context, id, submittedBy string, offset
|
||||
return UploadProgress{}, invalidf("build context must be a gzip-compressed tarball (.tar.gz)")
|
||||
}
|
||||
}
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id)
|
||||
// Remembered blob sizes: a part is one of many, and what is finally stored is
|
||||
// checked against the store itself on completion.
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id, false)
|
||||
if err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
@@ -798,8 +832,10 @@ func (m *Manager) ReapStaleParts(olderThan time.Duration) (int, error) {
|
||||
// the sums cannot drift from what is actually occupying the volume (including
|
||||
// blobs uploaded before any budget existed). A staged chunked upload counts as
|
||||
// well, so parts spread over several pending submissions cannot hold more than
|
||||
// the budget allows.
|
||||
func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string) (user, total int64, err error) {
|
||||
// the budget allows. With fresh unset, a blob size read or written within
|
||||
// blobSizeTTL is used as it stands (blobSize). A size that cannot be read is
|
||||
// ErrStoreUnavailable, for the caller to try again.
|
||||
func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string, fresh bool) (user, total int64, err error) {
|
||||
subs, err := m.Store.ListSubmissions(ctx)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
@@ -808,14 +844,14 @@ func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string
|
||||
if s.ID == excludeID {
|
||||
continue
|
||||
}
|
||||
n, _, err := m.Blobs.Size(ctx, s.ID)
|
||||
n, err := m.blobSize(ctx, s.ID, fresh)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
if m.Parts != nil {
|
||||
staged, _, err := m.Parts.Size(s.ID)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
return 0, 0, fmt.Errorf("%w: size of the staged upload of %s: %v", ErrStoreUnavailable, s.ID, err)
|
||||
}
|
||||
n += staged
|
||||
}
|
||||
@@ -827,6 +863,43 @@ func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string
|
||||
return user, total, nil
|
||||
}
|
||||
|
||||
// blobSize is id's stored blob size: the one remembered from within blobSizeTTL
|
||||
// unless fresh, else read from the store and remembered.
|
||||
func (m *Manager) blobSize(ctx context.Context, id string, fresh bool) (int64, error) {
|
||||
now := m.now()
|
||||
if !fresh {
|
||||
m.sizesMu.Lock()
|
||||
k, ok := m.sizes[id]
|
||||
m.sizesMu.Unlock()
|
||||
if ok && now.Sub(k.at) < blobSizeTTL {
|
||||
return k.n, nil
|
||||
}
|
||||
}
|
||||
n, _, err := m.Blobs.Size(ctx, id)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("%w: size of the context of %s: %v", ErrStoreUnavailable, id, err)
|
||||
}
|
||||
m.noteSizeAt(id, n, now)
|
||||
return n, nil
|
||||
}
|
||||
|
||||
func (m *Manager) noteSize(id string, n int64) { m.noteSizeAt(id, n, m.now()) }
|
||||
|
||||
func (m *Manager) noteSizeAt(id string, n int64, at time.Time) {
|
||||
m.sizesMu.Lock()
|
||||
defer m.sizesMu.Unlock()
|
||||
if m.sizes == nil {
|
||||
m.sizes = map[string]knownSize{}
|
||||
}
|
||||
m.sizes[id] = knownSize{n: n, at: at}
|
||||
}
|
||||
|
||||
func (m *Manager) forgetSize(id string) {
|
||||
m.sizesMu.Lock()
|
||||
defer m.sizesMu.Unlock()
|
||||
delete(m.sizes, id)
|
||||
}
|
||||
|
||||
// ReapRejected deletes the uploaded context of every submission rejected more
|
||||
// than olderThan ago, keeping the row. It returns how many blobs it deleted and
|
||||
// carries on past a blob it cannot delete, reporting the first such error.
|
||||
@@ -848,6 +921,7 @@ func (m *Manager) ReapRejected(ctx context.Context, olderThan time.Duration) (in
|
||||
_, ok, err := m.Blobs.Size(ctx, s.ID)
|
||||
if err == nil && ok {
|
||||
if err = m.Blobs.Delete(ctx, s.ID); err == nil {
|
||||
m.noteSize(s.ID, 0)
|
||||
reaped++
|
||||
}
|
||||
}
|
||||
@@ -1133,7 +1207,9 @@ func (m *Manager) deleteBlob(ctx context.Context, id string) error {
|
||||
if m.Blobs == nil {
|
||||
return nil
|
||||
}
|
||||
if err := m.Blobs.Delete(ctx, id); err != nil {
|
||||
err := m.Blobs.Delete(ctx, id)
|
||||
m.forgetSize(id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("submit: submission removed, but its uploaded context could not be deleted (it may remain on the uploads store): %w", err)
|
||||
}
|
||||
return nil
|
||||
|
||||
Reference in new issue
Block a user