feat(submit): give the upload lane a lifecycle — withdraw + admin delete (#76)

Nothing ever removed a submission: users could not retract a pending row, no
route deleted blobs or rows, and the reaper never touches uploads — so every
upload accumulated on the 5 GiB PVC forever and the only cleanup was SQL or
kubectl against the store.

- Blobs.Delete on both transports (local: RemoveAll of the id-namespaced dir,
  id re-validated at the boundary; S3: idempotent object DELETE).
- Store: DeleteSubmission (admin, any status) and DeletePendingSubmission
  (owner+pending CAS — a reviewed row can never be withdrawn out from under
  its build).
- Manager.Delete / Manager.Withdraw delete the ROW first (under the CAS for
  withdraw) and the blob after, so a live row can never point at a reaped
  blob; a cleanup failure names the orphan explicitly instead of failing mute.
- API: DELETE /me/submissions/{id} (withdraw, app tier) and
  DELETE /api/v1/submissions/{id} (admin) both return the row as it was;
  audit events submission.withdraw / submission.delete; openapi documents both
  paths; admin route pinned in the admin-only table.
- Panel: two-step withdraw on a pending row (frees the pending slot and the
  storage budget); two-step delete on every admin row; zh/en copy; wire tests.

Unit: submit (withdraw happy path / wrong owner / reviewed row / no transport /
blob-cleanup failure), local+S3 delete idempotence, api handlers (200/404/409/
503 + route tier); pgint: withdraw CAS + admin delete exactly-once.
go vet/go test/gofmt clean; panel vitest 120 + typecheck green.
This commit is contained in:
Lemon-miaow committed 2026-09-24 10:24:23 +08:00
1 parent ad4d256d8f
commit 854320ac3f
20 files changed
+828 -45

No files matched your search

+14
View File
@@ -149,5 +149,19 @@ func (s *LocalContextStore) Size(_ context.Context, id string) (int64, bool, err
}
}
// Delete removes everything stored for id — the blob plus its id-namespaced
// directory (a stray temp file from an interrupted upload goes with it) — after
// the submission row is gone. Idempotent: nothing stored is success.
func (s *LocalContextStore) Delete(_ context.Context, id string) error {
dir, err := s.dir(id)
if err != nil {
return err
}
if err := os.RemoveAll(dir); err != nil {
return fmt.Errorf("submit: remove context blob: %w", err)
}
return nil
}
// Compile-time proof that the filesystem store satisfies the Blobs transport.
var _ Blobs = (*LocalContextStore)(nil)
+29
View File
@@ -109,6 +109,35 @@ func TestLocalContextStorePutOverwrites(t *testing.T) {
}
}
// Delete removes the blob and its id-namespaced directory, and is idempotent —
// the retry-safety the withdraw/delete cleanup depends on.
func TestLocalContextStoreDelete(t *testing.T) {
base := t.TempDir()
s := &LocalContextStore{Base: base}
ctx := context.Background()
if _, err := s.Put(ctx, "sub-abc", strings.NewReader("\x1f\x8bbytes")); err != nil {
t.Fatalf("Put: %v", err)
}
if err := s.Delete(ctx, "sub-abc"); err != nil {
t.Fatalf("Delete: %v", err)
}
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || ok {
t.Fatalf("Exists after Delete = (%v, %v), want (false, nil)", ok, err)
}
if _, err := os.Stat(filepath.Join(base, "sub-abc")); !os.IsNotExist(err) {
t.Fatalf("per-submission dir still present after Delete (err=%v)", err)
}
// Idempotent: deleting nothing is success, so a retried cleanup cannot fail.
if err := s.Delete(ctx, "sub-abc"); err != nil {
t.Fatalf("second Delete = %v, want nil (idempotent)", err)
}
// The same path guard as Put/Open.
if err := s.Delete(ctx, "../etc"); err == nil {
t.Fatal("Delete must reject an unsafe id")
}
}
func TestLocalContextStoreRejectsUnsafeID(t *testing.T) {
base := t.TempDir()
s := &LocalContextStore{Base: base}
+25 -8
View File
@@ -57,14 +57,16 @@ func (s *PGStore) ListSubmissionsBy(ctx context.Context, submittedBy string) ([]
return s.querySubmissions(ctx, q, submittedBy)
}
// cas executes a compare-and-set UPDATE and reports whether THIS call moved the
// row. The `status = 'pending_review'` guard is the actual CAS predicate and is
// deliberately kept INLINE in each caller's query — it is security-visible, so a
// reader auditing "can a non-pending row be flipped?" must see it next to the SET.
// cas only folds the shared ExecContext + RowsAffected tail so the two reviewers
// (approve, reject) cannot drift in how they report a lost race or a RowsAffected
// error. n == 0 means a concurrent review already won the row — reported as
// won=false (never an error), which the Manager maps to ErrAlreadyReviewed.
// cas executes a single-statement compare-and-set — an UPDATE or DELETE whose
// WHERE clause is the predicate — and reports whether THIS call moved a row. The
// `status = 'pending_review'` guard is the actual CAS predicate in each caller's
// query and is deliberately kept INLINE — it is security-visible, so a reader
// auditing "can a non-pending row be flipped or deleted?" must see it next to the
// SET or DELETE. cas only folds the shared ExecContext + RowsAffected tail so
// the reviewers and deleters cannot drift in how they report a lost race or a
// RowsAffected error. n == 0 means a concurrent actor already won the row —
// reported as won=false (never an error), which the Manager maps to
// ErrAlreadyReviewed / ErrNotFound.
func (s *PGStore) cas(ctx context.Context, q string, args ...any) (bool, error) {
res, err := s.db.ExecContext(ctx, q, args...)
if err != nil {
@@ -107,6 +109,21 @@ func (s *PGStore) LinkBuild(ctx context.Context, id, buildID string) error {
return nil
}
// DeleteSubmission is the admin delete: any status, one row, no predicate beyond
// the id. False means the id was already gone (a concurrent delete won).
func (s *PGStore) DeleteSubmission(ctx context.Context, id string) (bool, error) {
const q = `DELETE FROM image_submissions WHERE id = $1`
return s.cas(ctx, q, id)
}
// DeletePendingSubmission is the withdraw CAS: owner + pending_review must both
// still hold, so a reviewed submission can never be deleted through this path.
func (s *PGStore) DeletePendingSubmission(ctx context.Context, id, submittedBy string) (bool, error) {
const q = `DELETE FROM image_submissions
WHERE id = $1 AND submitted_by = $2 AND status = 'pending_review'`
return s.cas(ctx, q, id, submittedBy)
}
func (s *PGStore) querySubmissions(ctx context.Context, q string, args ...any) ([]Submission, error) {
rows, err := s.db.QueryContext(ctx, q, args...)
if err != nil {
+15
View File
@@ -21,6 +21,7 @@ type s3Client interface {
PutObject(ctx context.Context, bucket, object string, reader io.Reader, size int64, opts minio.PutObjectOptions) (minio.UploadInfo, error)
StatObject(ctx context.Context, bucket, object string, opts minio.StatObjectOptions) (minio.ObjectInfo, error)
GetObject(ctx context.Context, bucket, object string, opts minio.GetObjectOptions) (s3Object, error)
RemoveObject(ctx context.Context, bucket, object string, opts minio.RemoveObjectOptions) error
}
// s3Object is the handle GetObject yields: a stream whose Stat performs the HEAD
@@ -205,6 +206,20 @@ func (s *S3ContextStore) Size(ctx context.Context, id string) (int64, bool, erro
return info.Size, true, nil
}
// Delete removes the stored object for the withdrawn/deleted submission. S3's
// DELETE is idempotent — removing an absent key succeeds — which is exactly the
// contract the cleanup path needs on a retry.
func (s *S3ContextStore) Delete(ctx context.Context, id string) error {
key, err := s.keyFor(id)
if err != nil {
return err
}
if err := s.client.RemoveObject(ctx, s.bucket, key, minio.RemoveObjectOptions{}); err != nil {
return fmt.Errorf("submit: remove context blob: %w", err)
}
return 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
+40 -3
View File
@@ -15,9 +15,10 @@ import (
// error for a missing stat, so S3ContextStore's key derivation and not-found
// handling are exercised without a live bucket.
type fakeS3 struct {
objects map[string][]byte
putErr error
statErr error // when set, StatObject returns it (e.g. auth rejected / bucket missing)
objects map[string][]byte
putErr error
statErr error // when set, StatObject returns it (e.g. auth rejected / bucket missing)
removeErr error
}
func (f *fakeS3) PutObject(_ context.Context, bucket, object string, r io.Reader, _ int64, _ minio.PutObjectOptions) (minio.UploadInfo, error) {
@@ -45,6 +46,16 @@ func (f *fakeS3) StatObject(_ context.Context, bucket, object string, _ minio.St
return minio.ObjectInfo{}, minio.ErrorResponse{Code: "NoSuchKey", StatusCode: http.StatusNotFound}
}
// RemoveObject mirrors S3's idempotent DELETE: removing a key (absent or not)
// succeeds unless removeErr injects a failure.
func (f *fakeS3) RemoveObject(_ context.Context, bucket, object string, _ minio.RemoveObjectOptions) error {
if f.removeErr != nil {
return f.removeErr
}
delete(f.objects, bucket+"/"+object)
return nil
}
// fakeS3Object is the object handle fakeS3.GetObject yields: Stat mirrors
// StatObject's not-found behaviour, Read serves the stored bytes.
type fakeS3Object struct {
@@ -142,6 +153,32 @@ func TestS3ContextStorePutAndExists(t *testing.T) {
}
}
// Delete removes the derived object and is idempotent (S3 DELETE of an absent
// key succeeds), so a retried cleanup after a partial failure cannot stick.
func TestS3ContextStoreDelete(t *testing.T) {
fake := &fakeS3{}
s := &S3ContextStore{client: fake, bucket: "felis-uploads", prefix: "builds"}
ctx := context.Background()
if _, err := s.Put(ctx, "sub-abc", strings.NewReader("\x1f\x8bbytes")); err != nil {
t.Fatalf("Put: %v", err)
}
if err := s.Delete(ctx, "sub-abc"); err != nil {
t.Fatalf("Delete: %v", err)
}
if ok, err := s.Exists(ctx, "sub-abc"); err != nil || ok {
t.Fatalf("Exists after Delete = (%v, %v), want (false, nil)", ok, err)
}
if err := s.Delete(ctx, "sub-abc"); err != nil {
t.Fatalf("second Delete = %v, want nil (idempotent)", err)
}
// A transport failure is surfaced, not swallowed.
fake.removeErr = errors.New("s3 unavailable")
if err := s.Delete(ctx, "sub-abc"); err == nil {
t.Fatal("Delete with a failing transport = nil, want error")
}
}
// Open serves the stored object's bytes and maps a missing key to ErrBlobNotFound
// (the internal fetch route's 404), eagerly — before the caller reads a byte.
func TestS3ContextStoreOpen(t *testing.T) {
+96
View File
@@ -203,6 +203,17 @@ type Store interface {
// (Approve surfaces that distinctly so remediation does not double-build; see
// the Approve ordering note and the package KNOWN-LIMITATION).
LinkBuild(ctx context.Context, id, buildID string) error
// DeleteSubmission removes a submission row outright — the admin delete path.
// Any status is deletable: the row is the business record and the caller has
// chosen to retire it. Reports whether a row was actually deleted; false
// means a concurrent delete won, which the Manager surfaces as ErrNotFound.
DeleteSubmission(ctx context.Context, id string) (bool, error)
// DeletePendingSubmission is the submitter's withdraw CAS: it deletes the row
// only while it is still owned by submittedBy AND still pending_review, so a
// concurrent approve/reject wins or loses cleanly and a reviewed submission
// can never be withdrawn out from under its build. Reports whether THIS call
// deleted the row.
DeletePendingSubmission(ctx context.Context, id, submittedBy string) (bool, error)
}
// Builds is the slice of the build subsystem the approval lane drives. An
@@ -235,6 +246,11 @@ type Blobs interface {
// 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)
// Delete removes everything stored for id — the withdrawal/review-cleanup
// path. It is idempotent: deleting nothing is success, so a retried cleanup
// never fails on absence. Callers reap the blob only AFTER the row is gone
// (see Manager.deleteBlob), so a live row can never point at a reaped blob.
Delete(ctx context.Context, id string) 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
@@ -711,6 +727,86 @@ func (m *Manager) Reject(ctx context.Context, id, reviewedBy, reason string) (*S
return sub, nil
}
// Withdraw retracts the submitter's own pending submission: the row is deleted
// (CAS-guarded on owner + pending_review, so a concurrent review can never be
// undercut) and its uploaded context is reaped, freeing both the user's pending
// slot and their storage budget for a fresh submission. Once reviewed the
// context is frozen — an approved build may already be consuming it — so a
// non-pending row reports ErrAlreadyReviewed, exactly like UploadContext, and
// another user's id stays invisible (ErrNotFound), like every other owner-scoped
// operation on the lane.
func (m *Manager) Withdraw(ctx context.Context, id, submittedBy string) (*Submission, error) {
if strings.TrimSpace(submittedBy) == "" {
return nil, invalidf("submitter identity is required")
}
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
}
won, err := m.Store.DeletePendingSubmission(ctx, id, submittedBy)
if err != nil {
return nil, err
}
if !won {
// A concurrent approve/reject (or delete) moved the row between the load
// and the CAS. The context must NOT be reaped: it belongs to the review
// outcome now.
return nil, ErrAlreadyReviewed
}
if err := m.deleteBlob(ctx, id); err != nil {
return nil, err
}
return sub, nil
}
// Delete retires any submission outright (the admin path): the row and its
// stored context are both removed, any status. The reviewer identity is recorded
// by the API's audit event, not on the row — the row no longer exists. Deleting
// an approved submission whose build is still running can fail that build (its
// context fetch answers 404); the admin has explicitly chosen to retire the
// artifact. A concurrent delete reports ErrNotFound, the same outcome the
// initial load would have given had it lost the race.
func (m *Manager) Delete(ctx context.Context, id string) (*Submission, error) {
sub, err := m.Store.GetSubmission(ctx, id)
if err != nil {
return nil, err
}
deleted, err := m.Store.DeleteSubmission(ctx, id)
if err != nil {
return nil, err
}
if !deleted {
return nil, ErrNotFound
}
if err := m.deleteBlob(ctx, id); err != nil {
return nil, err
}
return sub, nil
}
// deleteBlob reaps a submission's stored context after its row is gone. The row
// is deleted FIRST (for Withdraw, under the status CAS), so a cleanup failure
// here can never leave a live row pointing at a reaped blob — the failure mode
// is the other direction: the row is retired and the blob is orphaned on the
// uploads store, which the returned error names explicitly so the operator knows
// 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.Blobs == nil {
return nil
}
if err := m.Blobs.Delete(ctx, id); err != nil {
return fmt.Errorf("submit: submission removed, but its uploaded context could not be deleted (it may remain on the uploads store): %w", err)
}
return nil
}
// List returns every submission, newest first (the admin review queue).
func (m *Manager) List(ctx context.Context) ([]Submission, error) {
return m.Store.ListSubmissions(ctx)
+183
View File
@@ -25,6 +25,7 @@ type fakeBlobs struct {
putErr error
existsErr error
sizeErr error
deleteErr error
forceExists *bool // overrides the stored-map lookup for the approve-gate tests
}
@@ -65,6 +66,15 @@ func (f *fakeBlobs) Size(_ context.Context, id string) (int64, bool, error) {
return int64(len(b)), true, nil
}
// Delete mirrors the real stores: idempotent, nothing stored is success.
func (f *fakeBlobs) Delete(_ context.Context, id string) error {
if f.deleteErr != nil {
return f.deleteErr
}
delete(f.stored, id)
return nil
}
func (f *fakeBlobs) Open(_ context.Context, id string) (io.ReadCloser, error) {
b, ok := f.stored[id]
if !ok {
@@ -90,6 +100,7 @@ type fakeStore struct {
approveErr error
rejectErr error
linkErr error
deleteErr error
linked []string // "id=buildID" recorder
}
@@ -190,6 +201,29 @@ func (f *fakeStore) LinkBuild(_ context.Context, id, buildID string) error {
return nil
}
func (f *fakeStore) DeleteSubmission(_ context.Context, id string) (bool, error) {
if f.deleteErr != nil {
return false, f.deleteErr
}
if _, ok := f.subs[id]; !ok {
return false, nil // concurrent delete won / nonexistent
}
delete(f.subs, id)
return true, nil
}
func (f *fakeStore) DeletePendingSubmission(_ context.Context, id, by string) (bool, error) {
if f.deleteErr != nil {
return false, f.deleteErr
}
s, ok := f.subs[id]
if !ok || s.SubmittedBy != by || s.Status != StatusPendingReview {
return false, nil
}
delete(f.subs, id)
return true, nil
}
// fakeBuilds records Submit calls and can inject a failure.
type fakeBuilds struct {
submitErr error
@@ -819,6 +853,155 @@ func TestCreatePendingCountFailureSurfaces(t *testing.T) {
}
}
// Withdraw retracts the submitter's own pending submission: the row AND its
// uploaded context are gone — which is what frees the pending slot and the
// storage budget for a fresh submission.
func TestWithdrawDeletesRowAndBlob(t *testing.T) {
m, st, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(ctx, seed.ID, "user-1", strings.NewReader(gzBody("bytes"))); err != nil {
t.Fatalf("UploadContext: %v", err)
}
sub, err := m.Withdraw(ctx, seed.ID, "user-1")
if err != nil {
t.Fatalf("Withdraw: %v", err)
}
if sub.ID != seed.ID || sub.Status != StatusPendingReview {
t.Fatalf("withdrawn row = %+v, want the pending row back", sub)
}
if _, ok := st.subs[seed.ID]; ok {
t.Fatal("withdraw must delete the row")
}
if _, ok := fb.stored[seed.ID]; ok {
t.Fatal("withdraw must reap the uploaded context")
}
}
// Another user's id is invisible on the withdraw path (404, not 403), and
// nothing is touched.
func TestWithdrawNotOwnerIsNotFound(t *testing.T) {
m, st, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(ctx, seed.ID, "user-1", strings.NewReader(gzBody("x"))); err != nil {
t.Fatalf("UploadContext: %v", err)
}
if _, err := m.Withdraw(ctx, seed.ID, "user-2"); !errors.Is(err, ErrNotFound) {
t.Fatalf("err = %v, want ErrNotFound", err)
}
if _, ok := st.subs[seed.ID]; !ok {
t.Fatal("a non-owner withdraw must not delete the row")
}
if _, ok := fb.stored[seed.ID]; !ok {
t.Fatal("a non-owner withdraw must not reap the context")
}
}
// A reviewed submission is frozen: the withdraw CAS refuses it (409), keeping
// the context a running/queued build may still consume.
func TestWithdrawReviewedIsAlreadyReviewed(t *testing.T) {
m, st, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(ctx, seed.ID, "user-1", strings.NewReader(gzBody("x"))); err != nil {
t.Fatalf("UploadContext: %v", err)
}
if _, err := m.Reject(ctx, seed.ID, "[email protected]", "no"); err != nil {
t.Fatalf("Reject: %v", err)
}
if _, err := m.Withdraw(ctx, seed.ID, "user-1"); !errors.Is(err, ErrAlreadyReviewed) {
t.Fatalf("err = %v, want ErrAlreadyReviewed", err)
}
if _, ok := st.subs[seed.ID]; !ok {
t.Fatal("a reviewed row must survive a withdraw attempt")
}
if _, ok := fb.stored[seed.ID]; !ok {
t.Fatal("a reviewed row's context must survive a withdraw attempt")
}
}
// The admin delete retires ANY status, reaping the context with it.
func TestDeleteAnyStatusReapsContext(t *testing.T) {
m, st, _ := newManager()
fb := newFakeBlobs()
m.Blobs = fb
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.UploadContext(ctx, seed.ID, "user-1", strings.NewReader(gzBody("x"))); err != nil {
t.Fatalf("UploadContext: %v", err)
}
if _, err := m.Reject(ctx, seed.ID, "[email protected]", "no"); err != nil {
t.Fatalf("Reject: %v", err)
}
sub, err := m.Delete(ctx, seed.ID)
if err != nil {
t.Fatalf("Delete: %v", err)
}
if sub.Status != StatusRejected {
t.Fatalf("returned status = %q, want the row as it was before the delete", sub.Status)
}
if _, ok := st.subs[seed.ID]; ok {
t.Fatal("delete must remove the row")
}
if _, ok := fb.stored[seed.ID]; ok {
t.Fatal("delete must reap the context")
}
}
// Deleting an unknown id — including the second delete of one already gone —
// is ErrNotFound (404), not a crash and not a silent success.
func TestDeleteUnknownIsNotFound(t *testing.T) {
m, _, _ := newManager()
m.Blobs = newFakeBlobs()
if _, err := m.Delete(context.Background(), "sub-nope"); !errors.Is(err, ErrNotFound) {
t.Fatalf("err = %v, want ErrNotFound", err)
}
}
// A deployment without an upload transport has no blobs to reap: delete still
// retires the row.
func TestDeleteWithoutTransport(t *testing.T) {
m, st, _ := newManager() // Blobs nil
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
if _, err := m.Delete(ctx, seed.ID); err != nil {
t.Fatalf("Delete: %v", err)
}
if _, ok := st.subs[seed.ID]; ok {
t.Fatal("delete must remove the row even with no transport")
}
}
// A blob-cleanup failure after the row is gone must surface loudly (naming what
// was left behind), never masquerade as success.
func TestDeleteBlobCleanupFailureSurfaces(t *testing.T) {
m, st, _ := newManager()
fb := newFakeBlobs()
fb.deleteErr = errors.New("pvc read-only")
m.Blobs = fb
ctx := context.Background()
seed, _ := m.Create(ctx, CreateRequest{DisplayName: "Pack", SubmittedBy: "user-1"})
_, err := m.Delete(ctx, seed.ID)
if err == nil || !strings.Contains(err.Error(), "could not be deleted") {
t.Fatalf("err = %v, want the cleanup failure named", err)
}
if _, ok := st.subs[seed.ID]; ok {
t.Fatal("the row is deleted before the blob reap; it must stay deleted")
}
}
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.