Arm the shared cooldown for rate limits tunneled through HTTP 200
Azure is the only zero-data-retention route for the gpt-5.6 family, so its capacity 429s arrive frequently and OpenRouter forwards them inside an HTTP 200 envelope. Those bypassed the rate controller entirely: a paced run kept sending a request every three seconds into a throttled endpoint, failing row by row. An in-envelope 429 now records the same escalating cooldown as a transport 429, so later acquisitions fail fast until the deadline passes.
This commit is contained in:
@@ -365,6 +365,12 @@ func (c *Client) complete(ctx context.Context, gate *ratelimit.Controller, r com
|
|||||||
Code int `json:"code"`
|
Code int `json:"code"`
|
||||||
}
|
}
|
||||||
_ = json.Unmarshal(envelope.Error, &detail)
|
_ = json.Unmarshal(envelope.Error, &detail)
|
||||||
|
if detail.Code == http.StatusTooManyRequests {
|
||||||
|
// An upstream rate limit tunneled through HTTP 200 must arm the
|
||||||
|
// same cooldown as a transport 429: later Acquire calls fail fast
|
||||||
|
// instead of pacing more requests into a throttled endpoint.
|
||||||
|
return "", gate.ReportLimit()
|
||||||
|
}
|
||||||
if detail.Code != 0 {
|
if detail.Code != 0 {
|
||||||
return "", fmt.Errorf("AI provider reported an error (code %d)", detail.Code)
|
return "", fmt.Errorf("AI provider reported an error (code %d)", detail.Code)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ import (
|
|||||||
"reflect"
|
"reflect"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"finance-duck/internal/domain"
|
"finance-duck/internal/domain"
|
||||||
"finance-duck/internal/ratelimit"
|
"finance-duck/internal/ratelimit"
|
||||||
@@ -358,6 +359,30 @@ func TestMalformedEnvelopesRejected(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// An upstream rate limit tunneled inside an HTTP 200 envelope must arm the
|
||||||
|
// shared cooldown like a transport 429: the next classification fails fast
|
||||||
|
// instead of pacing another request into a throttled endpoint.
|
||||||
|
func TestEnvelope429ArmsSharedCooldown(t *testing.T) {
|
||||||
|
f, d := fixture()
|
||||||
|
calls := 0
|
||||||
|
c := mockClient(t, func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
calls++
|
||||||
|
_, _ = io.WriteString(w, `{"error":{"code":429,"message":"private"},"choices":[]}`)
|
||||||
|
})
|
||||||
|
c.rate.Store(&ratelimit.Controller{InitialBackoff: time.Minute})
|
||||||
|
_, err := c.Classify(context.Background(), f, d, true)
|
||||||
|
var limit *ratelimit.RateLimitError
|
||||||
|
if err == nil || !errors.As(err, &limit) || strings.Contains(err.Error(), "private") {
|
||||||
|
t.Fatalf("envelope 429 not reported as a rate limit: %v", err)
|
||||||
|
}
|
||||||
|
if _, err = c.Classify(context.Background(), f, d, true); err == nil || !errors.As(err, &limit) {
|
||||||
|
t.Fatalf("cooldown not armed: %v", err)
|
||||||
|
}
|
||||||
|
if calls != 1 {
|
||||||
|
t.Fatalf("throttled endpoint was contacted again: %d calls", calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
type failingTransport struct{}
|
type failingTransport struct{}
|
||||||
|
|
||||||
func (failingTransport) RoundTrip(*http.Request) (*http.Response, error) {
|
func (failingTransport) RoundTrip(*http.Request) (*http.Response, error) {
|
||||||
|
|||||||
@@ -109,6 +109,43 @@ func (g *Controller) Release() {
|
|||||||
<-g.active
|
<-g.active
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// recordLimit escalates the consecutive-failure backoff, retains the cooldown
|
||||||
|
// and learns spacing. Callers hold the Acquire gate, like Do's 429 branch.
|
||||||
|
func (g *Controller) recordLimit(header string) *RateLimitError {
|
||||||
|
if g.backoff <= 0 {
|
||||||
|
g.backoff = g.InitialBackoff
|
||||||
|
if g.backoff <= 0 {
|
||||||
|
g.backoff = time.Second
|
||||||
|
}
|
||||||
|
} else if g.backoff >= maxBackoff/2 {
|
||||||
|
g.backoff = max(g.backoff, maxBackoff)
|
||||||
|
} else {
|
||||||
|
g.backoff *= 2
|
||||||
|
}
|
||||||
|
fallback := max(g.backoff, g.MinimumInterval, g.learnedInterval)
|
||||||
|
limit := retryLimit(header, time.Now(), fallback)
|
||||||
|
g.mu.Lock()
|
||||||
|
g.limit = limit
|
||||||
|
g.mu.Unlock()
|
||||||
|
// Keep the most conservative learned cadence for this controller's
|
||||||
|
// lifetime, capped at 30 seconds. The actual provider deadline is never
|
||||||
|
// capped; persistent failures separately escalate up to 15 minutes.
|
||||||
|
learned := maxLearnedInterval
|
||||||
|
if !limit.unbounded {
|
||||||
|
learned = min(learned, time.Until(limit.next))
|
||||||
|
}
|
||||||
|
g.learnedInterval = max(g.learnedInterval, learned)
|
||||||
|
return limit
|
||||||
|
}
|
||||||
|
|
||||||
|
// ReportLimit records a rate limit the provider communicated outside the HTTP
|
||||||
|
// status — typically inside an HTTP 200 error envelope — so later Acquire
|
||||||
|
// calls fail fast during the cooldown exactly as after a transport HTTP 429.
|
||||||
|
// It must be called while holding an Acquire, like Do.
|
||||||
|
func (g *Controller) ReportLimit() *RateLimitError {
|
||||||
|
return g.recordLimit("")
|
||||||
|
}
|
||||||
|
|
||||||
// retryLimit never converts a positive overflowing delay into a short wait.
|
// retryLimit never converts a positive overflowing delay into a short wait.
|
||||||
// Delays beyond time.Duration's range disable retries rather than truncate the
|
// Delays beyond time.Duration's range disable retries rather than truncate the
|
||||||
// provider's instruction. HTTP dates retain their absolute timestamp unchanged.
|
// provider's instruction. HTTP dates retain their absolute timestamp unchanged.
|
||||||
@@ -202,29 +239,7 @@ func (g *Controller) Do(ctx context.Context, attempt func(context.Context) (*htt
|
|||||||
}
|
}
|
||||||
return resp, nil
|
return resp, nil
|
||||||
}
|
}
|
||||||
if g.backoff <= 0 {
|
limit := g.recordLimit(resp.Header.Get("Retry-After"))
|
||||||
g.backoff = g.InitialBackoff
|
|
||||||
if g.backoff <= 0 {
|
|
||||||
g.backoff = time.Second
|
|
||||||
}
|
|
||||||
} else if g.backoff >= maxBackoff/2 {
|
|
||||||
g.backoff = max(g.backoff, maxBackoff)
|
|
||||||
} else {
|
|
||||||
g.backoff *= 2
|
|
||||||
}
|
|
||||||
fallback := max(g.backoff, g.MinimumInterval, g.learnedInterval)
|
|
||||||
limit := retryLimit(resp.Header.Get("Retry-After"), time.Now(), fallback)
|
|
||||||
g.mu.Lock()
|
|
||||||
g.limit = limit
|
|
||||||
g.mu.Unlock()
|
|
||||||
// Keep the most conservative learned cadence for this controller's
|
|
||||||
// lifetime, capped at 30 seconds. The actual provider deadline is never
|
|
||||||
// capped; persistent failures separately escalate up to 15 minutes.
|
|
||||||
learned := maxLearnedInterval
|
|
||||||
if !limit.unbounded {
|
|
||||||
learned = min(learned, time.Until(limit.next))
|
|
||||||
}
|
|
||||||
g.learnedInterval = max(g.learnedInterval, learned)
|
|
||||||
// Never read or expose provider errors, and release each response before
|
// Never read or expose provider errors, and release each response before
|
||||||
// any sleep or retry. Other responses are processed by the caller.
|
// any sleep or retry. Other responses are processed by the caller.
|
||||||
resp.Body.Close()
|
resp.Body.Close()
|
||||||
|
|||||||
Reference in New Issue
Block a user