feat(submit): 上传存储全局上限、磁盘余量检查,被拒上下文 7 天后回收
This commit is contained in:
8 files changed
+302
-23
No files matched your search
@@ -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.
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
@@ -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"))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
+105
-20
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user