fix(api): make OTP-start throttle atomic to close concurrent-burst bypass
The email-OTP resend cooldown checked the window with a peek (allowed) and only recorded it after delivery. For OTP that throttle is the sole defense and each admitted send is a real, non-idempotent email, so a burst of truly concurrent starts all passed the peek before any recorded and every one mailed: N concurrent starts bombed a mailbox with N codes. Add an atomic reserve/release pair to cooldownLimiter: reserve checks and records the window in one critical section under the mutex, so a concurrent burst yields exactly one winner; release rolls a reservation back only if it is still the current one, so a slow failing caller never clobbers a newer holder. handleEmailOTPStart now reserves both the principal and the recipient key up front and defers a rollback that frees both windows on any mint, create, or delivery error — preserving the old "a failed send does not consume the cooldown" property, now race-free. The wake path keeps allowed→record: its real gate is the running cap and its side effect (SetDesiredState) is idempotent, so the peek gap is harmless there. Tests: a frozen-clock gate-mailer fires 8 concurrent starts for one victim from one principal and asserts exactly one mail and one 202; a flaky-mailer test proves a failed delivery releases the window so an immediate retry in the same instant is admitted.
This commit is contained in:
3 files changed
+182
-8
No files matched your search
@@ -421,6 +421,45 @@ func (c *cooldownLimiter) record(name string) {
|
|||||||
c.last[name] = c.now()
|
c.last[name] = c.now()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// reserve atomically checks name's cooldown AND, if the window is open, records it
|
||||||
|
// in the same critical section, returning the reservation time and true. Unlike
|
||||||
|
// allowed→record there is no gap between the check and the commit, so a burst of
|
||||||
|
// truly concurrent callers yields exactly one winner. Use it where the throttle is
|
||||||
|
// the SOLE defense and each admitted call has a non-idempotent side effect (an OTP
|
||||||
|
// email): an allowed peek would let N goroutines pass together before any records
|
||||||
|
// and bomb a mailbox. The wake path can stay on allowed→record because its real
|
||||||
|
// gate is the running cap and its side effect (SetDesiredState) is idempotent. A
|
||||||
|
// non-positive window disables the throttle (the reservation is a no-op).
|
||||||
|
func (c *cooldownLimiter) reserve(name string, window time.Duration) (time.Time, bool) {
|
||||||
|
if window <= 0 {
|
||||||
|
return time.Time{}, true
|
||||||
|
}
|
||||||
|
c.mu.Lock()
|
||||||
|
defer c.mu.Unlock()
|
||||||
|
if last, ok := c.last[name]; ok && c.now().Sub(last) < window {
|
||||||
|
return time.Time{}, false
|
||||||
|
}
|
||||||
|
t := c.now()
|
||||||
|
c.last[name] = t
|
||||||
|
return t, true
|
||||||
|
}
|
||||||
|
|
||||||
|
// release rolls back a reservation made at reservedAt, but only if it is still the
|
||||||
|
// current one — a later reserve that superseded it is left intact. It lets a caller
|
||||||
|
// undo its hold when a downstream step fails, so a failed mint or delivery never
|
||||||
|
// consumes the window, without a slow failing caller clobbering a newer holder. A
|
||||||
|
// zero reservedAt (a disabled-window reserve) matches nothing and is a no-op.
|
||||||
|
func (c *cooldownLimiter) release(name string, reservedAt time.Time) {
|
||||||
|
if reservedAt.IsZero() {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
c.mu.Lock()
|
||||||
|
defer c.mu.Unlock()
|
||||||
|
if last, ok := c.last[name]; ok && last.Equal(reservedAt) {
|
||||||
|
delete(c.last, name)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ---- running-server cap ----
|
// ---- running-server cap ----
|
||||||
|
|
||||||
// withinRunningCap reports whether waking info's server is allowed under the
|
// withinRunningCap reports whether waking info's server is allowed under the
|
||||||
|
|||||||
@@ -121,16 +121,36 @@ func (a *API) handleEmailOTPStart(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, r, newError(http.StatusBadRequest, "bad_request", "a valid email is required"))
|
writeError(w, r, newError(http.StatusBadRequest, "bad_request", "a valid email is required"))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Throttle sends on both the caller and the recipient before minting anything,
|
// Atomically reserve the cooldown on both the caller and the recipient BEFORE
|
||||||
// so a refused request mints no code and mails nothing. The recipient key is
|
// minting, so a burst of truly concurrent starts yields exactly one winner. Here
|
||||||
// lower-cased so case variants of one address can't sidestep the per-mailbox cap.
|
// the throttle is the sole defense and each admitted send is a real, non-idempotent
|
||||||
|
// email, so an allowed→record peek would let N goroutines slip past together and
|
||||||
|
// bomb a mailbox. The recipient key is lower-cased so case variants of one address
|
||||||
|
// can't sidestep the per-mailbox cap. If any later step fails the deferred rollback
|
||||||
|
// frees both windows, so a failed mint or delivery never consumes the cooldown —
|
||||||
|
// the same property the old record-after-send gave, now race-free.
|
||||||
userKey, emailKey := "user:"+p.UserID, "email:"+strings.ToLower(email)
|
userKey, emailKey := "user:"+p.UserID, "email:"+strings.ToLower(email)
|
||||||
lim := a.otpLimiter()
|
lim := a.otpLimiter()
|
||||||
if !lim.allowed(userKey, otpResendCooldown) || !lim.allowed(emailKey, otpResendCooldown) {
|
userAt, ok := lim.reserve(userKey, otpResendCooldown)
|
||||||
|
if !ok {
|
||||||
writeError(w, r, newError(http.StatusTooManyRequests, "otp_resend_cooldown",
|
writeError(w, r, newError(http.StatusTooManyRequests, "otp_resend_cooldown",
|
||||||
"a code was sent recently; wait a moment before requesting another"))
|
"a code was sent recently; wait a moment before requesting another"))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
emailAt, ok := lim.reserve(emailKey, otpResendCooldown)
|
||||||
|
if !ok {
|
||||||
|
lim.release(userKey, userAt)
|
||||||
|
writeError(w, r, newError(http.StatusTooManyRequests, "otp_resend_cooldown",
|
||||||
|
"a code was sent recently; wait a moment before requesting another"))
|
||||||
|
return
|
||||||
|
}
|
||||||
|
committed := false
|
||||||
|
defer func() {
|
||||||
|
if !committed {
|
||||||
|
lim.release(userKey, userAt)
|
||||||
|
lim.release(emailKey, emailAt)
|
||||||
|
}
|
||||||
|
}()
|
||||||
code, err := newEmailOTP()
|
code, err := newEmailOTP()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
writeError(w, r, err)
|
writeError(w, r, err)
|
||||||
@@ -150,10 +170,9 @@ func (a *API) handleEmailOTPStart(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, r, err)
|
writeError(w, r, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Start both cooldowns only after a code was actually sent: a failed mint or
|
// The send succeeded: keep both reservations (the deferred rollback becomes a
|
||||||
// delivery above must not consume the throttle, mirroring the wake path.
|
// no-op) so the cooldown windows stand.
|
||||||
lim.record(userKey)
|
committed = true
|
||||||
lim.record(emailKey)
|
|
||||||
a.audit(r, auditActor(p), "account.email.otp_sent", "")
|
a.audit(r, auditActor(p), "account.email.otp_sent", "")
|
||||||
writeJSON(w, http.StatusAccepted, map[string]any{
|
writeJSON(w, http.StatusAccepted, map[string]any{
|
||||||
"sent": true,
|
"sent": true,
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
@@ -413,3 +415,117 @@ func TestEmailOTPMailerError(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// flakyMailer fails its first SendOTP, then succeeds, so a test can observe whether
|
||||||
|
// a failed delivery left the cooldown consumed (a retry would be wrongly throttled)
|
||||||
|
// or released (the retry is admitted, as it must be).
|
||||||
|
type flakyMailer struct {
|
||||||
|
calls int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *flakyMailer) SendOTP(_ context.Context, _, _ string) error {
|
||||||
|
m.calls++
|
||||||
|
if m.calls == 1 {
|
||||||
|
return errors.New("smtp down")
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestEmailOTPStartFailedDeliveryReleasesCooldown covers the new reserve→rollback
|
||||||
|
// path: a send that reserves the cooldown but then fails to deliver must release it,
|
||||||
|
// so the very next attempt at the same instant is admitted rather than 429'd. Without
|
||||||
|
// the deferred release a transient SMTP blip would lock a player out for the whole
|
||||||
|
// window — strictly worse than the throttle is meant to be.
|
||||||
|
func TestEmailOTPStartFailedDeliveryReleasesCooldown(t *testing.T) {
|
||||||
|
user := &Principal{UserID: "u1", Email: "[email protected]", Role: "user"}
|
||||||
|
mailer := &flakyMailer{}
|
||||||
|
api := newTestAPI(newFakeRepo(), newFakeCluster())
|
||||||
|
api.External = staticExternal{p: user}
|
||||||
|
api.Mailer = mailer
|
||||||
|
api.Now = func() time.Time { return time.Unix(1_700_000_000, 0) } // frozen: same window
|
||||||
|
eh := api.ExternalHandler()
|
||||||
|
|
||||||
|
if w := do(eh, "POST", "/api/v1/account/email/start", `{"email":"[email protected]"}`, nil); w.Code < 500 {
|
||||||
|
t.Fatalf("first send (mailer fails): code = %d, want 5xx (%s)", w.Code, w.Body.String())
|
||||||
|
}
|
||||||
|
// Same instant, same caller and recipient: had the failed send burned the window
|
||||||
|
// this would be a 429. The rollback frees it, so the retry delivers.
|
||||||
|
if w := do(eh, "POST", "/api/v1/account/email/start", `{"email":"[email protected]"}`, nil); w.Code != http.StatusAccepted {
|
||||||
|
t.Fatalf("retry after failed delivery: code = %d, want 202 (the failed send must release the cooldown) (%s)", w.Code, w.Body.String())
|
||||||
|
}
|
||||||
|
if mailer.calls != 2 {
|
||||||
|
t.Errorf("mailer calls = %d, want 2 (one failed, one delivered)", mailer.calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// gateMailer blocks every SendOTP until all concurrent callers have arrived, making
|
||||||
|
// any check-then-act window in the throttle deterministically observable instead of
|
||||||
|
// scheduler-dependent. It is the committed counterpart of the adversarial burst probe.
|
||||||
|
type gateMailer struct {
|
||||||
|
calls int64
|
||||||
|
entered chan struct{}
|
||||||
|
release chan struct{}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *gateMailer) SendOTP(_ context.Context, _, _ string) error {
|
||||||
|
atomic.AddInt64(&m.calls, 1)
|
||||||
|
m.entered <- struct{}{}
|
||||||
|
<-m.release
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestEmailOTPStartConcurrentBurstBounded is the regression guard for the email-bomb
|
||||||
|
// closure under concurrency. net/http serves each request on its own goroutine, so a
|
||||||
|
// throttle that peeks then records in two steps lets a burst of starts for one victim
|
||||||
|
// from one principal all slip through together. The atomic reserve admits exactly one;
|
||||||
|
// the rest get 429. (Runs without -race — the gate makes the race deterministic.)
|
||||||
|
func TestEmailOTPStartConcurrentBurstBounded(t *testing.T) {
|
||||||
|
const n = 8
|
||||||
|
mailer := &gateMailer{entered: make(chan struct{}, n), release: make(chan struct{})}
|
||||||
|
api := newTestAPI(newFakeRepo(), newFakeCluster())
|
||||||
|
api.External = staticExternal{p: &Principal{UserID: "u1", Email: "[email protected]", Role: "user"}}
|
||||||
|
api.Mailer = mailer
|
||||||
|
api.Now = func() time.Time { return time.Unix(1_700_000_000, 0) } // frozen: one shared window
|
||||||
|
eh := api.ExternalHandler()
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
codes := make([]int, n)
|
||||||
|
for i := 0; i < n; i++ {
|
||||||
|
wg.Add(1)
|
||||||
|
go func(i int) {
|
||||||
|
defer wg.Done()
|
||||||
|
w := do(eh, "POST", "/api/v1/account/email/start", `{"email":"[email protected]"}`, nil)
|
||||||
|
codes[i] = w.Code
|
||||||
|
}(i)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Drain whoever reached the mailer, then release them. A correct throttle admits
|
||||||
|
// exactly one; a buggy one lets several arrive before any records, and they pile
|
||||||
|
// up here within the deadline.
|
||||||
|
deadline := time.After(2 * time.Second)
|
||||||
|
reached := 0
|
||||||
|
loop:
|
||||||
|
for reached < n {
|
||||||
|
select {
|
||||||
|
case <-mailer.entered:
|
||||||
|
reached++
|
||||||
|
case <-deadline:
|
||||||
|
break loop
|
||||||
|
}
|
||||||
|
}
|
||||||
|
close(mailer.release)
|
||||||
|
wg.Wait()
|
||||||
|
|
||||||
|
accepted := 0
|
||||||
|
for _, c := range codes {
|
||||||
|
if c == http.StatusAccepted {
|
||||||
|
accepted++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if got := atomic.LoadInt64(&mailer.calls); got != 1 {
|
||||||
|
t.Errorf("burst delivered %d mails for one victim from one account in one window; want exactly 1 (accepted=%d, reached=%d)", got, accepted, reached)
|
||||||
|
}
|
||||||
|
if accepted != 1 {
|
||||||
|
t.Errorf("burst accepted %d starts; want exactly 1 (the rest must be 429)", accepted)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in new issue
Block a user