feat(submit): 模组上传改为分片续传并显示进度,经 Cloudflare 边缘也能传满 1GiB
This commit is contained in:
28 files changed
+2543
-154
No files matched your search
@@ -616,6 +616,13 @@ func (a *API) externalAPIRoutes() []apiRoute {
|
||||
// App-tier and owner-scoped (the id must belong to the principal), exactly
|
||||
// like the create/list routes above.
|
||||
{Method: "POST", Pattern: "/api/v1/me/submissions/{id}/context", h: a.handleUploadSubmissionContext},
|
||||
// The chunked form of that upload, for a context larger than one request
|
||||
// carries through the edge (Cloudflare refuses bodies over 100 MB): GET
|
||||
// reports the staged length (the resume point), PUT ?offset= appends one
|
||||
// part, POST .../complete stores the staged whole. Same owner scoping.
|
||||
{Method: "GET", Pattern: "/api/v1/me/submissions/{id}/context/upload", h: a.handleContextUploadStatus},
|
||||
{Method: "PUT", Pattern: "/api/v1/me/submissions/{id}/context/upload", h: a.handleContextUploadPart},
|
||||
{Method: "POST", Pattern: "/api/v1/me/submissions/{id}/context/upload/complete", h: a.handleContextUploadComplete},
|
||||
// Withdraw the caller's OWN pending submission: the row and its uploaded
|
||||
// context are deleted, freeing the pending slot and storage budget. Same
|
||||
// owner-scoping as the upload route — a reviewed submission is frozen (409)
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"testing"
|
||||
|
||||
"felis.lolicon.best/internal/build"
|
||||
"felis.lolicon.best/internal/submit"
|
||||
"felis.lolicon.best/internal/updates"
|
||||
"sigs.k8s.io/yaml"
|
||||
)
|
||||
@@ -34,26 +35,27 @@ func TestOpenAPISchemasMatchWireStructs(t *testing.T) {
|
||||
}
|
||||
|
||||
pairs := map[string]any{
|
||||
"ServerInfo": ServerInfo{},
|
||||
"FleetServer": fleetServerView{},
|
||||
"MyServerView": MyServerView{},
|
||||
"BackupView": BackupView{},
|
||||
"Build": build.Build{},
|
||||
"Image": build.Image{},
|
||||
"BuildScan": buildScanView{},
|
||||
"ScanSummary": build.ScanSummary{},
|
||||
"ScanPolicy": build.ScanPolicy{},
|
||||
"ScanFinding": build.ScanFinding{},
|
||||
"Submission": submissionView{},
|
||||
"UserView": UserView{},
|
||||
"UserDetail": UserDetail{},
|
||||
"QuotaView": QuotaView{},
|
||||
"SessionView": SessionView{},
|
||||
"PasskeyCredential": passkeyCredentialView{},
|
||||
"UpdateWindow": updateWindow{},
|
||||
"DBBackupStatus": dbBackupView{},
|
||||
"UpdateReport": updateReportView{},
|
||||
"UpdateComponent": updates.ComponentStatus{},
|
||||
"ServerInfo": ServerInfo{},
|
||||
"FleetServer": fleetServerView{},
|
||||
"MyServerView": MyServerView{},
|
||||
"BackupView": BackupView{},
|
||||
"Build": build.Build{},
|
||||
"Image": build.Image{},
|
||||
"BuildScan": buildScanView{},
|
||||
"ScanSummary": build.ScanSummary{},
|
||||
"ScanPolicy": build.ScanPolicy{},
|
||||
"ScanFinding": build.ScanFinding{},
|
||||
"Submission": submissionView{},
|
||||
"ContextUploadProgress": submit.UploadProgress{},
|
||||
"UserView": UserView{},
|
||||
"UserDetail": UserDetail{},
|
||||
"QuotaView": QuotaView{},
|
||||
"SessionView": SessionView{},
|
||||
"PasskeyCredential": passkeyCredentialView{},
|
||||
"UpdateWindow": updateWindow{},
|
||||
"DBBackupStatus": dbBackupView{},
|
||||
"UpdateReport": updateReportView{},
|
||||
"UpdateComponent": updates.ComponentStatus{},
|
||||
}
|
||||
for name, v := range pairs {
|
||||
s, ok := doc.Components.Schemas[name]
|
||||
|
||||
+108
-11
@@ -36,6 +36,14 @@ type SubmissionService interface {
|
||||
// submission at the platform-derived context ref. submittedBy is the principal,
|
||||
// never the body, so a user can only upload to a submission they own.
|
||||
UploadContext(ctx context.Context, id, submittedBy string, r io.Reader) (*submit.Submission, error)
|
||||
// UploadStatus, UploadPart and CompleteUpload are the chunked form of
|
||||
// UploadContext, for a context larger than one request carries through the
|
||||
// edge (Cloudflare refuses bodies over 100 MB): parts are staged in order, the
|
||||
// staged length is the resume point, and completion stores the whole through
|
||||
// UploadContext's checks. Same owner scoping.
|
||||
UploadStatus(ctx context.Context, id, submittedBy string) (submit.UploadProgress, error)
|
||||
UploadPart(ctx context.Context, id, submittedBy string, offset int64, r io.Reader) (submit.UploadProgress, error)
|
||||
CompleteUpload(ctx context.Context, id, submittedBy string) (*submit.Submission, error)
|
||||
// ListBy returns one page of one user's submissions, newest first (the "my
|
||||
// uploads" view). The scope is submittedBy, whatever opts says.
|
||||
ListBy(ctx context.Context, submittedBy string, opts submit.ListOpts) (submit.Page, error)
|
||||
@@ -173,26 +181,105 @@ func (a *API) handleUploadSubmissionContext(w http.ResponseWriter, r *http.Reque
|
||||
// yields exactly one admitted stream per replica. The rollback keeps a failed
|
||||
// upload (aborted transfer, wrong format, spent quota) from burning the
|
||||
// window, so a legit retry after a genuine failure is not punished.
|
||||
lim := a.submitLimiter()
|
||||
reservedAt, ok := lim.reserve(submissionUploadKey+p.UserID, a.SubmitUploadCooldown)
|
||||
commit, release, ok := a.reserveUpload(w, r, p.UserID)
|
||||
if !ok {
|
||||
writeError(w, r, newError(http.StatusTooManyRequests, "submission_cooldown",
|
||||
"an upload was accepted recently; wait a moment before uploading again"))
|
||||
return
|
||||
}
|
||||
committed := false
|
||||
defer func() {
|
||||
if !committed {
|
||||
lim.release(submissionUploadKey+p.UserID, reservedAt)
|
||||
}
|
||||
}()
|
||||
defer release()
|
||||
id := r.PathValue("id")
|
||||
sub, err := a.Submissions.UploadContext(r.Context(), id, p.UserID, r.Body)
|
||||
if err != nil {
|
||||
writeSubmitError(w, r, err)
|
||||
return
|
||||
}
|
||||
committed = true
|
||||
commit()
|
||||
a.audit(r, "submission.upload", sub.ID)
|
||||
writeJSON(w, http.StatusOK, sub)
|
||||
}
|
||||
|
||||
// reserveUpload claims the per-user upload cooldown for userID, answering 429
|
||||
// when an upload landed within it. commit keeps the reservation; release, run
|
||||
// deferred, gives it back unless commit ran.
|
||||
func (a *API) reserveUpload(w http.ResponseWriter, r *http.Request, userID string) (commit, release func(), ok bool) {
|
||||
lim := a.submitLimiter()
|
||||
reservedAt, ok := lim.reserve(submissionUploadKey+userID, a.SubmitUploadCooldown)
|
||||
if !ok {
|
||||
writeError(w, r, newError(http.StatusTooManyRequests, "submission_cooldown",
|
||||
"an upload was accepted recently; wait a moment before uploading again"))
|
||||
return nil, nil, false
|
||||
}
|
||||
committed := false
|
||||
return func() { committed = true }, func() {
|
||||
if !committed {
|
||||
lim.release(submissionUploadKey+userID, reservedAt)
|
||||
}
|
||||
}, true
|
||||
}
|
||||
|
||||
// handleContextUploadStatus reports how far the caller's chunked upload of a
|
||||
// submission has come (app-tier, owner-scoped like the single upload). The
|
||||
// panel reads it before the first part and again after a failed one, and sends
|
||||
// the next part from received.
|
||||
func (a *API) handleContextUploadStatus(w http.ResponseWriter, r *http.Request) {
|
||||
if a.Submissions == nil {
|
||||
writeError(w, r, errSubmissionsUnavailable)
|
||||
return
|
||||
}
|
||||
p := principalFromContext(r.Context())
|
||||
prog, err := a.Submissions.UploadStatus(r.Context(), r.PathValue("id"), p.UserID)
|
||||
if err != nil {
|
||||
writeSubmitError(w, r, err)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, prog)
|
||||
}
|
||||
|
||||
// handleContextUploadPart appends one part of the caller's chunked upload. The
|
||||
// body is the raw bytes; ?offset= is where they start: 0 starts over, anything
|
||||
// else must equal the staged length (409 upload_offset_mismatch otherwise). A
|
||||
// part is small enough for any edge, so no cooldown applies here: the staged
|
||||
// total is bounded by the context cap and the storage budget, and completion
|
||||
// holds the cooldown.
|
||||
func (a *API) handleContextUploadPart(w http.ResponseWriter, r *http.Request) {
|
||||
if a.Submissions == nil {
|
||||
writeError(w, r, errSubmissionsUnavailable)
|
||||
return
|
||||
}
|
||||
offset, err := strconv.ParseInt(r.URL.Query().Get("offset"), 10, 64)
|
||||
if err != nil {
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_request",
|
||||
"offset must be the byte position the part starts at"))
|
||||
return
|
||||
}
|
||||
p := principalFromContext(r.Context())
|
||||
prog, err := a.Submissions.UploadPart(r.Context(), r.PathValue("id"), p.UserID, offset, r.Body)
|
||||
if err != nil {
|
||||
writeSubmitError(w, r, err)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusOK, prog)
|
||||
}
|
||||
|
||||
// handleContextUploadComplete stores the caller's staged upload as the
|
||||
// submission's context. It holds the per-user upload cooldown and is audited
|
||||
// like the single upload, since this is where a context lands.
|
||||
func (a *API) handleContextUploadComplete(w http.ResponseWriter, r *http.Request) {
|
||||
if a.Submissions == nil {
|
||||
writeError(w, r, errSubmissionsUnavailable)
|
||||
return
|
||||
}
|
||||
p := principalFromContext(r.Context())
|
||||
commit, release, ok := a.reserveUpload(w, r, p.UserID)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
defer release()
|
||||
sub, err := a.Submissions.CompleteUpload(r.Context(), r.PathValue("id"), p.UserID)
|
||||
if err != nil {
|
||||
writeSubmitError(w, r, err)
|
||||
return
|
||||
}
|
||||
commit()
|
||||
a.audit(r, "submission.upload", sub.ID)
|
||||
writeJSON(w, http.StatusOK, sub)
|
||||
}
|
||||
@@ -431,6 +518,7 @@ var errSubmissionsUnavailable = newError(http.StatusServiceUnavailable, "submiss
|
||||
// server-side fault that collapses to 500 via writeError. The lane deliberately
|
||||
// does not surface those as 4xx: the client did nothing wrong.
|
||||
func writeSubmitError(w http.ResponseWriter, r *http.Request, err error) {
|
||||
var mismatch *submit.OffsetMismatchError
|
||||
switch {
|
||||
case errors.Is(err, submit.ErrInvalid):
|
||||
writeError(w, r, newError(http.StatusBadRequest, "bad_request", "%s", err.Error()))
|
||||
@@ -453,6 +541,15 @@ func writeSubmitError(w http.ResponseWriter, r *http.Request, err error) {
|
||||
case errors.Is(err, submit.ErrUploadsUnavailable):
|
||||
writeError(w, r, newError(http.StatusServiceUnavailable, "uploads_unavailable",
|
||||
"modpack upload transport is not configured"))
|
||||
case errors.Is(err, submit.ErrUploadBusy):
|
||||
writeError(w, r, newError(http.StatusConflict, "upload_busy",
|
||||
"another request is still writing this upload; read where it stands and continue from there"))
|
||||
case errors.As(err, &mismatch):
|
||||
writeError(w, r, newError(http.StatusConflict, "upload_offset_mismatch",
|
||||
"the upload holds %d bytes; send the part that starts there", mismatch.Received))
|
||||
case errors.Is(err, submit.ErrPartTooLarge):
|
||||
writeError(w, r, newError(http.StatusRequestEntityTooLarge, "part_too_large",
|
||||
"the part is larger than part_max_bytes in the upload status"))
|
||||
default:
|
||||
writeError(w, r, err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,136 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/submit"
|
||||
)
|
||||
|
||||
func TestContextUploadPartForwardsOffsetBodyAndPrincipal(t *testing.T) {
|
||||
fs := &fakeSubmissions{progress: submit.UploadProgress{Received: 8, PartMaxBytes: 33554432, MaxContextBytes: 1073741824}}
|
||||
api := appSubAPI(fs)
|
||||
w := do(api.ExternalHandler(), "PUT", "/api/v1/me/submissions/sub-9/context/upload?offset=4", "abcd",
|
||||
ctHeader("application/octet-stream"))
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("code = %d, want 200 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if got := w.Body.String(); got != `{"received":8,"part_max_bytes":33554432,"max_context_bytes":1073741824}`+"\n" {
|
||||
t.Fatalf("body = %s", got)
|
||||
}
|
||||
if fs.chunkID != "sub-9" || fs.chunkBy != "user-7" || fs.partOffset != 4 || fs.partBody != "abcd" {
|
||||
t.Fatalf("forwarded id=%q by=%q offset=%d body=%q, want sub-9 user-7 4 abcd", fs.chunkID, fs.chunkBy, fs.partOffset, fs.partBody)
|
||||
}
|
||||
}
|
||||
|
||||
func TestContextUploadPartNeedsAnOffset(t *testing.T) {
|
||||
for _, target := range []string{
|
||||
"/api/v1/me/submissions/sub-9/context/upload",
|
||||
"/api/v1/me/submissions/sub-9/context/upload?offset=four",
|
||||
} {
|
||||
fs := &fakeSubmissions{}
|
||||
w := do(appSubAPI(fs).ExternalHandler(), "PUT", target, "abcd", nil)
|
||||
if w.Code != http.StatusBadRequest || decodeErr(t, w) != "bad_request" {
|
||||
t.Fatalf("%s: code = %d (%s), want 400 bad_request", target, w.Code, w.Body.String())
|
||||
}
|
||||
if fs.chunkID != "" {
|
||||
t.Fatalf("%s: the part reached the lane", target)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestContextUploadStatusReportsTheStagedLength(t *testing.T) {
|
||||
fs := &fakeSubmissions{progress: submit.UploadProgress{Received: 50331648, PartMaxBytes: 33554432, MaxContextBytes: 1073741824}}
|
||||
w := do(appSubAPI(fs).ExternalHandler(), "GET", "/api/v1/me/submissions/sub-9/context/upload", "", nil)
|
||||
if w.Code != http.StatusOK || w.Body.String() != `{"received":50331648,"part_max_bytes":33554432,"max_context_bytes":1073741824}`+"\n" {
|
||||
t.Fatalf("status = %d %s", w.Code, w.Body.String())
|
||||
}
|
||||
if fs.chunkID != "sub-9" || fs.chunkBy != "user-7" {
|
||||
t.Fatalf("asked for id=%q by=%q, want sub-9 user-7", fs.chunkID, fs.chunkBy)
|
||||
}
|
||||
}
|
||||
|
||||
func TestContextUploadErrors(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
err error
|
||||
code int
|
||||
want string
|
||||
msg string
|
||||
}{
|
||||
{"offset mismatch", &submit.OffsetMismatchError{Received: 12}, 409, "upload_offset_mismatch", "the upload holds 12 bytes"},
|
||||
{"busy", submit.ErrUploadBusy, 409, "upload_busy", ""},
|
||||
{"part too large", fmt.Errorf("submit: write upload part: %w", submit.ErrPartTooLarge), 413, "part_too_large", ""},
|
||||
{"not owned", submit.ErrNotFound, 404, "not_found", ""},
|
||||
{"no part store", submit.ErrUploadsUnavailable, 503, "uploads_unavailable", ""},
|
||||
} {
|
||||
fs := &fakeSubmissions{chunkErr: tc.err}
|
||||
w := do(appSubAPI(fs).ExternalHandler(), "PUT", "/api/v1/me/submissions/sub-9/context/upload?offset=12", "abcd", nil)
|
||||
if w.Code != tc.code || decodeErr(t, w) != tc.want {
|
||||
t.Errorf("%s: %d %s, want %d %s", tc.name, w.Code, w.Body.String(), tc.code, tc.want)
|
||||
}
|
||||
if tc.msg != "" && !strings.Contains(w.Body.String(), tc.msg) {
|
||||
t.Errorf("%s: body %s does not say %q", tc.name, w.Body.String(), tc.msg)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Completion is where a context lands: it holds the per-user upload cooldown and
|
||||
// writes the same audit event as the single upload, and a failed completion
|
||||
// gives the cooldown back.
|
||||
func TestContextUploadCompleteHoldsTheCooldownAndAudits(t *testing.T) {
|
||||
repo := newFakeRepo()
|
||||
fs := &fakeSubmissions{}
|
||||
api := appSubAPI(fs)
|
||||
api.Repo = repo
|
||||
api.SubmitUploadCooldown = time.Minute
|
||||
eh := api.ExternalHandler()
|
||||
|
||||
fs.completeErr = submit.ErrUploadsUnavailable
|
||||
if w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context/upload/complete", "", nil); w.Code != http.StatusServiceUnavailable {
|
||||
t.Fatalf("failed completion: code = %d, want 503 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
fs.completeErr = nil
|
||||
w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context/upload/complete", "", nil)
|
||||
if w.Code != http.StatusOK {
|
||||
t.Fatalf("completion right after a failed one: code = %d, want 200 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
if fs.chunkID != "sub-9" || fs.chunkBy != "user-7" {
|
||||
t.Fatalf("completed id=%q by=%q, want sub-9 user-7", fs.chunkID, fs.chunkBy)
|
||||
}
|
||||
var audited int
|
||||
for _, e := range repo.audits {
|
||||
if e.Action == "submission.upload" {
|
||||
audited++
|
||||
}
|
||||
}
|
||||
if audited != 1 {
|
||||
t.Fatalf("submission.upload audits = %d, want 1 (only the completion that stored): %+v", audited, repo.audits)
|
||||
}
|
||||
w = do(eh, "POST", "/api/v1/me/submissions/sub-9/context/upload/complete", "", nil)
|
||||
if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "submission_cooldown" {
|
||||
t.Fatalf("second completion in the window: %d %s, want 429 submission_cooldown", w.Code, w.Body.String())
|
||||
}
|
||||
// A part never waits on the cooldown.
|
||||
if w := do(eh, "PUT", "/api/v1/me/submissions/sub-9/context/upload?offset=0", "\x1f\x8b", nil); w.Code != http.StatusOK {
|
||||
t.Fatalf("part inside the cooldown: code = %d, want 200 (%s)", w.Code, w.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestContextUploadWithoutServiceIs503(t *testing.T) {
|
||||
api := appSubAPI(&fakeSubmissions{})
|
||||
api.Submissions = nil
|
||||
eh := api.ExternalHandler()
|
||||
for _, rq := range [][2]string{
|
||||
{"GET", "/api/v1/me/submissions/sub-9/context/upload"},
|
||||
{"PUT", "/api/v1/me/submissions/sub-9/context/upload?offset=0"},
|
||||
{"POST", "/api/v1/me/submissions/sub-9/context/upload/complete"},
|
||||
} {
|
||||
if w := do(eh, rq[0], rq[1], "", nil); w.Code != http.StatusServiceUnavailable {
|
||||
t.Errorf("%s %s = %d, want 503", rq[0], rq[1], w.Code)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -55,6 +55,35 @@ type fakeSubmissions struct {
|
||||
|
||||
approvedDigest string
|
||||
openDigest string
|
||||
|
||||
// The chunked upload: what the handler forwarded, and canned outcomes.
|
||||
chunkID string
|
||||
chunkBy string
|
||||
partOffset int64
|
||||
partBody string
|
||||
progress submit.UploadProgress
|
||||
chunkErr error
|
||||
completeErr error
|
||||
}
|
||||
|
||||
func (f *fakeSubmissions) UploadStatus(_ context.Context, id, submittedBy string) (submit.UploadProgress, error) {
|
||||
f.chunkID, f.chunkBy = id, submittedBy
|
||||
return f.progress, f.chunkErr
|
||||
}
|
||||
|
||||
func (f *fakeSubmissions) UploadPart(_ context.Context, id, submittedBy string, offset int64, r io.Reader) (submit.UploadProgress, error) {
|
||||
f.chunkID, f.chunkBy, f.partOffset = id, submittedBy, offset
|
||||
b, _ := io.ReadAll(r)
|
||||
f.partBody = string(b)
|
||||
return f.progress, f.chunkErr
|
||||
}
|
||||
|
||||
func (f *fakeSubmissions) CompleteUpload(_ context.Context, id, submittedBy string) (*submit.Submission, error) {
|
||||
f.chunkID, f.chunkBy = id, submittedBy
|
||||
if f.completeErr != nil {
|
||||
return nil, f.completeErr
|
||||
}
|
||||
return &submit.Submission{ID: id, SubmittedBy: submittedBy, Status: submit.StatusPendingReview}, nil
|
||||
}
|
||||
|
||||
func (f *fakeSubmissions) Create(_ context.Context, req submit.CreateRequest) (*submit.Submission, error) {
|
||||
|
||||
@@ -170,13 +170,6 @@ func (a AuthConfig) EffectiveClientIPHeader() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
// BehindCloudflare reports whether requests reach the API through the
|
||||
// Cloudflare edge: an Access audience (set only by the edge setup) or
|
||||
// CF-Connecting-IP as the client address header.
|
||||
func (a AuthConfig) BehindCloudflare() bool {
|
||||
return strings.EqualFold(a.EffectiveClientIPHeader(), "CF-Connecting-IP")
|
||||
}
|
||||
|
||||
// K8sConfig is the [k8s] table.
|
||||
type K8sConfig struct {
|
||||
Namespace string `toml:"namespace"`
|
||||
@@ -253,10 +246,10 @@ type RegistryConfig struct {
|
||||
// own; this bounds the sum, which on k3s local-path is the only bound, since
|
||||
// the uploads PVC's size is not enforced there. Empty keeps 4Gi.
|
||||
UserUploadsMaxBytes string `toml:"user_uploads_max_bytes"`
|
||||
// ContextMaxBytes caps one uploaded build context, as a quantity ("95Mi").
|
||||
// Empty keeps 1Gi, except behind the Cloudflare edge (see BehindCloudflare),
|
||||
// whose proxy refuses request bodies over 100 MB before they reach the API;
|
||||
// there it keeps 95Mi, so the API's own 413 is what the uploader sees.
|
||||
// ContextMaxBytes caps one uploaded build context, as a quantity ("512Mi").
|
||||
// Empty keeps 1Gi. The Cloudflare edge refuses a single request body over
|
||||
// 100 MB; the panel sends a context in 32 MiB parts
|
||||
// (/api/v1/me/submissions/{id}/context/upload), so the cap holds behind it.
|
||||
ContextMaxBytes string `toml:"context_max_bytes"`
|
||||
// S3 configures the object-store backend for user_uploads_context when it is an
|
||||
// s3:// base (the alternative to a local uploads path). It mirrors
|
||||
|
||||
@@ -58,16 +58,21 @@ const DefaultUploadsMinFree = 0.10
|
||||
// CheckRoom refuses an upload of up to need bytes that could push Base's
|
||||
// filesystem below its free floor.
|
||||
func (s *LocalContextStore) CheckRoom(need int64) error {
|
||||
return checkRoom(s.Base, need, s.MinFree)
|
||||
}
|
||||
|
||||
// checkRoom refuses need more bytes under dir when they could push its filesystem
|
||||
// below the minFree share (DefaultUploadsMinFree when 0), wrapping ErrUploadsFull.
|
||||
func checkRoom(dir string, need int64, minFree float64) error {
|
||||
var st syscall.Statfs_t
|
||||
if err := syscall.Statfs(s.Base, &st); err != nil {
|
||||
return fmt.Errorf("submit: measure the uploads store %s: %w", s.Base, err)
|
||||
if err := syscall.Statfs(dir, &st); err != nil {
|
||||
return fmt.Errorf("submit: measure the uploads store %s: %w", dir, err)
|
||||
}
|
||||
bsize := uint64(st.Bsize) // uint32 on darwin
|
||||
total, avail := uint64(st.Blocks)*bsize, uint64(st.Bavail)*bsize
|
||||
if total == 0 {
|
||||
return nil
|
||||
}
|
||||
minFree := s.MinFree
|
||||
if minFree <= 0 {
|
||||
minFree = DefaultUploadsMinFree
|
||||
}
|
||||
|
||||
@@ -0,0 +1,238 @@
|
||||
package submit
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// DefaultPartMaxBytes caps one part of a chunked upload. A single request body is
|
||||
// bounded by whatever proxy fronts the API — the Cloudflare edge refuses bodies
|
||||
// over 100 MB on the Free and Pro plans — so a large modpack arrives as a run of
|
||||
// parts. 32 MiB stays well under that limit and keeps a retried part cheap.
|
||||
const DefaultPartMaxBytes = 32 << 20
|
||||
|
||||
// StalePartRetention is how long a staged upload may sit untouched before the
|
||||
// reaper deletes it: long enough to resume after a lost connection or a laptop
|
||||
// lid, short enough that abandoned parts do not hold the uploads volume.
|
||||
const StalePartRetention = 24 * time.Hour
|
||||
|
||||
// partSuffix names a staged upload on disk: {Dir}/{id}.part.
|
||||
const partSuffix = ".part"
|
||||
|
||||
// ErrUploadBusy reports that another request is writing (or completing) the same
|
||||
// staged upload. Parts go in order, one at a time; the API answers 409 and the
|
||||
// client asks where the upload stands before it sends again.
|
||||
var ErrUploadBusy = errors.New("submit: another request is writing this upload")
|
||||
|
||||
// ErrPartTooLarge reports a part longer than the part cap. The API answers 413.
|
||||
var ErrPartTooLarge = errors.New("submit: an upload part exceeds the part size limit")
|
||||
|
||||
// OffsetMismatchError reports a part that does not start where the staged upload
|
||||
// ends. Received is where it does end, so the client resumes from there.
|
||||
type OffsetMismatchError struct {
|
||||
Received int64
|
||||
}
|
||||
|
||||
func (e *OffsetMismatchError) Error() string {
|
||||
return fmt.Sprintf("submit: the upload holds %d bytes; send the part that starts there", e.Received)
|
||||
}
|
||||
|
||||
// PartStore stages a chunked upload until Manager.CompleteUpload hands the
|
||||
// assembled bytes to Blobs. Parts are appended in order to {Dir}/{id}.part, and
|
||||
// the file's length is the resume point: a client that lost its connection asks
|
||||
// for it and carries on from there. The store serializes the requests for one id
|
||||
// in process, which is enough for the single felis-api replica (its Deployment is
|
||||
// Recreate, never two pods at once).
|
||||
type PartStore struct {
|
||||
// Dir holds the staged uploads. cmd/felis puts it on the uploads volume, so a
|
||||
// staged upload survives an API restart and the room check sees the same disk.
|
||||
Dir string
|
||||
// MinFree is the share of Dir's filesystem a part must leave free; 0 uses
|
||||
// DefaultUploadsMinFree.
|
||||
MinFree float64
|
||||
|
||||
mu sync.Mutex
|
||||
busy map[string]bool
|
||||
}
|
||||
|
||||
// hold claims id for one request; the returned func releases it. A second claim
|
||||
// while the first is held fails with ErrUploadBusy.
|
||||
func (s *PartStore) hold(id string) (func(), error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.busy[id] {
|
||||
return nil, ErrUploadBusy
|
||||
}
|
||||
if s.busy == nil {
|
||||
s.busy = map[string]bool{}
|
||||
}
|
||||
s.busy[id] = true
|
||||
return func() {
|
||||
s.mu.Lock()
|
||||
delete(s.busy, id)
|
||||
s.mu.Unlock()
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *PartStore) path(id string) (string, error) {
|
||||
if !idRE.MatchString(id) {
|
||||
return "", fmt.Errorf("submit: invalid submission id %q", id)
|
||||
}
|
||||
return filepath.Join(s.Dir, id+partSuffix), nil
|
||||
}
|
||||
|
||||
// CheckRoom refuses a part of up to need bytes that could push Dir's filesystem
|
||||
// below its free floor.
|
||||
func (s *PartStore) CheckRoom(need int64) error {
|
||||
if err := os.MkdirAll(s.Dir, 0o750); err != nil {
|
||||
return fmt.Errorf("submit: mkdir upload parts dir: %w", err)
|
||||
}
|
||||
return checkRoom(s.Dir, need, s.MinFree)
|
||||
}
|
||||
|
||||
// Size reports how many bytes are staged for id; nothing staged is (0, false, nil).
|
||||
func (s *PartStore) Size(id string) (int64, bool, error) {
|
||||
p, err := s.path(id)
|
||||
if err != nil {
|
||||
return 0, false, err
|
||||
}
|
||||
fi, err := os.Stat(p)
|
||||
switch {
|
||||
case err == nil:
|
||||
return fi.Size(), true, nil
|
||||
case os.IsNotExist(err):
|
||||
return 0, false, nil
|
||||
default:
|
||||
return 0, false, fmt.Errorf("submit: stat staged upload: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Append writes r to id's staged upload at offset and returns the new length.
|
||||
// Offset 0 starts the upload over; any other offset must equal the staged length,
|
||||
// or the call fails with *OffsetMismatchError naming it. A part that fails to
|
||||
// arrive whole is cut back off, so the staged bytes are always a prefix of what
|
||||
// the client sent.
|
||||
func (s *PartStore) Append(id string, offset int64, r io.Reader) (int64, error) {
|
||||
p, err := s.path(id)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
release, err := s.hold(id)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer release()
|
||||
if err := os.MkdirAll(s.Dir, 0o750); err != nil {
|
||||
return 0, fmt.Errorf("submit: mkdir upload parts dir: %w", err)
|
||||
}
|
||||
var f *os.File
|
||||
if offset == 0 {
|
||||
f, err = os.OpenFile(p, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o640)
|
||||
} else {
|
||||
f, err = os.OpenFile(p, os.O_WRONLY, 0)
|
||||
if os.IsNotExist(err) {
|
||||
return 0, &OffsetMismatchError{Received: 0}
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("submit: open staged upload: %w", err)
|
||||
}
|
||||
fi, err := f.Stat()
|
||||
if err != nil {
|
||||
f.Close()
|
||||
return 0, fmt.Errorf("submit: stat staged upload: %w", err)
|
||||
}
|
||||
if fi.Size() != offset {
|
||||
f.Close()
|
||||
return 0, &OffsetMismatchError{Received: fi.Size()}
|
||||
}
|
||||
if _, err := f.Seek(offset, io.SeekStart); err != nil {
|
||||
f.Close()
|
||||
return 0, fmt.Errorf("submit: seek staged upload: %w", err)
|
||||
}
|
||||
n, err := io.Copy(f, r)
|
||||
if err == nil {
|
||||
err = f.Sync()
|
||||
}
|
||||
if err != nil {
|
||||
f.Truncate(offset)
|
||||
f.Close()
|
||||
return 0, fmt.Errorf("submit: write upload part: %w", err)
|
||||
}
|
||||
if err := f.Close(); err != nil {
|
||||
return 0, fmt.Errorf("submit: close staged upload: %w", err)
|
||||
}
|
||||
return offset + n, nil
|
||||
}
|
||||
|
||||
// open returns id's staged upload for reading, or an ErrInvalid error when
|
||||
// nothing is staged. The caller must already hold id.
|
||||
func (s *PartStore) open(id string) (*os.File, error) {
|
||||
p, err := s.path(id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f, err := os.Open(p)
|
||||
if os.IsNotExist(err) {
|
||||
return nil, invalidf("no upload in progress for this submission; send its parts first")
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("submit: open staged upload: %w", err)
|
||||
}
|
||||
return f, nil
|
||||
}
|
||||
|
||||
// Delete removes id's staged upload. Nothing staged is success.
|
||||
func (s *PartStore) Delete(id string) error {
|
||||
p, err := s.path(id)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.Remove(p); err != nil && !os.IsNotExist(err) {
|
||||
return fmt.Errorf("submit: remove staged upload: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Reap deletes every staged upload last written before cutoff, skipping one a
|
||||
// request holds right now. It returns how many it deleted and carries on past
|
||||
// one it cannot delete, reporting the first such error.
|
||||
func (s *PartStore) Reap(cutoff time.Time) (int, error) {
|
||||
entries, err := os.ReadDir(s.Dir)
|
||||
if os.IsNotExist(err) {
|
||||
return 0, nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("submit: list staged uploads: %w", err)
|
||||
}
|
||||
var reaped int
|
||||
var firstErr error
|
||||
for _, e := range entries {
|
||||
id, ok := strings.CutSuffix(e.Name(), partSuffix)
|
||||
if !ok || e.IsDir() || !idRE.MatchString(id) {
|
||||
continue
|
||||
}
|
||||
fi, err := e.Info()
|
||||
if err != nil || !fi.ModTime().Before(cutoff) {
|
||||
continue
|
||||
}
|
||||
release, err := s.hold(id)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
err = s.Delete(id)
|
||||
release()
|
||||
if err == nil {
|
||||
reaped++
|
||||
} else if firstErr == nil {
|
||||
firstErr = err
|
||||
}
|
||||
}
|
||||
return reaped, firstErr
|
||||
}
|
||||
@@ -0,0 +1,338 @@
|
||||
package submit
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// newChunkedManager wires a Manager with an in-memory Blobs, a PartStore on a
|
||||
// temp dir and 4-byte parts, and files one pending submission for user-1.
|
||||
func newChunkedManager(t *testing.T) (*Manager, *fakeStore, *fakeBlobs, string) {
|
||||
t.Helper()
|
||||
m, st, _ := newManager()
|
||||
fb := newFakeBlobs()
|
||||
m.Blobs = fb
|
||||
// A near-zero floor keeps the room check off this machine's own disk usage.
|
||||
m.Parts = &PartStore{Dir: t.TempDir(), MinFree: 1e-9}
|
||||
m.PartMaxBytes = 4
|
||||
sub, err := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return m, st, fb, sub.ID
|
||||
}
|
||||
|
||||
func sendPart(t *testing.T, m *Manager, id string, offset int64, part string) int64 {
|
||||
t.Helper()
|
||||
p, err := m.UploadPart(context.Background(), id, "user-1", offset, strings.NewReader(part))
|
||||
if err != nil {
|
||||
t.Fatalf("UploadPart(offset %d, %q): %v", offset, part, err)
|
||||
}
|
||||
return p.Received
|
||||
}
|
||||
|
||||
func staged(t *testing.T, m *Manager, id string) int64 {
|
||||
t.Helper()
|
||||
p, err := m.UploadStatus(context.Background(), id, "user-1")
|
||||
if err != nil {
|
||||
t.Fatalf("UploadStatus: %v", err)
|
||||
}
|
||||
return p.Received
|
||||
}
|
||||
|
||||
func TestChunkedUploadAssemblesAndStores(t *testing.T) {
|
||||
m, _, fb, id := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
|
||||
p, err := m.UploadStatus(ctx, id, "user-1")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if p != (UploadProgress{Received: 0, PartMaxBytes: 4, MaxContextBytes: 1 << 30}) {
|
||||
t.Fatalf("status before any part = %+v", p)
|
||||
}
|
||||
for _, step := range []struct {
|
||||
offset int64
|
||||
part string
|
||||
want int64
|
||||
}{
|
||||
{0, "\x1f\x8b\x08\x00", 4},
|
||||
{4, "abcd", 8},
|
||||
{8, "efgh", 12},
|
||||
{12, "ij", 14},
|
||||
} {
|
||||
if got := sendPart(t, m, id, step.offset, step.part); got != step.want {
|
||||
t.Fatalf("after the part at %d: received %d, want %d", step.offset, got, step.want)
|
||||
}
|
||||
}
|
||||
if _, ok := fb.stored[id]; ok {
|
||||
t.Fatal("parts reached the blob store before the upload completed")
|
||||
}
|
||||
|
||||
sub, err := m.CompleteUpload(ctx, id, "user-1")
|
||||
if err != nil {
|
||||
t.Fatalf("CompleteUpload: %v", err)
|
||||
}
|
||||
if got := string(fb.stored[id]); got != "\x1f\x8b\x08\x00abcdefghij" {
|
||||
t.Fatalf("stored %q, want the parts in order", got)
|
||||
}
|
||||
if sub.ContextSHA256 != "df86a051af3d8615de097952ac01bc094781053bcda61aaac848edb98068209d" {
|
||||
t.Fatalf("recorded digest %s", sub.ContextSHA256)
|
||||
}
|
||||
if got := staged(t, m, id); got != 0 {
|
||||
t.Fatalf("staged bytes after completion = %d, want the staged copy gone", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadResumesWhereTheStagedBytesEnd(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
|
||||
var mismatch *OffsetMismatchError
|
||||
_, err := m.UploadPart(ctx, id, "user-1", 4, strings.NewReader("abcd"))
|
||||
if !errors.As(err, &mismatch) || mismatch.Received != 0 {
|
||||
t.Fatalf("a part past an empty upload = %v, want an offset mismatch at 0", err)
|
||||
}
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
_, err = m.UploadPart(ctx, id, "user-1", 8, strings.NewReader("efgh"))
|
||||
if !errors.As(err, &mismatch) || mismatch.Received != 4 {
|
||||
t.Fatalf("a part that skips ahead = %v, want an offset mismatch at 4", err)
|
||||
}
|
||||
if got := sendPart(t, m, id, 4, "abcd"); got != 8 {
|
||||
t.Fatalf("the part at the staged length: received %d, want 8", got)
|
||||
}
|
||||
// Offset 0 starts over.
|
||||
if got := sendPart(t, m, id, 0, "\x1f\x8b"); got != 2 {
|
||||
t.Fatalf("restart at 0: received %d, want 2", got)
|
||||
}
|
||||
}
|
||||
|
||||
type failingReader struct {
|
||||
data string
|
||||
done bool
|
||||
}
|
||||
|
||||
func (f *failingReader) Read(p []byte) (int, error) {
|
||||
if f.done {
|
||||
return 0, errors.New("connection reset")
|
||||
}
|
||||
f.done = true
|
||||
return copy(p, f.data), nil
|
||||
}
|
||||
|
||||
func TestChunkedUploadCutsBackAPartThatBrokeOff(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
|
||||
_, err := m.UploadPart(context.Background(), id, "user-1", 4, &failingReader{data: "ab"})
|
||||
if err == nil {
|
||||
t.Fatal("a part that broke off was accepted")
|
||||
}
|
||||
if got := staged(t, m, id); got != 4 {
|
||||
t.Fatalf("staged after a broken part = %d, want 4 (the half part cut back off)", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadFirstPartMustBeGzip(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
_, err := m.UploadPart(context.Background(), id, "user-1", 0, strings.NewReader("PK\x03\x04"))
|
||||
if !errors.Is(err, ErrInvalid) {
|
||||
t.Fatalf("a zip first part = %v, want ErrInvalid", err)
|
||||
}
|
||||
if got := staged(t, m, id); got != 0 {
|
||||
t.Fatalf("staged after a refused first part = %d, want 0", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadPartCap(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
|
||||
_, err := m.UploadPart(context.Background(), id, "user-1", 4, strings.NewReader("abcde"))
|
||||
if !errors.Is(err, ErrPartTooLarge) {
|
||||
t.Fatalf("a 5-byte part over a 4-byte cap = %v, want ErrPartTooLarge", err)
|
||||
}
|
||||
if got := staged(t, m, id); got != 4 {
|
||||
t.Fatalf("staged after an oversize part = %d, want 4", got)
|
||||
}
|
||||
}
|
||||
|
||||
// The staged total meets the same cap as a whole context, with the same error.
|
||||
func TestChunkedUploadTotalCappedLikeAWholeContext(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
m.MaxContextBytes = 10
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
sendPart(t, m, id, 4, "abcd")
|
||||
|
||||
_, err := m.UploadPart(context.Background(), id, "user-1", 8, strings.NewReader("efg"))
|
||||
if !errors.Is(err, ErrInvalid) || errors.Is(err, ErrPartTooLarge) {
|
||||
t.Fatalf("a part past the context cap = %v, want the context-too-large ErrInvalid", err)
|
||||
}
|
||||
if got := sendPart(t, m, id, 8, "ef"); got != 10 {
|
||||
t.Fatalf("a part ending exactly at the cap: received %d, want 10", got)
|
||||
}
|
||||
}
|
||||
|
||||
// Staged bytes count against the storage budget, so parts parked on several
|
||||
// pending submissions cannot hold more than stored contexts could.
|
||||
func TestChunkedUploadStagedBytesCountAgainstTheBudget(t *testing.T) {
|
||||
m, _, fb, a := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
m.MaxStoredBytesPerUser = 10
|
||||
m.PartMaxBytes = 16
|
||||
sendPart(t, m, a, 0, "\x1f\x8b\x08\x00abcd")
|
||||
|
||||
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")); !errors.Is(err, ErrQuotaExceeded) {
|
||||
t.Fatalf("3 bytes beside 8 staged under a 10-byte budget = %v, want ErrQuotaExceeded", err)
|
||||
}
|
||||
if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader("\x1f\x8b")); err != nil {
|
||||
t.Fatalf("2 bytes beside 8 staged = %v, want accepted", err)
|
||||
}
|
||||
// And the other way round: B's stored 2 bytes bound what A may still stage.
|
||||
if _, err := m.UploadPart(ctx, a, "user-1", 8, strings.NewReader("e")); !errors.Is(err, ErrQuotaExceeded) {
|
||||
t.Fatalf("staging past the budget = %v, want ErrQuotaExceeded", err)
|
||||
}
|
||||
if _, ok := fb.stored[a]; ok {
|
||||
t.Fatal("staging stored a blob")
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadOwnerAndStatus(t *testing.T) {
|
||||
m, st, _, id := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
|
||||
if _, err := m.UploadStatus(ctx, id, "user-2"); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("another user's status = %v, want ErrNotFound", err)
|
||||
}
|
||||
if _, err := m.UploadPart(ctx, id, "user-2", 4, strings.NewReader("abcd")); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("another user's part = %v, want ErrNotFound", err)
|
||||
}
|
||||
if _, err := m.CompleteUpload(ctx, id, "user-2"); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("another user's completion = %v, want ErrNotFound", err)
|
||||
}
|
||||
st.subs[id].Status = StatusRejected
|
||||
if _, err := m.UploadPart(ctx, id, "user-1", 4, strings.NewReader("abcd")); !errors.Is(err, ErrAlreadyReviewed) {
|
||||
t.Fatalf("a part for a reviewed submission = %v, want ErrAlreadyReviewed", err)
|
||||
}
|
||||
if _, err := m.CompleteUpload(ctx, id, "user-1"); !errors.Is(err, ErrAlreadyReviewed) {
|
||||
t.Fatalf("completing a reviewed submission = %v, want ErrAlreadyReviewed", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadWithoutPartStoreIsUnavailable(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
m.Parts = nil
|
||||
if _, err := m.UploadStatus(context.Background(), id, "user-1"); !errors.Is(err, ErrUploadsUnavailable) {
|
||||
t.Fatalf("status without a part store = %v, want ErrUploadsUnavailable", err)
|
||||
}
|
||||
if _, err := m.UploadPart(context.Background(), id, "user-1", 0, strings.NewReader("\x1f\x8b")); !errors.Is(err, ErrUploadsUnavailable) {
|
||||
t.Fatalf("part without a part store = %v, want ErrUploadsUnavailable", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCompleteUploadWithNothingStagedIsInvalid(t *testing.T) {
|
||||
m, _, fb, id := newChunkedManager(t)
|
||||
if _, err := m.CompleteUpload(context.Background(), id, "user-1"); !errors.Is(err, ErrInvalid) {
|
||||
t.Fatalf("completing an empty upload = %v, want ErrInvalid", err)
|
||||
}
|
||||
if len(fb.stored) != 0 {
|
||||
t.Fatal("an empty completion stored a blob")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCompleteUploadFailureKeepsTheStagedBytes(t *testing.T) {
|
||||
m, _, fb, id := newChunkedManager(t)
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
sendPart(t, m, id, 4, "ab")
|
||||
fb.putErr = errors.New("object store unreachable")
|
||||
|
||||
if _, err := m.CompleteUpload(context.Background(), id, "user-1"); err == nil {
|
||||
t.Fatal("completion succeeded with the blob store down")
|
||||
}
|
||||
if got := staged(t, m, id); got != 6 {
|
||||
t.Fatalf("staged after a failed completion = %d, want 6 kept for the retry", got)
|
||||
}
|
||||
fb.putErr = nil
|
||||
if _, err := m.CompleteUpload(context.Background(), id, "user-1"); err != nil {
|
||||
t.Fatalf("retried completion: %v", err)
|
||||
}
|
||||
if got := string(fb.stored[id]); got != "\x1f\x8b\x08\x00ab" {
|
||||
t.Fatalf("stored %q after the retry", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestChunkedUploadOneRequestAtATime(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
release, err := m.Parts.hold(id)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := m.UploadPart(context.Background(), id, "user-1", 4, strings.NewReader("ab")); !errors.Is(err, ErrUploadBusy) {
|
||||
t.Fatalf("a part while another is writing = %v, want ErrUploadBusy", err)
|
||||
}
|
||||
if _, err := m.CompleteUpload(context.Background(), id, "user-1"); !errors.Is(err, ErrUploadBusy) {
|
||||
t.Fatalf("completing while a part is writing = %v, want ErrUploadBusy", err)
|
||||
}
|
||||
release()
|
||||
if got := sendPart(t, m, id, 4, "ab"); got != 6 {
|
||||
t.Fatalf("the part after release: received %d, want 6", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithdrawDeletesTheStagedUpload(t *testing.T) {
|
||||
m, _, _, id := newChunkedManager(t)
|
||||
sendPart(t, m, id, 0, "\x1f\x8b\x08\x00")
|
||||
if _, err := m.Withdraw(context.Background(), id, "user-1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(m.Parts.Dir, id+".part")); !os.IsNotExist(err) {
|
||||
t.Fatalf("staged upload after withdraw: stat err = %v, want it gone", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReapStaleParts(t *testing.T) {
|
||||
m, _, _, oldID := newChunkedManager(t)
|
||||
ctx := context.Background()
|
||||
fresh, err := m.Create(ctx, CreateRequest{DisplayName: "Fresh", SubmittedBy: "user-1"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sendPart(t, m, oldID, 0, "\x1f\x8b\x08\x00")
|
||||
sendPart(t, m, fresh.ID, 0, "\x1f\x8b\x08\x00")
|
||||
if err := os.Chtimes(filepath.Join(m.Parts.Dir, oldID+".part"), testNow.Add(-25*time.Hour), testNow.Add(-25*time.Hour)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Chtimes(filepath.Join(m.Parts.Dir, fresh.ID+".part"), testNow.Add(-23*time.Hour), testNow.Add(-23*time.Hour)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// A stray file that is not a staged upload is left alone.
|
||||
if err := os.WriteFile(filepath.Join(m.Parts.Dir, "notes"), []byte("x"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
os.Chtimes(filepath.Join(m.Parts.Dir, "notes"), testNow.Add(-48*time.Hour), testNow.Add(-48*time.Hour))
|
||||
|
||||
n, err := m.ReapStaleParts(StalePartRetention)
|
||||
if err != nil || n != 1 {
|
||||
t.Fatalf("ReapStaleParts = %d, %v; want 1", n, err)
|
||||
}
|
||||
if got := staged(t, m, oldID); got != 0 {
|
||||
t.Fatalf("the 25-hour-old upload still holds %d bytes", got)
|
||||
}
|
||||
if got := staged(t, m, fresh.ID); got != 4 {
|
||||
t.Fatalf("the 23-hour-old upload holds %d bytes, want 4 kept", got)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(m.Parts.Dir, "notes")); err != nil {
|
||||
t.Fatalf("the reaper touched a file that is not a staged upload: %v", err)
|
||||
}
|
||||
}
|
||||
+192
-38
@@ -64,6 +64,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -379,6 +380,13 @@ type Manager struct {
|
||||
// MaxStoredBytesTotal overrides the budget for every user's stored contexts
|
||||
// together; 0 uses defaultMaxStoredBytesTotal.
|
||||
MaxStoredBytesTotal int64
|
||||
// Parts stages chunked uploads (UploadPart, CompleteUpload). Nil leaves only
|
||||
// the single-request UploadContext, and the chunked calls return
|
||||
// ErrUploadsUnavailable.
|
||||
Parts *PartStore
|
||||
// PartMaxBytes overrides the cap on one chunked-upload part; 0 uses
|
||||
// DefaultPartMaxBytes.
|
||||
PartMaxBytes int64
|
||||
|
||||
Now func() time.Time
|
||||
IDGen func() string
|
||||
@@ -394,6 +402,13 @@ func (m *Manager) maxContextBytes() int64 {
|
||||
return defaultMaxContextBytes
|
||||
}
|
||||
|
||||
func (m *Manager) partMaxBytes() int64 {
|
||||
if m.PartMaxBytes > 0 {
|
||||
return m.PartMaxBytes
|
||||
}
|
||||
return DefaultPartMaxBytes
|
||||
}
|
||||
|
||||
func (m *Manager) maxPendingPerUser() int {
|
||||
if m.MaxPendingPerUser > 0 {
|
||||
return m.MaxPendingPerUser
|
||||
@@ -575,18 +590,10 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i
|
||||
if m.Blobs == nil {
|
||||
return nil, ErrUploadsUnavailable
|
||||
}
|
||||
|
||||
sub, err := m.Store.GetSubmission(ctx, id)
|
||||
sub, err := m.ownPending(ctx, id, submittedBy)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if sub.SubmittedBy != submittedBy {
|
||||
// Not the owner: invisible, so the endpoint cannot confirm the id exists.
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
if sub.Status != StatusPendingReview {
|
||||
return nil, ErrAlreadyReviewed
|
||||
}
|
||||
|
||||
// Sniff the gzip magic before touching the store so a wrong-format upload fails
|
||||
// fast, without persisting anything or reading the whole body.
|
||||
@@ -595,34 +602,10 @@ 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)")
|
||||
}
|
||||
|
||||
// Per-user storage budget: sum the bytes this user's OTHER submissions
|
||||
// already hold (excluding this id, whose blob a re-upload supersedes) and cap
|
||||
// the write at whatever remains. cappedReader trips on the first byte past
|
||||
// the limit, so the store never persists a blob that would exceed the budget
|
||||
// (it removes its temp file on the copy error) and the failure surfaces as a
|
||||
// 403, not a 500. The read-then-write pair is not atomic in this 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.
|
||||
used, total, err := m.storedBytes(ctx, submittedBy, id)
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
remaining, budgetErr := m.maxStoredBytesPerUser()-used, errStorageQuota
|
||||
if left := m.maxStoredBytesTotal() - total; left < remaining {
|
||||
remaining, budgetErr = left, fmt.Errorf("%w: every user's uploads together reached the %d-byte limit", ErrUploadsFull, m.maxStoredBytesTotal())
|
||||
}
|
||||
if remaining <= 0 {
|
||||
return nil, budgetErr
|
||||
}
|
||||
limit, over := m.maxContextBytes(), errContextTooLarge
|
||||
if remaining < limit {
|
||||
// The budget binds before the single-blob cap: an upload tripping here is
|
||||
// refused as a spent allowance (or a full store), never as a malformed
|
||||
// request.
|
||||
limit, over = remaining, budgetErr
|
||||
}
|
||||
// The upload's size is unknown until it ends, so the room check assumes the
|
||||
// most it may write.
|
||||
if rc, ok := m.Blobs.(RoomChecker); ok {
|
||||
@@ -648,12 +631,174 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i
|
||||
return sub, nil
|
||||
}
|
||||
|
||||
// ownPending loads id for an upload by submittedBy: another user's submission is
|
||||
// ErrNotFound (the endpoint cannot confirm the id exists), a reviewed one
|
||||
// ErrAlreadyReviewed.
|
||||
func (m *Manager) ownPending(ctx context.Context, id, submittedBy string) (*Submission, error) {
|
||||
sub, err := m.Store.GetSubmission(ctx, id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if sub.SubmittedBy != submittedBy {
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
if sub.Status != StatusPendingReview {
|
||||
return nil, ErrAlreadyReviewed
|
||||
}
|
||||
return sub, nil
|
||||
}
|
||||
|
||||
// uploadLimit is how many bytes id's context may hold, and the error an upload
|
||||
// past it fails with. Per-user storage budget: sum the bytes this user's OTHER
|
||||
// submissions already hold (excluding this id, whose blob a re-upload
|
||||
// supersedes) and cap the write at whatever remains. cappedReader trips on the
|
||||
// first byte past the limit, so the store never persists a blob that would
|
||||
// exceed the budget (it removes its temp file on the copy error) and the failure
|
||||
// surfaces as a 403, not a 500. The read-then-write pair is not atomic in this
|
||||
// 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)
|
||||
if err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
remaining, budgetErr := m.maxStoredBytesPerUser()-used, errStorageQuota
|
||||
if left := m.maxStoredBytesTotal() - total; left < remaining {
|
||||
remaining, budgetErr = left, fmt.Errorf("%w: every user's uploads together reached the %d-byte limit", ErrUploadsFull, m.maxStoredBytesTotal())
|
||||
}
|
||||
if remaining <= 0 {
|
||||
return 0, nil, budgetErr
|
||||
}
|
||||
limit, over := m.maxContextBytes(), errContextTooLarge
|
||||
if remaining < limit {
|
||||
// The budget binds before the single-blob cap: an upload tripping here is
|
||||
// refused as a spent allowance (or a full store), never as a malformed
|
||||
// request.
|
||||
limit, over = remaining, budgetErr
|
||||
}
|
||||
return limit, over, nil
|
||||
}
|
||||
|
||||
// UploadProgress is where a chunked upload stands: Received bytes are staged,
|
||||
// the next part starts there and carries at most PartMaxBytes, and the whole
|
||||
// context may reach MaxContextBytes.
|
||||
type UploadProgress struct {
|
||||
Received int64 `json:"received"`
|
||||
PartMaxBytes int64 `json:"part_max_bytes"`
|
||||
MaxContextBytes int64 `json:"max_context_bytes"`
|
||||
}
|
||||
|
||||
func (m *Manager) chunked(ctx context.Context, id, submittedBy string) error {
|
||||
if strings.TrimSpace(submittedBy) == "" {
|
||||
return invalidf("submitter identity is required")
|
||||
}
|
||||
if m.Blobs == nil || m.Parts == nil {
|
||||
return ErrUploadsUnavailable
|
||||
}
|
||||
_, err := m.ownPending(ctx, id, submittedBy)
|
||||
return err
|
||||
}
|
||||
|
||||
// UploadStatus reports how far the caller's chunked upload of id has come, so a
|
||||
// client that lost its connection resumes where the staged bytes end. Nothing
|
||||
// staged reads as Received 0.
|
||||
func (m *Manager) UploadStatus(ctx context.Context, id, submittedBy string) (UploadProgress, error) {
|
||||
if err := m.chunked(ctx, id, submittedBy); err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
n, _, err := m.Parts.Size(id)
|
||||
if err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
return UploadProgress{Received: n, PartMaxBytes: m.partMaxBytes(), MaxContextBytes: m.maxContextBytes()}, nil
|
||||
}
|
||||
|
||||
// UploadPart appends one part of the caller's chunked upload of id, starting at
|
||||
// offset: 0 starts over, anything else must equal what is staged
|
||||
// (*OffsetMismatchError otherwise). The first part must open with the gzip
|
||||
// magic. A part is capped at PartMaxBytes (ErrPartTooLarge), and the staged
|
||||
// total at the same limit UploadContext applies to a whole context, with the
|
||||
// same errors, so a chunked upload cannot stage more than a single request could
|
||||
// store. Nothing reaches Blobs until CompleteUpload.
|
||||
func (m *Manager) UploadPart(ctx context.Context, id, submittedBy string, offset int64, r io.Reader) (UploadProgress, error) {
|
||||
if err := m.chunked(ctx, id, submittedBy); err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
if offset < 0 {
|
||||
return UploadProgress{}, invalidf("offset must not be negative")
|
||||
}
|
||||
br := bufio.NewReader(r)
|
||||
if offset == 0 {
|
||||
if magic, err := br.Peek(2); err != nil || magic[0] != 0x1f || magic[1] != 0x8b {
|
||||
return UploadProgress{}, invalidf("build context must be a gzip-compressed tarball (.tar.gz)")
|
||||
}
|
||||
}
|
||||
limit, over, err := m.uploadLimit(ctx, submittedBy, id)
|
||||
if err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
left, partOver := m.partMaxBytes(), error(ErrPartTooLarge)
|
||||
if room := limit - offset; room <= left {
|
||||
left, partOver = max(room, 0), over
|
||||
}
|
||||
if err := m.Parts.CheckRoom(left); err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
n, err := m.Parts.Append(id, offset, &cappedReader{r: br, left: left, over: partOver})
|
||||
if err != nil {
|
||||
return UploadProgress{}, err
|
||||
}
|
||||
return UploadProgress{Received: n, PartMaxBytes: m.partMaxBytes(), MaxContextBytes: m.maxContextBytes()}, nil
|
||||
}
|
||||
|
||||
// CompleteUpload stores the caller's staged upload of id as its build context,
|
||||
// through UploadContext, so it meets every check a single-request upload does
|
||||
// (format, size, budget, room) and records the digest the same way. The staged
|
||||
// bytes are deleted once stored; after a failure they stay, for a retry or a
|
||||
// fresh start at offset 0, until the reaper takes them.
|
||||
func (m *Manager) CompleteUpload(ctx context.Context, id, submittedBy string) (*Submission, error) {
|
||||
if err := m.chunked(ctx, id, submittedBy); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
release, err := m.Parts.hold(id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer release()
|
||||
f, err := m.Parts.open(id)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sub, err := m.UploadContext(ctx, id, submittedBy, f)
|
||||
f.Close()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := m.Parts.Delete(id); err != nil {
|
||||
// Stored and recorded; the leftover copy only waits for the reaper.
|
||||
log.Printf("submit: %v", err)
|
||||
}
|
||||
return sub, nil
|
||||
}
|
||||
|
||||
// ReapStaleParts deletes every staged upload untouched for longer than olderThan.
|
||||
func (m *Manager) ReapStaleParts(olderThan time.Duration) (int, error) {
|
||||
if m.Parts == nil {
|
||||
return 0, nil
|
||||
}
|
||||
return m.Parts.Reap(m.now().Add(-olderThan))
|
||||
}
|
||||
|
||||
// storedBytes sums the stored-blob sizes of submittedBy's submissions (user) and
|
||||
// of everyone's (total), excluding excludeID — the submission a pending re-upload
|
||||
// is about to replace, whose bytes must not be counted twice. Sizes are read from
|
||||
// the blob store itself, the same source of truth uploads/approval consult, so
|
||||
// the sums cannot drift from what is actually occupying the volume (including
|
||||
// blobs uploaded before any budget existed).
|
||||
// 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) {
|
||||
subs, err := m.Store.ListSubmissions(ctx)
|
||||
if err != nil {
|
||||
@@ -663,12 +808,16 @@ func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string
|
||||
if s.ID == excludeID {
|
||||
continue
|
||||
}
|
||||
n, ok, err := m.Blobs.Size(ctx, s.ID)
|
||||
n, _, err := m.Blobs.Size(ctx, s.ID)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
if !ok {
|
||||
continue
|
||||
if m.Parts != nil {
|
||||
staged, _, err := m.Parts.Size(s.ID)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
n += staged
|
||||
}
|
||||
total += n
|
||||
if s.SubmittedBy == submittedBy {
|
||||
@@ -976,6 +1125,11 @@ func (m *Manager) Delete(ctx context.Context, id string) (*Submission, error) {
|
||||
// exactly what is left behind. A deployment with no upload transport (Blobs nil)
|
||||
// has no blobs to reap.
|
||||
func (m *Manager) deleteBlob(ctx context.Context, id string) error {
|
||||
if m.Parts != nil {
|
||||
if err := m.Parts.Delete(id); err != nil {
|
||||
return fmt.Errorf("submit: submission removed, but its unfinished upload could not be deleted: %w", err)
|
||||
}
|
||||
}
|
||||
if m.Blobs == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user