From 879b1777f5668f259ec0ba54a08618c3f9d52de2 Mon Sep 17 00:00:00 2001 From: Minseong Choi Date: Tue, 30 Jun 2026 23:04:07 +0900 Subject: [PATCH] fix(api): make OTP-start throttle atomic to close concurrent-burst bypass MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- internal/api/api.go | 39 ++++++++ internal/api/handlers_email_otp.go | 35 +++++-- internal/api/handlers_email_otp_test.go | 116 ++++++++++++++++++++++++ 3 files changed, 182 insertions(+), 8 deletions(-) diff --git a/internal/api/api.go b/internal/api/api.go index 92ac72d..dfbc00f 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -421,6 +421,45 @@ func (c *cooldownLimiter) record(name string) { 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 ---- // withinRunningCap reports whether waking info's server is allowed under the diff --git a/internal/api/handlers_email_otp.go b/internal/api/handlers_email_otp.go index 26ddb6e..81bf074 100644 --- a/internal/api/handlers_email_otp.go +++ b/internal/api/handlers_email_otp.go @@ -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")) return } - // Throttle sends on both the caller and the recipient before minting anything, - // so a refused request mints no code and mails nothing. The recipient key is - // lower-cased so case variants of one address can't sidestep the per-mailbox cap. + // Atomically reserve the cooldown on both the caller and the recipient BEFORE + // minting, so a burst of truly concurrent starts yields exactly one winner. Here + // 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) 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", "a code was sent recently; wait a moment before requesting another")) 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() if err != nil { writeError(w, r, err) @@ -150,10 +170,9 @@ func (a *API) handleEmailOTPStart(w http.ResponseWriter, r *http.Request) { writeError(w, r, err) return } - // Start both cooldowns only after a code was actually sent: a failed mint or - // delivery above must not consume the throttle, mirroring the wake path. - lim.record(userKey) - lim.record(emailKey) + // The send succeeded: keep both reservations (the deferred rollback becomes a + // no-op) so the cooldown windows stand. + committed = true a.audit(r, auditActor(p), "account.email.otp_sent", "") writeJSON(w, http.StatusAccepted, map[string]any{ "sent": true, diff --git a/internal/api/handlers_email_otp_test.go b/internal/api/handlers_email_otp_test.go index 8477ca0..af75113 100644 --- a/internal/api/handlers_email_otp_test.go +++ b/internal/api/handlers_email_otp_test.go @@ -5,6 +5,8 @@ import ( "errors" "net/http" "net/http/httptest" + "sync" + "sync/atomic" "testing" "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: "u1@example.net", 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":"player@example.net"}`, 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":"player@example.net"}`, 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: "u1@example.net", 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":"victim@example.net"}`, 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) + } +}