From ad4d256d8f947fc1d03902160115f959bf97a621 Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Thu, 24 Sep 2026 10:14:42 +0800 Subject: [PATCH] =?UTF-8?q?fix(submit):=20bound=20the=20untrusted=20upload?= =?UTF-8?q?=20lane=20=E2=80=94=20per-user=20caps=20+=20throttles=20(#75)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A logged-in user could file submissions without bound and stream a 1 GiB context per submission. The only limits were the single-blob size cap and the 5 GiB uploads PVC (platform/workloads.go); nothing counted a user's rows or bytes, so one account could fill the volume and every other user's upload would start failing. - Create: per-user pending_review cap (default 5) — the review queue cannot be parked full of one account's rows. Check-then-insert, documented soft. - UploadContext: per-user stored-context budget (default 2 GiB) charged against the blob store's REAL sizes (new Blobs.Size on local/S3 stores), so the sum cannot drift from the volume; the write is capped at the remaining budget, so the excess is refused before it is persisted, and a re-upload is charged only for its new bytes. - API: per-user create/upload throttles (30s/15s, cmd/felis-wired) on a dedicated cooldown keyspace, reserve→release so a failed attempt never burns the window and a burst collapses to one winner; ErrQuotaExceeded → 403 submission_quota_exceeded (distinct from the 400 an oversize blob gets), 429 submission_cooldown for the throttles. - Panel: zh/en copy for both codes; openapi documents 403/429 on the two user routes; pgint covers the pending-queue count. Unit tests: submit package (cap, budget boundary/exact-fit/replacement, oversize-vs-quota split) and api handlers (quota 403 both paths, throttle 429 + recovery + failure-release). go vet/go test/gofmt clean; panel vitest 118 + typecheck green. --- cmd/felis/api.go | 5 + docs/openapi.yaml | 23 +++- internal/api/api.go | 26 ++++ internal/api/submissions.go | 62 ++++++++- internal/api/submissions_test.go | 110 +++++++++++++++ internal/pgint/pgint_test.go | 12 ++ internal/submit/blobstore.go | 20 +++ internal/submit/pgstore.go | 8 ++ internal/submit/s3store.go | 18 +++ internal/submit/submit.go | 144 ++++++++++++++++++-- internal/submit/submit_test.go | 148 +++++++++++++++++++++ panel/src/i18n/resources/en-US/errors.json | 2 + panel/src/i18n/resources/zh-CN/errors.json | 2 + panel/src/lib/api.test.ts | 8 ++ panel/src/lib/api.ts | 4 + 15 files changed, 577 insertions(+), 15 deletions(-) diff --git a/cmd/felis/api.go b/cmd/felis/api.go index af987c4..d0343a6 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -291,6 +291,11 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { AdminHostname: cfg.Auth.AdminHostname, PanelHostname: cfg.Auth.PanelHostname, WakeCooldown: 30 * time.Second, + // The user-modpack lane's per-user throttles: a create spaces out + // review-queue rows, an upload spaces out (up to 1 GiB) context streams. + // Separate keys, so the normal create→upload sequence stays immediate. + SubmitCreateCooldown: 30 * time.Second, + SubmitUploadCooldown: 15 * time.Second, // Bound concurrent console/build-log SSE streams per principal. Generous enough // for legitimate multi-tab / multi-server watching, while capping how many // upstream follow connections a single caller can tie up if their streams stall. diff --git a/docs/openapi.yaml b/docs/openapi.yaml index 45205c3..0692b34 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -4266,6 +4266,15 @@ paths: $ref: '#/components/responses/BadRequest' '401': $ref: '#/components/responses/Unauthorized' + '403': + description: >- + The per-user submission allowance is spent — too many of the + caller's submissions are awaiting review, or their stored-upload + budget is full (submission_quota_exceeded). + '429': + description: >- + A submission was created within the per-user cooldown window + (submission_cooldown). '503': $ref: '#/components/responses/ServiceUnavailable' get: @@ -4307,8 +4316,10 @@ paths: principal; a submission the caller does not own is reported as 404, so this endpoint cannot upload to or probe another user's submission. Only a pending_review submission accepts a context (409 otherwise); a wrong-format - or oversize body is rejected with 400. Returns 503 when the deployment's - context store has no implemented upload transport. + or oversize body is rejected with 400, and an upload that would push the + caller past their per-user stored-context budget is refused with 403 + before the excess is persisted. Returns 503 when the deployment's context + store has no implemented upload transport. x-felis-face: [external] x-felis-tier: app security: [{ accessJWT: [] }] @@ -4329,10 +4340,18 @@ paths: $ref: '#/components/responses/BadRequest' '401': $ref: '#/components/responses/Unauthorized' + '403': + description: >- + The upload would exceed the caller's per-user stored-context budget + (submission_quota_exceeded). '404': $ref: '#/components/responses/NotFound' '409': $ref: '#/components/responses/Conflict' + '429': + description: >- + An upload was accepted within the per-user cooldown window + (submission_cooldown). '503': $ref: '#/components/responses/ServiceUnavailable' diff --git a/internal/api/api.go b/internal/api/api.go index 2827d09..a9ec66f 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -123,6 +123,15 @@ type API struct { // on the wake lever). Zero disables throttling. WakeCooldown time.Duration + // SubmitCreateCooldown / SubmitUploadCooldown throttle the user-modpack + // submission lane per user: create bounds how quickly review-queue rows can + // appear, upload bounds how often a user may stream a (up to 1 GiB) build + // context. The keys are separate, so the lane's normal shape — create, then + // upload — is never blocked by its own throttle. Zero disables each lever + // (the same idiom as WakeCooldown); cmd/felis wires positive values. + SubmitCreateCooldown time.Duration + SubmitUploadCooldown time.Duration + // MaxRunningServers caps how many servers may be desired-Running cluster-wide // (spec §9.1: the concurrency-上限 lever hanging on the same wake chokepoint as // cooldown and autostartPolicy). Zero — the default — disables it: §9.2 wires @@ -158,6 +167,9 @@ type API struct { otpCooldownOnce sync.Once otpCooldown *cooldownLimiter + submitCooldownOnce sync.Once + submitCooldown *cooldownLimiter + streamCapOnce sync.Once streamCap *streamLimiter } @@ -204,6 +216,20 @@ func (a *API) otpLimiter() *cooldownLimiter { return a.otpCooldown } +// submitLimiter lazily builds a SEPARATE cooldown limiter for the user-modpack +// submission lane, so its throttles never share state with the wake or OTP +// keyspaces. One limiter backs both levers with prefixed keys (see the +// submissionCreateKey/UploadKey constants), so create and upload never contend +// with each other. Like the other cooldowns it is process-local; with multiple +// api replicas the effective spacing is per-replica, the same accepted +// KNOWN-LIMITATION the OTP resend throttle carries. +func (a *API) submitLimiter() *cooldownLimiter { + a.submitCooldownOnce.Do(func() { + a.submitCooldown = &cooldownLimiter{now: a.now, last: map[string]time.Time{}} + }) + return a.submitCooldown +} + // streamGate lazily builds the per-principal SSE stream cap bound to // MaxStreamsPerPrincipal. A zero cap yields a disabled limiter that admits every // stream, so a deployment (or test) that leaves it unset pays nothing. diff --git a/internal/api/submissions.go b/internal/api/submissions.go index 0034c17..4d39100 100644 --- a/internal/api/submissions.go +++ b/internal/api/submissions.go @@ -59,6 +59,15 @@ type createSubmissionRequest struct { DisplayName string `json:"display_name"` } +// The submission lane's two cooldown keys, prefixed into the shared submit +// limiter's per-user keys. Create and upload are separate levers on purpose: +// creating a submission and then immediately uploading its context is the lane's +// normal shape, so one must never consume the other's window. +const ( + submissionCreateKey = "create:" + submissionUploadKey = "upload:" +) + // rejectSubmissionRequest is the POST /submissions/{id}/reject body. A reason is // required (the submit layer rejects an empty one with 400). type rejectSubmissionRequest struct { @@ -74,6 +83,26 @@ func (a *API) handleCreateSubmission(w http.ResponseWriter, r *http.Request) { return } p := principalFromContext(r.Context()) + // Reserve the per-user create cooldown BEFORE the store write. Unlike the + // wake lever's allowed→record (whose real gate is the running cap and whose + // effect is idempotent), a create is a non-idempotent row insertion with no + // other bound on its rate, so a burst of truly concurrent creates must yield + // exactly one winner per window. The deferred rollback frees the window + // whenever the create fails — a 400 typo, a spent quota, a store error — so + // only a row that was actually recorded consumes it. + lim := a.submitLimiter() + reservedAt, ok := lim.reserve(submissionCreateKey+p.UserID, a.SubmitCreateCooldown) + if !ok { + writeError(w, r, newError(http.StatusTooManyRequests, "submission_cooldown", + "a submission was created recently; wait a moment before creating another")) + return + } + committed := false + defer func() { + if !committed { + lim.release(submissionCreateKey+p.UserID, reservedAt) + } + }() var body createSubmissionRequest if err := decodeJSON(w, r, &body); err != nil { writeError(w, r, err) @@ -87,6 +116,7 @@ func (a *API) handleCreateSubmission(w http.ResponseWriter, r *http.Request) { writeSubmitError(w, r, err) return } + committed = true a.audit(r, p.Email, "submission.create", sub.ID) writeJSON(w, http.StatusCreated, sub) } @@ -107,12 +137,34 @@ func (a *API) handleUploadSubmissionContext(w http.ResponseWriter, r *http.Reque return } p := principalFromContext(r.Context()) + // Reserve the per-user upload cooldown BEFORE streaming. The body is the + // expensive part (up to the 1 GiB blob cap), so without a reservation the + // throttle would never bound the resource it exists for: a caller could + // repeatedly start long uploads and abort them. Reserving also collapses the + // lane's parallel overshoot — a burst of concurrent uploads from one user + // 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) + 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) + } + }() 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 a.audit(r, p.Email, "submission.upload", sub.ID) writeJSON(w, http.StatusOK, sub) } @@ -250,9 +302,10 @@ var errSubmissionsUnavailable = newError(http.StatusServiceUnavailable, "submiss // writeSubmitError maps submit-package errors onto HTTP status codes. Only the // business sentinels are client-facing: a validation failure is 400, a missing -// submission is 404, an already-reviewed submission is 409, and an unconfigured -// upload transport is 503 (the store this deployment set has no implemented -// transport — an honest "not available here", not a client error). Everything +// submission is 404, an already-reviewed submission is 409, a spent per-user +// allowance is 403 (the same status the server-resource quota answers with), and +// an unconfigured upload transport is 503 (the store this deployment set has no +// implemented transport — an honest "not available here", not a client error). Everything // else — including a build.ErrInvalid raised by the pre-CAS build.Validate (a // platform registry/context MISCONFIGURATION, never client input, since every // build input is platform-derived) and a post-CAS Submit hand-off failure — is a @@ -267,6 +320,9 @@ func writeSubmitError(w http.ResponseWriter, r *http.Request, err error) { case errors.Is(err, submit.ErrAlreadyReviewed): writeError(w, r, newError(http.StatusConflict, "already_reviewed", "submission has already been reviewed")) + case errors.Is(err, submit.ErrQuotaExceeded): + writeError(w, r, newError(http.StatusForbidden, "submission_quota_exceeded", + "submission quota reached")) case errors.Is(err, submit.ErrBlobNotFound): writeError(w, r, newError(http.StatusNotFound, "not_found", "no context uploaded for this submission")) case errors.Is(err, submit.ErrUploadsUnavailable): diff --git a/internal/api/submissions_test.go b/internal/api/submissions_test.go index 5629d55..89aafc3 100644 --- a/internal/api/submissions_test.go +++ b/internal/api/submissions_test.go @@ -9,6 +9,7 @@ import ( "net/http" "strings" "testing" + "time" "felis.lolicon.best/internal/build" "felis.lolicon.best/internal/submit" @@ -567,3 +568,112 @@ func TestAdminSubmissionContextRoute(t *testing.T) { } }) } + +// A spent per-user allowance is 403 submission_quota_exceeded on both the create +// and the upload path — distinctly NOT the 400 a malformed request gets, and not +// the 429 the cooldown answers with. +func TestSubmissionQuotaIs403(t *testing.T) { + t.Run("create", func(t *testing.T) { + fs := &fakeSubmissions{createErr: fmt.Errorf("%w: 5 submissions are already awaiting review", submit.ErrQuotaExceeded)} + w := do(appSubAPI(fs).ExternalHandler(), "POST", "/api/v1/me/submissions", `{"display_name":"Pack"}`, nil) + if w.Code != http.StatusForbidden { + t.Fatalf("code = %d, want 403 (%s)", w.Code, w.Body.String()) + } + if got := decodeErr(t, w); got != "submission_quota_exceeded" { + t.Errorf("error code = %q, want submission_quota_exceeded", got) + } + }) + t.Run("upload", func(t *testing.T) { + fs := &fakeSubmissions{uploadErr: fmt.Errorf("%w: exceeds your remaining storage allowance", submit.ErrQuotaExceeded)} + w := do(appSubAPI(fs).ExternalHandler(), "POST", "/api/v1/me/submissions/sub-9/context", "\x1f\x8bdata", nil) + if w.Code != http.StatusForbidden { + t.Fatalf("code = %d, want 403 (%s)", w.Code, w.Body.String()) + } + if got := decodeErr(t, w); got != "submission_quota_exceeded" { + t.Errorf("error code = %q, want submission_quota_exceeded", got) + } + }) +} + +// The per-user create cooldown bounds review-queue growth: a second create in +// the same window is 429 submission_cooldown and never reaches the service; the +// window recovers afterwards. +func TestCreateSubmissionRateLimited(t *testing.T) { + fs := &fakeSubmissions{} + api := appSubAPI(fs) + clock := time.Unix(1_700_000_000, 0) + api.Now = func() time.Time { return clock } + api.SubmitCreateCooldown = time.Minute + eh := api.ExternalHandler() + + if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"First"}`, nil); w.Code != http.StatusCreated { + t.Fatalf("first create: code = %d, want 201 (%s)", w.Code, w.Body.String()) + } + w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Second"}`, nil) + if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "submission_cooldown" { + t.Fatalf("immediate second create: code = %d body %s, want 429 submission_cooldown", w.Code, w.Body.String()) + } + // The gate sits before the body handling: even a malformed request is + // refused while the window is closed, so it cannot be used to probe. + if w := do(eh, "POST", "/api/v1/me/submissions", `{`, nil); w.Code != http.StatusTooManyRequests { + t.Fatalf("malformed create during cooldown: code = %d, want 429", w.Code) + } + clock = clock.Add(time.Minute + time.Second) + if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Third"}`, nil); w.Code != http.StatusCreated { + t.Fatalf("post-cooldown create: code = %d, want 201 (%s)", w.Code, w.Body.String()) + } +} + +// A failed create frees the window: only a row that was actually recorded burns +// the cooldown, so a validation typo is not punished with a wait. +func TestCreateSubmissionFailureDoesNotBurnCooldown(t *testing.T) { + fs := &fakeSubmissions{createErr: fmt.Errorf("%w: display name is required", submit.ErrInvalid)} + api := appSubAPI(fs) + api.Now = func() time.Time { return time.Unix(1_700_000_000, 0) } + api.SubmitCreateCooldown = time.Minute + eh := api.ExternalHandler() + + if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":""}`, nil); w.Code != http.StatusBadRequest { + t.Fatalf("failed create: code = %d, want 400", w.Code) + } + fs.createErr = nil + if w := do(eh, "POST", "/api/v1/me/submissions", `{"display_name":"Fixed"}`, nil); w.Code != http.StatusCreated { + t.Fatalf("retry at the same instant: code = %d, want 201 (%s)", w.Code, w.Body.String()) + } +} + +// The per-user upload cooldown bounds context streaming: a second upload in the +// same window is 429 submission_cooldown, and a FAILED upload frees the window +// for an immediate retry. +func TestUploadSubmissionContextRateLimited(t *testing.T) { + fs := &fakeSubmissions{} + api := appSubAPI(fs) + clock := time.Unix(1_700_000_000, 0) + api.Now = func() time.Time { return clock } + api.SubmitUploadCooldown = time.Minute + eh := api.ExternalHandler() + body := "\x1f\x8b\x08\x00 the modpack bytes" + + if w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusOK { + t.Fatalf("first upload: code = %d, want 200 (%s)", w.Code, w.Body.String()) + } + w := do(eh, "POST", "/api/v1/me/submissions/sub-9/context", body, nil) + if w.Code != http.StatusTooManyRequests || decodeErr(t, w) != "submission_cooldown" { + t.Fatalf("immediate second upload: code = %d body %s, want 429 submission_cooldown", w.Code, w.Body.String()) + } + + // A failed upload releases its reservation, so the user is not punished for + // a genuine failure (aborted transfer, spent quota) with a cooldown wait. + fs2 := &fakeSubmissions{uploadErr: submit.ErrUploadsUnavailable} + api2 := appSubAPI(fs2) + api2.Now = func() time.Time { return clock } + api2.SubmitUploadCooldown = time.Minute + eh2 := api2.ExternalHandler() + if w := do(eh2, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusServiceUnavailable { + t.Fatalf("failed upload: code = %d, want 503", w.Code) + } + fs2.uploadErr = nil + if w := do(eh2, "POST", "/api/v1/me/submissions/sub-9/context", body, nil); w.Code != http.StatusOK { + t.Fatalf("retry at the same instant after failure: code = %d, want 200 (%s)", w.Code, w.Body.String()) + } +} diff --git a/internal/pgint/pgint_test.go b/internal/pgint/pgint_test.go index d42b0c3..6f59753 100644 --- a/internal/pgint/pgint_test.go +++ b/internal/pgint/pgint_test.go @@ -980,6 +980,14 @@ func TestSubmitStoreContract(t *testing.T) { if err := s.CreateSubmission(ctx, sub); err != nil { t.Fatalf("CreateSubmission: %v", err) } + // The pending-queue count is the read behind the per-user pending cap: a + // fresh user starts at zero and a pending row counts. + if n, err := s.CountPendingSubmissionsBy(ctx, u.ID); err != nil || n != 1 { + t.Fatalf("CountPendingSubmissionsBy after create = (%d, %v), want (1, nil)", n, err) + } + if n, err := s.CountPendingSubmissionsBy(ctx, u.ID+"-nobody"); err != nil || n != 0 { + t.Fatalf("CountPendingSubmissionsBy for an unknown user = (%d, %v), want (0, nil)", n, err) + } got, err := s.GetSubmission(ctx, id) if err != nil { t.Fatalf("GetSubmission: %v", err) @@ -1012,6 +1020,10 @@ func TestSubmitStoreContract(t *testing.T) { if ok, err := s.RejectSubmission(ctx, id, "reviewer@example.net", "no", now); err != nil || ok { t.Fatalf("reject after approve = (%v, %v), want (false, nil)", ok, err) } + // A reviewed row leaves the pending queue. + if n, err := s.CountPendingSubmissionsBy(ctx, u.ID); err != nil || n != 0 { + t.Fatalf("CountPendingSubmissionsBy after approve = (%d, %v), want (0, nil)", n, err) + } got, _ = s.GetSubmission(ctx, id) if got.Status != submit.StatusApproved || got.ImageRef == "" || got.ReviewedBy != "reviewer@example.net" || got.ReviewedAt == nil { t.Fatalf("approved row = %+v", got) diff --git a/internal/submit/blobstore.go b/internal/submit/blobstore.go index ff20c85..c82403f 100644 --- a/internal/submit/blobstore.go +++ b/internal/submit/blobstore.go @@ -129,5 +129,25 @@ func (s *LocalContextStore) Open(_ context.Context, id string) (io.ReadCloser, e return f, nil } +// Size reports the stored blob's size — the accounting read behind the per-user +// storage budget. A missing blob is (0, false, nil): absence is not an error +// here, it is simply no bytes to count (the same distinction Exists draws for +// Approve). +func (s *LocalContextStore) Size(_ context.Context, id string) (int64, bool, error) { + dir, err := s.dir(id) + if err != nil { + return 0, false, err + } + fi, err := os.Stat(filepath.Join(dir, contextBlobName)) + 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 context blob: %w", err) + } +} + // Compile-time proof that the filesystem store satisfies the Blobs transport. var _ Blobs = (*LocalContextStore)(nil) diff --git a/internal/submit/pgstore.go b/internal/submit/pgstore.go index 781b258..58a0252 100644 --- a/internal/submit/pgstore.go +++ b/internal/submit/pgstore.go @@ -33,6 +33,14 @@ func (s *PGStore) CreateSubmission(ctx context.Context, sub *Submission) error { return err } +func (s *PGStore) CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error) { + const q = `SELECT count(*) FROM image_submissions + WHERE submitted_by = $1 AND status = 'pending_review'` + var n int + err := s.db.QueryRowContext(ctx, q, submittedBy).Scan(&n) + return n, err +} + func (s *PGStore) GetSubmission(ctx context.Context, id string) (*Submission, error) { const q = `SELECT ` + submissionColumns + ` FROM image_submissions WHERE id = $1` return scanSubmission(s.db.QueryRowContext(ctx, q, id)) diff --git a/internal/submit/s3store.go b/internal/submit/s3store.go index d5d1852..80c4f98 100644 --- a/internal/submit/s3store.go +++ b/internal/submit/s3store.go @@ -187,6 +187,24 @@ func (s *S3ContextStore) Exists(ctx context.Context, id string) (bool, error) { return true, nil } +// Size reports the stored blob's size — the accounting read behind the per-user +// storage budget. A missing object is (0, false, nil), the same distinction +// Exists draws. +func (s *S3ContextStore) Size(ctx context.Context, id string) (int64, bool, error) { + key, err := s.keyFor(id) + if err != nil { + return 0, false, err + } + info, err := s.client.StatObject(ctx, s.bucket, key, minio.StatObjectOptions{}) + if err != nil { + if isS3NotFound(err) { + return 0, false, nil + } + return 0, false, fmt.Errorf("submit: stat context blob: %w", err) + } + return info.Size, true, nil +} + // Open returns the stored context blob for id — the read side of the transport the // build Pod's fetch initContainer uses. minio's GetObject returns only once the // server answered with an object (it surfaces NoSuchKey up front), so a missing diff --git a/internal/submit/submit.go b/internal/submit/submit.go index 7d6fc6f..a06a0c7 100644 --- a/internal/submit/submit.go +++ b/internal/submit/submit.go @@ -86,6 +86,12 @@ var ( ErrInvalid = errors.New("submit: invalid request") ErrNotFound = errors.New("submit: submission not found") ErrAlreadyReviewed = errors.New("submit: submission already reviewed") + // ErrQuotaExceeded reports that the caller's upload allowance is spent — + // either too many of their submissions are already awaiting review, or their + // stored contexts already fill the per-user byte budget. The request is not + // malformed; the allowance is exhausted. The API maps it to 403, matching the + // server-resource quota's status (spec §7). + ErrQuotaExceeded = errors.New("submit: quota exceeded") // ErrUploadsUnavailable means this deployment configured a context store with // no implemented upload transport (a nil Manager.Blobs — e.g. an object-store // base with no client wired). UploadContext returns it so the endpoint reports @@ -108,14 +114,37 @@ func invalidf(format string, a ...any) error { // preserves that chain through to the API error mapper. var errContextTooLarge = fmt.Errorf("%w: build context exceeds the maximum allowed size", ErrInvalid) +// errStorageQuota trips when an upload would push the user past their per-user +// storage budget. It wraps ErrQuotaExceeded so the API answers 403, distinctly +// from errContextTooLarge's 400: the blob is fine, the allowance is spent. +var errStorageQuota = fmt.Errorf("%w: the upload exceeds your remaining storage allowance", ErrQuotaExceeded) + const ( maxDisplayName = 200 maxRejectReason = 1000 // defaultMaxContextBytes caps an uploaded build-context blob. Modpack contexts // (mods, configs, an occasional bundled world) are large, so the cap is // generous; it bounds what one untrusted upload can write to the uploads PVC, - // not a tight quota. Override per-Manager via MaxContextBytes. + // one blob at a time. The per-user budget below bounds the SUM across a user's + // uploads; this cap is what keeps any single write bounded. Override per-Manager + // via MaxContextBytes. defaultMaxContextBytes = 1 << 30 // 1 GiB + // defaultMaxPendingPerUser caps how many pending_review submissions one user + // may hold at once. Untrusted ingress has no natural bound — a logged-in + // player could otherwise file rows all day — and every pending row is a + // review-queue item an admin has to read, so the cap is small on purpose: + // enough to stage a couple of packs, far short of a flood. Override per + // Manager via MaxPendingPerUser. + defaultMaxPendingPerUser = 5 + // defaultMaxStoredBytesPerUser caps the total bytes one user's stored + // contexts may occupy on the uploads store. The uploads PVC renders at a + // fixed 5Gi (platform/workloads.go); without a per-user budget one account + // could fill it and every other user's upload would start failing. Two GiB + // leaves room for a couple of full-size modpacks (a single blob may be 1 GiB) + // while keeping a small user base from exhausting the volume; size the PVC + // above users × this budget before raising it. Override per Manager via + // MaxStoredBytesPerUser. + defaultMaxStoredBytesPerUser = 2 << 30 // 2 GiB ) // displayNameRE constrains the user-supplied label to a calm, single-line set: @@ -149,6 +178,10 @@ type Submission struct { type Store interface { // CreateSubmission inserts a pending_review row. CreateSubmission(ctx context.Context, s *Submission) error + // CountPendingSubmissionsBy reports how many of one user's submissions are + // still pending_review — the queue-length read behind the per-user pending + // cap in Create. + CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error) // GetSubmission loads one submission, or ErrNotFound. GetSubmission(ctx context.Context, id string) (*Submission, error) // ListSubmissions returns every submission, newest first (admin queue). @@ -196,6 +229,12 @@ type Builds interface { type Blobs interface { Put(ctx context.Context, id string, r io.Reader) (int64, error) Exists(ctx context.Context, id string) (bool, error) + // Size returns the stored blob's size in bytes; ok=false means no blob is + // stored for id. UploadContext sums this over a user's submissions to enforce + // the per-user storage budget, so it must report what is actually on the + // store — never a recorded number that could drift from it (a re-upload + // supersedes the previous blob in place). + Size(ctx context.Context, id string) (int64, bool, error) // Open returns the stored blob's bytes for the internal context-fetch route // the build Pod's initContainer dials (cmd/felis fetch-context). It returns an // error wrapping ErrBlobNotFound when no blob exists, so the route can answer @@ -239,6 +278,12 @@ type Manager struct { // MaxContextBytes overrides the uploaded-context size cap; 0 uses // defaultMaxContextBytes. MaxContextBytes int64 + // MaxPendingPerUser overrides how many of one user's submissions may await + // review at once; 0 uses defaultMaxPendingPerUser. + MaxPendingPerUser int + // MaxStoredBytesPerUser overrides the per-user stored-context budget; 0 uses + // defaultMaxStoredBytesPerUser. + MaxStoredBytesPerUser int64 Now func() time.Time IDGen func() string @@ -251,6 +296,20 @@ func (m *Manager) maxContextBytes() int64 { return defaultMaxContextBytes } +func (m *Manager) maxPendingPerUser() int { + if m.MaxPendingPerUser > 0 { + return m.MaxPendingPerUser + } + return defaultMaxPendingPerUser +} + +func (m *Manager) maxStoredBytesPerUser() int64 { + if m.MaxStoredBytesPerUser > 0 { + return m.MaxStoredBytesPerUser + } + return defaultMaxStoredBytesPerUser +} + func (m *Manager) now() time.Time { if m.Now != nil { return m.Now() @@ -344,6 +403,18 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e return nil, invalidf("submitter identity is required") } + // Per-user pending cap: every pending row is a review-queue item an admin + // must read, so one account may not park an unbounded number of them. The + // check and the insert are not atomic (two concurrent creates may jointly + // land one row over the cap); that is a soft overshoot of a queue-length + // lever, not a resource bound, so it is deliberately not worth a lock. + if pending, err := m.Store.CountPendingSubmissionsBy(ctx, req.SubmittedBy); err != nil { + return nil, err + } else if pending >= m.maxPendingPerUser() { + return nil, fmt.Errorf("%w: %d submissions are already awaiting review (limit %d)", + ErrQuotaExceeded, pending, m.maxPendingPerUser()) + } + id := m.newID() s := &Submission{ ID: id, @@ -373,7 +444,10 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e // - the context is mutable ONLY while pending_review — once approved the build // has already consumed it, once rejected it is dead; // - the body must be a gzip tarball (context.tar.gz) and is size-capped, so a -// wrong-format or oversize upload is rejected as a 400 without persisting. +// wrong-format or oversize upload is rejected as a 400 without persisting; +// - a user's stored contexts are budgeted (MaxStoredBytesPerUser): the write +// is capped at the remaining budget, so an upload that would exceed it is +// refused as a spent allowance (403) before the excess is persisted. // // A re-upload while still pending atomically supersedes the previous blob, so a // user can fix their pack before an admin reviews it. @@ -404,22 +478,72 @@ 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)") } - // Cap the size: cappedReader trips errContextTooLarge on the first byte past - // the limit, so the store never persists an oversize blob (it removes its temp - // file on the copy error) and the failure surfaces as a 400, not a 500. - if _, err := m.Blobs.Put(ctx, id, &cappedReader{r: br, left: m.maxContextBytes()}); err != nil { + // 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, err := m.storedBytes(ctx, submittedBy, id) + if err != nil { + return nil, err + } + remaining := m.maxStoredBytesPerUser() - used + if remaining <= 0 { + return nil, errStorageQuota + } + 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, never as a malformed request. + limit, over = remaining, errStorageQuota + } + if _, err := m.Blobs.Put(ctx, id, &cappedReader{r: br, left: limit, over: over}); err != nil { return nil, err } return sub, nil } +// storedBytes sums the stored-blob sizes of submittedBy's submissions, 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 sum cannot drift +// from what is actually occupying the volume (including blobs uploaded before +// any budget existed). +func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string) (int64, error) { + subs, err := m.Store.ListSubmissionsBy(ctx, submittedBy) + if err != nil { + return 0, err + } + var total int64 + for _, s := range subs { + if s.ID == excludeID { + continue + } + n, ok, err := m.Blobs.Size(ctx, s.ID) + if err != nil { + return 0, err + } + if ok { + total += n + } + } + return total, nil +} + // cappedReader passes through at most left bytes; the first byte beyond the limit -// trips errContextTooLarge. It reads one probe byte past the limit to tell an -// exactly-at-limit blob (accepted) from a larger one (rejected), so a stream of -// exactly the cap is never falsely rejected. +// trips over — errContextTooLarge for the single-blob cap, errStorageQuota when +// the per-user budget binds first. It reads one probe byte past the limit to tell +// an exactly-at-limit blob (accepted) from a larger one (rejected), so a stream +// of exactly the cap is never falsely rejected. type cappedReader struct { r io.Reader left int64 + over error } func (c *cappedReader) Read(p []byte) (int, error) { @@ -429,7 +553,7 @@ func (c *cappedReader) Read(p []byte) (int, error) { var probe [1]byte n, err := c.r.Read(probe[:]) if n > 0 { - return 0, errContextTooLarge + return 0, c.over } if err == nil { return 0, io.EOF diff --git a/internal/submit/submit_test.go b/internal/submit/submit_test.go index e52ef08..3ea0c5f 100644 --- a/internal/submit/submit_test.go +++ b/internal/submit/submit_test.go @@ -24,6 +24,7 @@ type fakeBlobs struct { stored map[string][]byte putErr error existsErr error + sizeErr error forceExists *bool // overrides the stored-map lookup for the approve-gate tests } @@ -52,6 +53,18 @@ func (f *fakeBlobs) Exists(_ context.Context, id string) (bool, error) { return ok, nil } +// Size mirrors the real stores: a missing blob is (0, false, nil). +func (f *fakeBlobs) Size(_ context.Context, id string) (int64, bool, error) { + if f.sizeErr != nil { + return 0, false, f.sizeErr + } + b, ok := f.stored[id] + if !ok { + return 0, false, nil + } + return int64(len(b)), true, nil +} + func (f *fakeBlobs) Open(_ context.Context, id string) (io.ReadCloser, error) { b, ok := f.stored[id] if !ok { @@ -73,6 +86,7 @@ type fakeStore struct { subs map[string]*Submission createErr error + countErr error approveErr error rejectErr error linkErr error @@ -91,6 +105,19 @@ func (f *fakeStore) CreateSubmission(_ context.Context, s *Submission) error { return nil } +func (f *fakeStore) CountPendingSubmissionsBy(_ context.Context, by string) (int, error) { + if f.countErr != nil { + return 0, f.countErr + } + var n int + for _, s := range f.subs { + if s.SubmittedBy == by && s.Status == StatusPendingReview { + n++ + } + } + return n, nil +} + func (f *fakeStore) GetSubmission(_ context.Context, id string) (*Submission, error) { s, ok := f.subs[id] if !ok { @@ -671,6 +698,127 @@ func TestUploadContextNoTransportUnavailable(t *testing.T) { } } +// The per-user pending cap bounds the review queue: at the limit a new create is +// refused with ErrQuotaExceeded (403), a reviewed row frees a slot, and another +// user's queue is unaffected. +func TestCreatePendingCapRejects(t *testing.T) { + m, _, _ := newManager() + m.MaxPendingPerUser = 2 + ctx := context.Background() + + for i := 0; i < 2; i++ { + if _, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}); err != nil { + t.Fatalf("create %d: %v", i+1, err) + } + } + _, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}) + if !errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("create past the cap = %v, want ErrQuotaExceeded", err) + } + + // A verdict moves the row out of pending_review, so the slot frees up. + if _, err := m.Reject(ctx, "sub-1", "admin@example.net", "out of scope"); err != nil { + t.Fatalf("Reject: %v", err) + } + if _, err := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}); err != nil { + t.Fatalf("create after review: %v", err) + } + + // The cap is per user, never global. + if _, err := m.Create(ctx, CreateRequest{DisplayName: "Other", SubmittedBy: "user-2"}); err != nil { + t.Fatalf("other user's first create: %v", err) + } +} + +// The per-user storage budget bounds the sum of a user's stored contexts. It is +// charged against the blob store's actual sizes, refuses at the boundary with +// ErrQuotaExceeded (403 — not the 400 a single oversize blob gets), and does not +// double-charge a re-upload of the same submission. +func TestUploadContextStorageBudget(t *testing.T) { + m, _, _ := newManager() + fb := newFakeBlobs() + m.Blobs = fb + m.MaxStoredBytesPerUser = 10 // tiny budget; gzBody is 4 magic bytes + payload + ctx := context.Background() + + a, err := m.Create(ctx, CreateRequest{DisplayName: "A", SubmittedBy: "user-1"}) + if err != nil { + t.Fatalf("create A: %v", err) + } + b, err := m.Create(ctx, CreateRequest{DisplayName: "B", SubmittedBy: "user-1"}) + if err != nil { + t.Fatalf("create B: %v", err) + } + + if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("a"))); err != nil { + t.Fatalf("first upload: %v", err) + } + // A's 5 bytes leave 5. B's exactly-5-byte blob must be accepted — the budget + // binds only past the limit, never at it. + if _, err := m.UploadContext(ctx, b.ID, "user-1", strings.NewReader(gzBody("b"))); err != nil { + t.Fatalf("upload exactly at the budget = %v, want accepted", err) + } + if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("far too much"))); err != nil { + // A re-upload is charged only for its NEW bytes (its old blob is + // superseded), and 16 bytes exceed the 5 bytes left after B. + if !errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("re-upload past the budget = %v, want ErrQuotaExceeded", err) + } + } else { + t.Fatal("re-upload past the budget was accepted") + } + // Nothing oversize persisted, and A's good blob was not clobbered. + if string(fb.stored[a.ID]) != gzBody("a") { + t.Fatalf("A's blob = %q, want the original (a failed re-upload must not replace it)", fb.stored[a.ID]) + } + + // B now holds 5 and the budget is full: a fresh submission's upload is + // refused up front, before reading any body. + c, err := m.Create(ctx, CreateRequest{DisplayName: "C", SubmittedBy: "user-1"}) + if err != nil { + t.Fatalf("create C: %v", err) + } + _, err = m.UploadContext(ctx, c.ID, "user-1", strings.NewReader(gzBody("c"))) + if !errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("upload with no budget left = %v, want ErrQuotaExceeded", err) + } + if _, ok := fb.stored[c.ID]; ok { + t.Fatal("a budget-refused upload must persist nothing") + } + + // A 5-byte replacement of A fits exactly (10 − B's 5), proving the + // replacement is not double-charged against its own old bytes. + if _, err := m.UploadContext(ctx, a.ID, "user-1", strings.NewReader(gzBody("z"))); err != nil { + t.Fatalf("budget-exact replacement = %v, want accepted", err) + } +} + +// A single oversize blob stays a 400 (ErrInvalid), distinct from the 403 the +// per-user budget answers with — the two failure classes must not collapse. +func TestUploadContextOversizeIsNotQuotaError(t *testing.T) { + m, _, _ := newManager() + fb := newFakeBlobs() + m.Blobs = fb + m.MaxContextBytes = 4 + seed, _ := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}) + + _, err := m.UploadContext(context.Background(), seed.ID, "user-1", strings.NewReader(gzBody("too big"))) + if !errors.Is(err, ErrInvalid) || errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("err = %v, want ErrInvalid and NOT ErrQuotaExceeded", err) + } +} + +// A store failure while counting the pending queue must surface as-is, never as +// a quota verdict that blames the user. +func TestCreatePendingCountFailureSurfaces(t *testing.T) { + m, st, _ := newManager() + st.countErr = errors.New("db is down") + _, err := m.Create(context.Background(), CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"}) + if err == nil || errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("err = %v, want the raw store failure", err) + } +} + func TestApproveRefusesMissingContext(t *testing.T) { // With a transport wired, approving a submission whose context was never // uploaded fails BEFORE the CAS: the row stays pending and no build starts. diff --git a/panel/src/i18n/resources/en-US/errors.json b/panel/src/i18n/resources/en-US/errors.json index 8f51810..2af285e 100644 --- a/panel/src/i18n/resources/en-US/errors.json +++ b/panel/src/i18n/resources/en-US/errors.json @@ -42,6 +42,8 @@ "build_unavailable": "Image builds aren't available right now.", "build_logs_unavailable": "Build logs aren't available right now.", "already_reviewed": "This submission has already been reviewed.", + "submission_quota_exceeded": "Your submission quota is full: too many pending reviews, or your stored uploads are at the limit.", + "submission_cooldown": "Too many submission requests — try again shortly.", "submissions_unavailable": "Submissions aren't available right now.", "uploads_unavailable": "Uploads aren't available right now.", "backup_unavailable": "Backups aren't available right now.", diff --git a/panel/src/i18n/resources/zh-CN/errors.json b/panel/src/i18n/resources/zh-CN/errors.json index 14f37b8..39d9941 100644 --- a/panel/src/i18n/resources/zh-CN/errors.json +++ b/panel/src/i18n/resources/zh-CN/errors.json @@ -42,6 +42,8 @@ "build_unavailable": "构建功能当前不可用。", "build_logs_unavailable": "构建日志暂时不可用。", "already_reviewed": "该提交已经审核过了。", + "submission_quota_exceeded": "你的提交配额已满:待审核提交过多,或已存上传总量达到上限。", + "submission_cooldown": "操作太频繁——请稍后再试。", "submissions_unavailable": "提交流程当前不可用。", "uploads_unavailable": "上传功能当前不可用。", "backup_unavailable": "备份功能当前不可用。", diff --git a/panel/src/lib/api.test.ts b/panel/src/lib/api.test.ts index e281262..aba7135 100644 --- a/panel/src/lib/api.test.ts +++ b/panel/src/lib/api.test.ts @@ -522,6 +522,14 @@ describe("image whitelist and builds wire shapes", () => { expect((opts as RequestInit).body).toBe(blob); expect((opts as RequestInit).headers).toEqual({ "Content-Type": "application/x-gzip" }); }); + + // The lane's two throttled outcomes (a spent allowance, a closed cooldown) + // must surface as their own copy, not the generic forbidden/error text. + it("maps the submission quota/cooldown codes to stable human copy", async () => { + const { humanizeError } = await import("./api"); + expect(humanizeError({ code: "submission_quota_exceeded" })).toMatch(/quota/i); + expect(humanizeError({ code: "submission_cooldown" })).toMatch(/try again/i); + }); }); describe("updates maintenance window", () => { diff --git a/panel/src/lib/api.ts b/panel/src/lib/api.ts index 92ca188..6c10113 100644 --- a/panel/src/lib/api.ts +++ b/panel/src/lib/api.ts @@ -722,6 +722,10 @@ export function humanizeError(e: unknown): string { return t("build_logs_unavailable"); case "already_reviewed": return t("already_reviewed"); + case "submission_quota_exceeded": + return t("submission_quota_exceeded"); + case "submission_cooldown": + return t("submission_cooldown"); case "submissions_unavailable": return t("submissions_unavailable"); case "uploads_unavailable":