fix(submit): 待审上限在事务内加咨询锁原子检查,跨副本不超额
This commit is contained in:
5 files changed
+114
-40
No files matched your search
+11
-2
@@ -95,10 +95,19 @@ worth revisiting.
|
|||||||
moved inside `ClaimServer` (advisory lock + re-check + UPDATE in one transaction),
|
moved inside `ClaimServer` (advisory lock + re-check + UPDATE in one transaction),
|
||||||
red-then-green in the pgint suite, which is exactly the real-Postgres harness this
|
red-then-green in the pgint suite, which is exactly the real-Postgres harness this
|
||||||
line was waiting for.
|
line was waiting for.
|
||||||
- `internal/api/api.go:671` — `cooldownLimiter` is process-local, so across N api
|
- `internal/api/api.go:773` — `cooldownLimiter` is process-local, so across N api
|
||||||
replicas a caller could draw up to N OTP codes per window. The intra-replica burst
|
replicas a caller could draw up to N OTP codes per window. The intra-replica burst
|
||||||
is closed; cross-replica bounding needs a shared store, out of scope for a
|
is closed; cross-replica bounding needs a shared store, out of scope for a
|
||||||
single-replica install.
|
single-replica install. Revisit before the api Deployment runs more than one
|
||||||
|
replica.
|
||||||
|
- `internal/submit/submit.go:524` — the per-user upload storage budget reads the
|
||||||
|
stored bytes, then writes. On one replica the API's per-user upload reservation
|
||||||
|
serializes it; across replicas a burst can overshoot by one blob per interleaved
|
||||||
|
upload, each still under the single-blob cap. The pending-submission cap no
|
||||||
|
longer has this shape: `CreateSubmission` counts and inserts under a
|
||||||
|
per-submitter advisory lock (pgint `TestSubmitPendingCapHoldsUnderConcurrency`).
|
||||||
|
Revisit with the cooldown above, before scaling api replicas: a reservation row
|
||||||
|
per upload in the same kind of transaction closes it.
|
||||||
- `internal/submit/submit.go:436` and `internal/submit/submit_test.go:351` — a
|
- `internal/submit/submit.go:436` and `internal/submit/submit_test.go:351` — a
|
||||||
post-CAS `Approve`
|
post-CAS `Approve`
|
||||||
failure leaves a row indistinguishable from the benign case, so `Approve` returns a
|
failure leaves a row indistinguishable from the benign case, so `Approve` returns a
|
||||||
|
|||||||
@@ -1109,6 +1109,48 @@ func TestRedeemPlayerBindCodeContract(t *testing.T) {
|
|||||||
|
|
||||||
// ---- submissions ---------------------------------------------------------------
|
// ---- submissions ---------------------------------------------------------------
|
||||||
|
|
||||||
|
// TestSubmitPendingCapHoldsUnderConcurrency: parallel creates by one user,
|
||||||
|
// through separate connections as separate replicas would make them, admit
|
||||||
|
// exactly the cap (build-supply-chain-16).
|
||||||
|
func TestSubmitPendingCapHoldsUnderConcurrency(t *testing.T) {
|
||||||
|
ctx := context.Background()
|
||||||
|
u := newUser(t, "user", "subcap")
|
||||||
|
s := submit.NewPGStore(db)
|
||||||
|
now := mustNow()
|
||||||
|
const limit, tries = 3, 12
|
||||||
|
sfx := suffix(t)
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
var mu sync.Mutex
|
||||||
|
admitted, refused := 0, 0
|
||||||
|
for i := 0; i < tries; i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(i int) {
|
||||||
|
defer wg.Done()
|
||||||
|
id := fmt.Sprintf("sub-cap-%d-%s", i, sfx)
|
||||||
|
_, err := s.CreateSubmission(ctx, &submit.Submission{ID: id, SubmittedBy: u.ID, DisplayName: "cap",
|
||||||
|
ContextRef: "s3://bucket/" + id, Status: submit.StatusPendingReview, CreatedAt: now}, limit)
|
||||||
|
mu.Lock()
|
||||||
|
defer mu.Unlock()
|
||||||
|
switch {
|
||||||
|
case err == nil:
|
||||||
|
admitted++
|
||||||
|
case errors.Is(err, submit.ErrQuotaExceeded):
|
||||||
|
refused++
|
||||||
|
default:
|
||||||
|
t.Errorf("create %d: %v", i, err)
|
||||||
|
}
|
||||||
|
}(i)
|
||||||
|
}
|
||||||
|
wg.Wait()
|
||||||
|
if admitted != limit || refused != tries-limit {
|
||||||
|
t.Fatalf("admitted %d, refused %d; want %d and %d", admitted, refused, limit, tries-limit)
|
||||||
|
}
|
||||||
|
if n, err := s.CountPendingSubmissionsBy(ctx, u.ID); err != nil || n != limit {
|
||||||
|
t.Fatalf("pending = (%d, %v), want %d", n, err, limit)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
func TestSubmitStoreContract(t *testing.T) {
|
func TestSubmitStoreContract(t *testing.T) {
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
u := newUser(t, "user", "sub")
|
u := newUser(t, "user", "sub")
|
||||||
@@ -1121,7 +1163,7 @@ func TestSubmitStoreContract(t *testing.T) {
|
|||||||
ContextRef: "s3://bucket/" + id + "/context.tar.gz",
|
ContextRef: "s3://bucket/" + id + "/context.tar.gz",
|
||||||
Status: submit.StatusPendingReview, CreatedAt: now,
|
Status: submit.StatusPendingReview, CreatedAt: now,
|
||||||
}
|
}
|
||||||
if err := s.CreateSubmission(ctx, sub); err != nil {
|
if _, err := s.CreateSubmission(ctx, sub, 5); err != nil {
|
||||||
t.Fatalf("CreateSubmission: %v", err)
|
t.Fatalf("CreateSubmission: %v", err)
|
||||||
}
|
}
|
||||||
// The pending-queue count is the read behind the per-user pending cap: a
|
// The pending-queue count is the read behind the per-user pending cap: a
|
||||||
@@ -1132,6 +1174,14 @@ func TestSubmitStoreContract(t *testing.T) {
|
|||||||
if n, err := s.CountPendingSubmissionsBy(ctx, u.ID+"-nobody"); err != nil || n != 0 {
|
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)
|
t.Fatalf("CountPendingSubmissionsBy for an unknown user = (%d, %v), want (0, nil)", n, err)
|
||||||
}
|
}
|
||||||
|
over := &submit.Submission{ID: id + "-over", SubmittedBy: u.ID, DisplayName: "over the cap",
|
||||||
|
ContextRef: "s3://bucket/over", Status: submit.StatusPendingReview, CreatedAt: now}
|
||||||
|
if n, err := s.CreateSubmission(ctx, over, 1); !errors.Is(err, submit.ErrQuotaExceeded) || n != 1 {
|
||||||
|
t.Fatalf("create at the cap = (%d, %v), want (1, ErrQuotaExceeded)", n, err)
|
||||||
|
}
|
||||||
|
if _, err := s.GetSubmission(ctx, over.ID); !errors.Is(err, submit.ErrNotFound) {
|
||||||
|
t.Fatalf("a refused create left a row: %v", err)
|
||||||
|
}
|
||||||
got, err := s.GetSubmission(ctx, id)
|
got, err := s.GetSubmission(ctx, id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("GetSubmission: %v", err)
|
t.Fatalf("GetSubmission: %v", err)
|
||||||
@@ -1200,11 +1250,11 @@ func TestSubmitStoreContract(t *testing.T) {
|
|||||||
// wrong owner or a reviewed row can never delete through it — and the admin
|
// wrong owner or a reviewed row can never delete through it — and the admin
|
||||||
// path deletes any status, exactly once.
|
// path deletes any status, exactly once.
|
||||||
id2 := "sub-w-" + suffix(t)
|
id2 := "sub-w-" + suffix(t)
|
||||||
if err := s.CreateSubmission(ctx, &submit.Submission{
|
if _, err := s.CreateSubmission(ctx, &submit.Submission{
|
||||||
ID: id2, SubmittedBy: u.ID, DisplayName: "withdraw me",
|
ID: id2, SubmittedBy: u.ID, DisplayName: "withdraw me",
|
||||||
ContextRef: "s3://bucket/" + id2 + "/context.tar.gz",
|
ContextRef: "s3://bucket/" + id2 + "/context.tar.gz",
|
||||||
Status: submit.StatusPendingReview, CreatedAt: now,
|
Status: submit.StatusPendingReview, CreatedAt: now,
|
||||||
}); err != nil {
|
}, 5); err != nil {
|
||||||
t.Fatalf("CreateSubmission(2): %v", err)
|
t.Fatalf("CreateSubmission(2): %v", err)
|
||||||
}
|
}
|
||||||
if ok, err := s.DeletePendingSubmission(ctx, id2, "someone-else"); err != nil || ok {
|
if ok, err := s.DeletePendingSubmission(ctx, id2, "someone-else"); err != nil || ok {
|
||||||
|
|||||||
@@ -24,13 +24,35 @@ var _ Store = (*PGStore)(nil)
|
|||||||
const submissionColumns = `id, submitted_by, display_name, context_ref, status,
|
const submissionColumns = `id, submitted_by, display_name, context_ref, status,
|
||||||
image_ref, build_id, reviewed_by, reject_reason, created_at, reviewed_at, context_sha256`
|
image_ref, build_id, reviewed_by, reject_reason, created_at, reviewed_at, context_sha256`
|
||||||
|
|
||||||
func (s *PGStore) CreateSubmission(ctx context.Context, sub *Submission) error {
|
// CreateSubmission counts and inserts in one transaction under a per-submitter
|
||||||
|
// advisory lock, so two creates by the same user, on one replica or several,
|
||||||
|
// serialize and the second sees the first's row. A hashtext collision between
|
||||||
|
// two users only serializes their creates.
|
||||||
|
func (s *PGStore) CreateSubmission(ctx context.Context, sub *Submission, maxPending int) (int, error) {
|
||||||
|
tx, err := s.db.BeginTx(ctx, nil)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
defer func() { _ = tx.Rollback() }()
|
||||||
|
if _, err := tx.ExecContext(ctx, `SELECT pg_advisory_xact_lock(hashtext('submission:' || $1))`, sub.SubmittedBy); err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
var pending int
|
||||||
|
if err := tx.QueryRowContext(ctx, `SELECT count(*) FROM image_submissions
|
||||||
|
WHERE submitted_by = $1 AND status = 'pending_review'`, sub.SubmittedBy).Scan(&pending); err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
if pending >= maxPending {
|
||||||
|
return pending, ErrQuotaExceeded
|
||||||
|
}
|
||||||
const q = `INSERT INTO image_submissions
|
const q = `INSERT INTO image_submissions
|
||||||
(id, submitted_by, display_name, context_ref, status, created_at)
|
(id, submitted_by, display_name, context_ref, status, created_at)
|
||||||
VALUES ($1, $2, $3, $4, $5, $6)`
|
VALUES ($1, $2, $3, $4, $5, $6)`
|
||||||
_, err := s.db.ExecContext(ctx, q,
|
if _, err := tx.ExecContext(ctx, q,
|
||||||
sub.ID, sub.SubmittedBy, sub.DisplayName, sub.ContextRef, string(sub.Status), sub.CreatedAt)
|
sub.ID, sub.SubmittedBy, sub.DisplayName, sub.ContextRef, string(sub.Status), sub.CreatedAt); err != nil {
|
||||||
return err
|
return pending, err
|
||||||
|
}
|
||||||
|
return pending, tx.Commit()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *PGStore) CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error) {
|
func (s *PGStore) CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error) {
|
||||||
|
|||||||
+13
-19
@@ -187,12 +187,12 @@ type Submission struct {
|
|||||||
// interface so the Manager is tested against an in-memory fake; the Postgres
|
// interface so the Manager is tested against an in-memory fake; the Postgres
|
||||||
// implementation (PGStore) is integration-tested only.
|
// implementation (PGStore) is integration-tested only.
|
||||||
type Store interface {
|
type Store interface {
|
||||||
// CreateSubmission inserts a pending_review row.
|
// CreateSubmission inserts a pending_review row unless its submitter already
|
||||||
CreateSubmission(ctx context.Context, s *Submission) error
|
// has maxPending rows pending_review, in which case nothing is written and
|
||||||
// CountPendingSubmissionsBy reports how many of one user's submissions are
|
// the error is ErrQuotaExceeded. It returns how many were pending before
|
||||||
// still pending_review — the queue-length read behind the per-user pending
|
// the insert. The count and the insert are one decision, so concurrent
|
||||||
// cap in Create.
|
// creates on any number of api replicas never land a row over the cap.
|
||||||
CountPendingSubmissionsBy(ctx context.Context, submittedBy string) (int, error)
|
CreateSubmission(ctx context.Context, s *Submission, maxPending int) (int, error)
|
||||||
// GetSubmission loads one submission, or ErrNotFound.
|
// GetSubmission loads one submission, or ErrNotFound.
|
||||||
GetSubmission(ctx context.Context, id string) (*Submission, error)
|
GetSubmission(ctx context.Context, id string) (*Submission, error)
|
||||||
// ListSubmissions returns every submission, newest first (admin queue).
|
// ListSubmissions returns every submission, newest first (admin queue).
|
||||||
@@ -449,18 +449,6 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e
|
|||||||
return nil, invalidf("submitter identity is required")
|
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()
|
id := m.newID()
|
||||||
s := &Submission{
|
s := &Submission{
|
||||||
ID: id,
|
ID: id,
|
||||||
@@ -470,7 +458,13 @@ func (m *Manager) Create(ctx context.Context, req CreateRequest) (*Submission, e
|
|||||||
Status: StatusPendingReview,
|
Status: StatusPendingReview,
|
||||||
CreatedAt: m.now(),
|
CreatedAt: m.now(),
|
||||||
}
|
}
|
||||||
if err := m.Store.CreateSubmission(ctx, s); err != nil {
|
// 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
|
||||||
|
// store checks and inserts as one step.
|
||||||
|
if pending, err := m.Store.CreateSubmission(ctx, s, m.maxPendingPerUser()); errors.Is(err, ErrQuotaExceeded) {
|
||||||
|
return nil, fmt.Errorf("%w: %d submissions are already awaiting review (limit %d)",
|
||||||
|
ErrQuotaExceeded, pending, m.maxPendingPerUser())
|
||||||
|
} else if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return s, nil
|
return s, nil
|
||||||
|
|||||||
@@ -113,25 +113,24 @@ type fakeStore struct {
|
|||||||
|
|
||||||
func newFakeStore() *fakeStore { return &fakeStore{subs: map[string]*Submission{}} }
|
func newFakeStore() *fakeStore { return &fakeStore{subs: map[string]*Submission{}} }
|
||||||
|
|
||||||
func (f *fakeStore) CreateSubmission(_ context.Context, s *Submission) error {
|
func (f *fakeStore) CreateSubmission(_ context.Context, s *Submission, maxPending int) (int, error) {
|
||||||
if f.createErr != nil {
|
|
||||||
return f.createErr
|
|
||||||
}
|
|
||||||
cp := *s
|
|
||||||
f.subs[s.ID] = &cp
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (f *fakeStore) CountPendingSubmissionsBy(_ context.Context, by string) (int, error) {
|
|
||||||
if f.countErr != nil {
|
if f.countErr != nil {
|
||||||
return 0, f.countErr
|
return 0, f.countErr
|
||||||
}
|
}
|
||||||
var n int
|
var n int
|
||||||
for _, s := range f.subs {
|
for _, x := range f.subs {
|
||||||
if s.SubmittedBy == by && s.Status == StatusPendingReview {
|
if x.SubmittedBy == s.SubmittedBy && x.Status == StatusPendingReview {
|
||||||
n++
|
n++
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if n >= maxPending {
|
||||||
|
return n, ErrQuotaExceeded
|
||||||
|
}
|
||||||
|
if f.createErr != nil {
|
||||||
|
return n, f.createErr
|
||||||
|
}
|
||||||
|
cp := *s
|
||||||
|
f.subs[s.ID] = &cp
|
||||||
return n, nil
|
return n, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user