diff --git a/cmd/felis/api.go b/cmd/felis/api.go index bd8acde..419c448 100644 --- a/cmd/felis/api.go +++ b/cmd/felis/api.go @@ -213,6 +213,13 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { ContextBaseURL: internalAPIBaseURL(), Blobs: blobs, } + if v := cfg.Registry.UserUploadsMaxBytes; v != "" { + if n, err := parseByteSize(v); err != nil || n <= 0 { + fmt.Fprintf(stderr, "felis api: [registry] user_uploads_max_bytes %q is not a positive size such as 4Gi; keeping the default\n", v) + } else { + submissions.MaxStoredBytesTotal = n + } + } // Restore subsystem (spec §7): the weak-SA restore Job mounts the target // world PVC + the backup PVC and runs `felis restore`. It needs deployment- @@ -404,6 +411,7 @@ func cmdAPI(args []string, stdout, stderr io.Writer) int { if pruner := registryPruner(cfg, builder.Store, a.Cluster, stderr); pruner != nil { go pruner.Loop(ctx, registryPruneInterval) } + go reapRejectedContexts(ctx, submissions, stderr) select { case <-ctx.Done(): @@ -619,6 +627,29 @@ func reconcileBuilds(ctx context.Context, b *build.Builder, stderr io.Writer) { } } +// reapRejectedContexts deletes, once an hour, the uploaded contexts of +// submissions rejected more than submit.RejectedContextRetention ago. Without it a +// rejected modpack keeps its bytes on the uploads store (and against its +// submitter's budget) until an admin deletes the row. +func reapRejectedContexts(ctx context.Context, m *submit.Manager, stderr io.Writer) { + t := time.NewTicker(time.Hour) + defer t.Stop() + for { + n, err := m.ReapRejected(ctx, submit.RejectedContextRetention) + if err != nil { + fmt.Fprintf(stderr, "felis api: reap rejected uploads: %v\n", err) + } + if n > 0 { + fmt.Fprintf(stderr, "felis api: deleted the uploaded contexts of %d rejected submission(s)\n", n) + } + select { + case <-ctx.Done(): + return + case <-t.C: + } + } +} + // registryPruneInterval spaces the registry pruner's runs. The registry-gc // sidecar sweeps once a day, so pruning more often only changes which sweep frees // a layer. diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index 6653583..2986842 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -2373,7 +2373,7 @@ persisted_registry_block() { out="$(awk ' /^[[:space:]]*\[/ { sect = $0; next } sect ~ /^[[:space:]]*\[registry\][[:space:]]*$/ && - /^[[:space:]]*(kaniko_image|trivy_image|trivy_db_repository|trivy_java_db_repository|build_cpu_limit|build_mem_limit|build_disk_limit|build_user_namespaces|build_runtime_class|max_concurrent_builds|user_uploads_context)[[:space:]]*=/ { print } + /^[[:space:]]*(kaniko_image|trivy_image|trivy_db_repository|trivy_java_db_repository|build_cpu_limit|build_mem_limit|build_disk_limit|build_user_namespaces|build_runtime_class|max_concurrent_builds|user_uploads_context|user_uploads_max_bytes)[[:space:]]*=/ { print } sect ~ /^[[:space:]]*\[registry\.s3\][[:space:]]*$/ && /^[[:space:]]*[A-Za-z_]+[[:space:]]*=/ { if (!s3hdr) { printf "[registry.s3]\n"; s3hdr = 1 } print diff --git a/internal/api/submissions.go b/internal/api/submissions.go index cf5ba5b..2832d33 100644 --- a/internal/api/submissions.go +++ b/internal/api/submissions.go @@ -377,8 +377,9 @@ 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, 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 +// allowance is 403 (the same status the server-resource quota answers with), a +// full uploads store (every user's uploads together at their cap, or the volume +// short of free space) is 507, 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 @@ -402,6 +403,9 @@ func writeSubmitError(w http.ResponseWriter, r *http.Request, err error) { "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.ErrUploadsFull): + writeError(w, r, newError(http.StatusInsufficientStorage, "uploads_full", + "the uploads store is full; an admin has to delete reviewed submissions before new uploads fit")) case errors.Is(err, submit.ErrUploadsUnavailable): writeError(w, r, newError(http.StatusServiceUnavailable, "uploads_unavailable", "modpack upload transport is not configured")) diff --git a/internal/config/config.go b/internal/config/config.go index 3aefd66..41aa7cc 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -207,6 +207,11 @@ type RegistryConfig struct { // archive (§19 WorldArchiver) and a build context (§16) are different artifacts // with different lifecycles, so the two must not share a store binding. UserUploadsContext string `toml:"user_uploads_context"` + // UserUploadsMaxBytes caps what every user's uploaded contexts may occupy + // together, as a quantity ("4Gi"). Each user also has a 2 GiB budget of their + // 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"` // 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 // ArchiveS3Config: Endpoint + Region locate the store and the *Ref fields NAME diff --git a/internal/submit/blobstore.go b/internal/submit/blobstore.go index 406bf5c..1461fd5 100644 --- a/internal/submit/blobstore.go +++ b/internal/submit/blobstore.go @@ -7,6 +7,7 @@ import ( "os" "path/filepath" "regexp" + "syscall" ) // contextBlobName is the fixed object name of a submission's build context under @@ -43,6 +44,39 @@ type LocalContextStore struct { // Base is the directory (uploads PVC mount) submission contexts are written // under. Each submission gets its own {Base}/{id}/ subdirectory. Base string + // MinFree is the share of Base's filesystem an upload must leave free; 0 uses + // DefaultUploadsMinFree. + MinFree float64 +} + +// DefaultUploadsMinFree is the share of the uploads filesystem an upload must +// leave free. On k3s local-path the uploads PVC is a directory on the node's +// disk, beside the worlds and the database, and below about a tenth free the +// kubelet starts evicting pods (the same floor backup.MinFreeAfter keeps). +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 { + 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) + } + 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 + } + floor := uint64(float64(total) * minFree) + if n := uint64(max(need, 0)); avail < n || avail-n < floor { + return fmt.Errorf("%w: %d MiB free of %d MiB, and an upload of up to %d MiB would leave less than %.0f%% free", + ErrUploadsFull, avail>>20, total>>20, n>>20, minFree*100) + } + return nil } // dir returns the per-submission directory, rejecting an id that could escape diff --git a/internal/submit/blobstore_test.go b/internal/submit/blobstore_test.go index ab3b2ff..69d4e54 100644 --- a/internal/submit/blobstore_test.go +++ b/internal/submit/blobstore_test.go @@ -156,3 +156,17 @@ func TestLocalContextStoreRejectsUnsafeID(t *testing.T) { t.Fatal("an unsafe id wrote outside Base") } } + +func TestLocalContextStoreCheckRoom(t *testing.T) { + s := &LocalContextStore{Base: t.TempDir(), MinFree: 1e-9} + if err := s.CheckRoom(1); err != nil { + t.Fatalf("CheckRoom(1 byte) = %v", err) + } + if err := s.CheckRoom(1 << 62); !errors.Is(err, ErrUploadsFull) { + t.Fatalf("CheckRoom(4 EiB) = %v, want ErrUploadsFull", err) + } + s.MinFree = 1 + if err := s.CheckRoom(0); !errors.Is(err, ErrUploadsFull) { + t.Fatalf("CheckRoom with a 100%% floor = %v, want ErrUploadsFull", err) + } +} diff --git a/internal/submit/submit.go b/internal/submit/submit.go index 33b6b3e..59d3462 100644 --- a/internal/submit/submit.go +++ b/internal/submit/submit.go @@ -125,6 +125,12 @@ var errContextTooLarge = fmt.Errorf("%w: build context exceeds the maximum allow // 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) +// ErrUploadsFull means the uploads store as a whole has no room: every user's +// contexts together reached MaxStoredBytesTotal, or the filesystem under a local +// store is close to full. No one's allowance is at fault, so the API answers 507 +// and an admin frees space by deleting reviewed submissions. +var ErrUploadsFull = errors.New("the uploads store is full") + const ( maxDisplayName = 200 maxRejectReason = 1000 @@ -151,8 +157,29 @@ const ( // above users × this budget before raising it. Override per Manager via // MaxStoredBytesPerUser. defaultMaxStoredBytesPerUser = 2 << 30 // 2 GiB + // defaultMaxStoredBytesTotal caps what every user's stored contexts occupy + // together. The per-user budget alone lets enough accounts fill any volume, + // and on k3s local-path the uploads PVC's 5Gi is a label, not a limit: the + // directory sits on the node's disk beside the worlds and the database. Four + // GiB keeps a default install inside that PVC. Override per Manager via + // MaxStoredBytesTotal ([registry] user_uploads_max_bytes). + defaultMaxStoredBytesTotal = 4 << 30 // 4 GiB ) +// RejectedContextRetention is how long a rejected submission keeps its uploaded +// context: long enough for the submitter to read the reason and an admin to look +// again. ReapRejected then deletes the blob; the row stays as the record. +// Approved contexts are kept, since they are what rebuilds the image after the +// registry is lost (docs/troubleshooting.md §9). +const RejectedContextRetention = 7 * 24 * time.Hour + +// RoomChecker is implemented by a Blobs that writes to a filesystem it shares +// with the node. Before an upload, CheckRoom confirms that need more bytes still +// leave the filesystem's floor free, and wraps ErrUploadsFull when they would not. +type RoomChecker interface { + CheckRoom(need int64) error +} + // displayNameRE constrains the user-supplied label to a calm, single-line set: // it is the only free-form string a submission carries, and although JSON // encoding already neutralises it in responses, rejecting control characters @@ -318,6 +345,9 @@ type Manager struct { // MaxStoredBytesPerUser overrides the per-user stored-context budget; 0 uses // defaultMaxStoredBytesPerUser. MaxStoredBytesPerUser int64 + // MaxStoredBytesTotal overrides the budget for every user's stored contexts + // together; 0 uses defaultMaxStoredBytesTotal. + MaxStoredBytesTotal int64 Now func() time.Time IDGen func() string @@ -344,6 +374,13 @@ func (m *Manager) maxStoredBytesPerUser() int64 { return defaultMaxStoredBytesPerUser } +func (m *Manager) maxStoredBytesTotal() int64 { + if m.MaxStoredBytesTotal > 0 { + return m.MaxStoredBytesTotal + } + return defaultMaxStoredBytesTotal +} + func (m *Manager) now() time.Time { if m.Now != nil { return m.Now() @@ -487,7 +524,10 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e // 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. +// refused as a spent allowance (403) before the excess is persisted; +// - so are everyone's together (MaxStoredBytesTotal), and a local store checks +// the filesystem keeps its free floor (RoomChecker): either refuses as +// ErrUploadsFull (507). // // A re-upload while still pending atomically supersedes the previous blob, so a // user can fix their pack before an admin reviews it. The upload hashes what it @@ -531,19 +571,30 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i // 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) + used, total, err := m.storedBytes(ctx, submittedBy, id) if err != nil { return nil, err } - remaining := m.maxStoredBytesPerUser() - used + 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, errStorageQuota + 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, never as a malformed request. - limit, over = remaining, errStorageQuota + // 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 { + if err := rc.CheckRoom(limit); err != nil { + return nil, err + } } h := sha256.New() if _, err := m.Blobs.Put(ctx, id, io.TeeReader(&cappedReader{r: br, left: limit, over: over}, h)); err != nil { @@ -563,31 +614,65 @@ func (m *Manager) UploadContext(ctx context.Context, id, submittedBy string, r i 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) +// 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). +func (m *Manager) storedBytes(ctx context.Context, submittedBy, excludeID string) (user, total int64, err error) { + subs, err := m.Store.ListSubmissions(ctx) if err != nil { - return 0, err + return 0, 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 + return 0, 0, err } - if ok { - total += n + if !ok { + continue + } + total += n + if s.SubmittedBy == submittedBy { + user += n } } - return total, nil + return user, total, nil +} + +// ReapRejected deletes the uploaded context of every submission rejected more +// than olderThan ago, keeping the row. It returns how many blobs it deleted and +// carries on past a blob it cannot delete, reporting the first such error. +func (m *Manager) ReapRejected(ctx context.Context, olderThan time.Duration) (int, error) { + if m.Blobs == nil { + return 0, nil + } + subs, err := m.Store.ListSubmissions(ctx) + if err != nil { + return 0, err + } + cutoff := m.now().Add(-olderThan) + var reaped int + var firstErr error + for _, s := range subs { + if s.Status != StatusRejected || s.ReviewedAt == nil || s.ReviewedAt.After(cutoff) { + continue + } + _, ok, err := m.Blobs.Size(ctx, s.ID) + if err == nil && ok { + if err = m.Blobs.Delete(ctx, s.ID); err == nil { + reaped++ + } + } + if err != nil && firstErr == nil { + firstErr = fmt.Errorf("submit: reap the context of rejected submission %s: %w", s.ID, err) + } + } + return reaped, firstErr } // cappedReader passes through at most left bytes; the first byte beyond the limit diff --git a/internal/submit/submit_test.go b/internal/submit/submit_test.go index 6a19221..de72977 100644 --- a/internal/submit/submit_test.go +++ b/internal/submit/submit_test.go @@ -1167,3 +1167,109 @@ func keysOf(m map[string][]byte) []string { } return out } + +// Every user's contexts together are budgeted too: once they fill it, the next +// upload is refused as a full store (507), whoever makes it and however little +// of their own allowance they used. +func TestUploadContextGlobalBudget(t *testing.T) { + m, _, _ := newManager() + fb := newFakeBlobs() + m.Blobs = fb + m.MaxStoredBytesPerUser = 100 + m.MaxStoredBytesTotal = 10 // gzBody("x") is 5 bytes + ctx := context.Background() + + for _, user := range []string{"user-1", "user-2"} { + sub, err := m.Create(ctx, CreateRequest{DisplayName: "P", SubmittedBy: user}) + if err != nil { + t.Fatal(err) + } + if _, err := m.UploadContext(ctx, sub.ID, user, strings.NewReader(gzBody("x"))); err != nil { + t.Fatalf("%s upload within the total = %v", user, err) + } + } + sub, err := m.Create(ctx, CreateRequest{DisplayName: "P", SubmittedBy: "user-3"}) + if err != nil { + t.Fatal(err) + } + _, err = m.UploadContext(ctx, sub.ID, "user-3", strings.NewReader(gzBody("x"))) + if !errors.Is(err, ErrUploadsFull) || errors.Is(err, ErrQuotaExceeded) { + t.Fatalf("upload into a full store = %v, want ErrUploadsFull", err) + } + if _, ok := fb.stored[sub.ID]; ok { + t.Fatal("a refused upload must persist nothing") + } +} + +type roomyBlobs struct { + *fakeBlobs + err error + need int64 +} + +func (r *roomyBlobs) CheckRoom(need int64) error { r.need = need; return r.err } + +// A store on the node's filesystem is asked for room for the most the upload may +// write before any of it is read. +func TestUploadContextChecksRoom(t *testing.T) { + m, _, _ := newManager() + rb := &roomyBlobs{fakeBlobs: newFakeBlobs(), err: fmt.Errorf("%w: disk", ErrUploadsFull)} + m.Blobs = rb + m.MaxContextBytes = 1000 + ctx := context.Background() + sub, err := m.Create(ctx, CreateRequest{DisplayName: "P", SubmittedBy: "user-1"}) + if err != nil { + t.Fatal(err) + } + if _, err := m.UploadContext(ctx, sub.ID, "user-1", strings.NewReader(gzBody("x"))); !errors.Is(err, ErrUploadsFull) { + t.Fatalf("upload onto a full disk = %v, want ErrUploadsFull", err) + } + if rb.need != 1000 { + t.Fatalf("room asked for %d bytes, want the 1000-byte cap", rb.need) + } + rb.err = nil + if _, err := m.UploadContext(ctx, sub.ID, "user-1", strings.NewReader(gzBody("x"))); err != nil { + t.Fatalf("upload with room = %v", err) + } +} + +// A rejected submission's context goes once the retention has passed; an approved +// one stays (it rebuilds the image after a registry loss), and so do the rows. +func TestReapRejected(t *testing.T) { + m, st, _ := newManager() + fb := newFakeBlobs() + m.Blobs = fb + ctx := context.Background() + ids := map[string]string{} + for _, name := range []string{"old-rejected", "new-rejected", "approved", "pending"} { + sub, err := m.Create(ctx, CreateRequest{DisplayName: name, SubmittedBy: "user-" + name}) + if err != nil { + t.Fatal(err) + } + if _, err := m.UploadContext(ctx, sub.ID, "user-"+name, strings.NewReader(gzBody(name))); err != nil { + t.Fatal(err) + } + ids[name] = sub.ID + } + old, recent := testNow.Add(-8*24*time.Hour), testNow.Add(-time.Hour) + st.subs[ids["old-rejected"]].Status, st.subs[ids["old-rejected"]].ReviewedAt = StatusRejected, &old + st.subs[ids["new-rejected"]].Status, st.subs[ids["new-rejected"]].ReviewedAt = StatusRejected, &recent + st.subs[ids["approved"]].Status, st.subs[ids["approved"]].ReviewedAt = StatusApproved, &old + + n, err := m.ReapRejected(ctx, RejectedContextRetention) + if err != nil || n != 1 { + t.Fatalf("ReapRejected = %d, %v; want 1", n, err) + } + for name, id := range ids { + _, kept := fb.stored[id] + if want := name != "old-rejected"; kept != want { + t.Errorf("%s context kept = %v, want %v", name, kept, want) + } + if _, ok := st.subs[id]; !ok { + t.Errorf("%s row deleted", name) + } + } + if n, err := m.ReapRejected(ctx, RejectedContextRetention); err != nil || n != 0 { + t.Fatalf("second ReapRejected = %d, %v; want 0", n, err) + } +}