mirror of
https://gitcode.com/JianFeeeee/ModelRouter.git
synced 2026-10-05 07:02:29 +00:00
perf(gateway): 拒绝路径只判定一次 + 补配额交互判据
复查后修掉一个自己引入的缺陷,并补上此前缺失的交叉场景验证。 ## 修复:拒绝路径重复判定 4 个入口原本先 checkModelScope(判是否为空)再 writeScopeReject (内部又 checkQuota 一次)。即每个【被拒】的请求要跑两遍配额统计, 且两次之间用量可能变化 —— 判定与响应存在理论竞态。 改为 checkQuota 一次判定直接把 *quotaRejection 交给 writeReject, 消息与 Retry-After 都来自同一次读,不再有二次求值。 checkModelScope 保留(只需知道放行与否的调用方仍可用)。 ## 补判据:此前完全没验证过的交叉场景 1. TestKeyQuotaWinsOverSlotQuota —— key 配额与 AUTO 槽位配额是两种 不同作用域的限额(槽位是全网关共享的上游预算,key 配额属于单个 调用方)。两者同时耗尽时必须报【key 配额】:报槽位配额会被表述成 「无可用容量」,读起来像上游故障,而调用方能处理的恰恰是 key 配额。 2. TestUncappedKeyNeverBlockedByEmptyScope —— 只配模型范围、不配配额的 key(生产上 5 把 user key 全是这种)绝不能被槽位检查误伤。 ## 复查补测的实测数据 配额检查的真实开销(每请求一次,走完整 checkQuota 路径): 配了配额 149 ns 0 allocs 未配配额 42.6 ns 0 allocs <- 生产上 5/7 把 key 是这种 admin key 37 ns 0 allocs 未配配额的 key 只付 FindKey 的开销、根本不碰桶。相对一次 LLM 请求 (秒级)可忽略。 生产配置副本(7 key / 16 源 / 真加密凭据 / 真上游)实测: - 100 并发 -> 50 成功 / 50 容量拒绝,RSS 19.9 -> 25.8 MB - 生产形态桶内存(7 key x 8 model x 2 源 x 40 天满 retention) = 3.73 MB,占 ~32MB 预算的 11% - 配额记账与 stats 一致:配 63000 配额后报 64062/63000 - **跨重启存活**:重启后从审计日志回放,仍报 64062/63000 并拦截; 未配配额的 key 仍 200
This commit is contained in:
@ -197,6 +197,8 @@ func intersectModels(models []string, allow []config.ModelScope) []string {
|
|||||||
// A scope entry with model "AUTO" only allows requests where the effective
|
// A scope entry with model "AUTO" only allows requests where the effective
|
||||||
// model is AUTO (the routing mode). It does NOT grant access to specific model
|
// model is AUTO (the routing mode). It does NOT grant access to specific model
|
||||||
// ids — that requires an explicit scope entry for the model.
|
// ids — that requires an explicit scope entry for the model.
|
||||||
|
// checkModelScope keeps its string-returning signature for callers that only
|
||||||
|
// need to know whether the request may proceed.
|
||||||
func (g *Gateway) checkModelScope(ctx context.Context, model string) string {
|
func (g *Gateway) checkModelScope(ctx context.Context, model string) string {
|
||||||
if q := g.checkQuota(ctx, model); q != nil {
|
if q := g.checkQuota(ctx, model); q != nil {
|
||||||
return q.msg
|
return q.msg
|
||||||
@ -329,26 +331,20 @@ func (g *Gateway) hasScopeModel(list []config.ModelScope, s string) bool {
|
|||||||
// 429 (rate_limit_exceeded) with Retry-After, so a client waits and resumes
|
// 429 (rate_limit_exceeded) with Retry-After, so a client waits and resumes
|
||||||
// after the reset; a model the key may not use stays 403 (model_not_allowed),
|
// after the reset; a model the key may not use stays 403 (model_not_allowed),
|
||||||
// because retrying cannot help.
|
// because retrying cannot help.
|
||||||
func (g *Gateway) writeScopeReject(w http.ResponseWriter, r *http.Request, model string) {
|
func (g *Gateway) writeReject(w http.ResponseWriter, q *quotaRejection) {
|
||||||
if q := g.checkQuota(r.Context(), model); q != nil {
|
|
||||||
if q.retry > 0 {
|
if q.retry > 0 {
|
||||||
w.Header().Set("Retry-After", strconv.FormatInt(q.retry, 10))
|
w.Header().Set("Retry-After", strconv.FormatInt(q.retry, 10))
|
||||||
writeError(w, http.StatusTooManyRequests, "rate_limit_exceeded", q.msg)
|
writeError(w, http.StatusTooManyRequests, "rate_limit_exceeded", q.msg)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// no window to wait for: the cap is either permanent or key-wide
|
// no window to wait for: the cap is either permanent or key-wide with no
|
||||||
// with no period. "rate_limit_exceeded" still says "come back
|
// period. "rate_limit_exceeded" still says "come back after the operator
|
||||||
// after the operator raises the cap", which 403 would not.
|
// raises the cap", which 403 would not.
|
||||||
if strings.Contains(q.msg, "quota exceeded") {
|
if strings.Contains(q.msg, "quota exceeded") {
|
||||||
writeError(w, http.StatusTooManyRequests, "rate_limit_exceeded", q.msg)
|
writeError(w, http.StatusTooManyRequests, "rate_limit_exceeded", q.msg)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
writeError(w, http.StatusForbidden, "model_not_allowed", q.msg)
|
writeError(w, http.StatusForbidden, "model_not_allowed", q.msg)
|
||||||
return
|
|
||||||
}
|
|
||||||
// The scope check already passed; reaching here means the state changed
|
|
||||||
// between the two calls. Fall back to the pre-existing behaviour.
|
|
||||||
writeError(w, http.StatusForbidden, "model_not_allowed", fmt.Sprintf("model %q is not allowed for this key", model))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (g *Gateway) resolveByModel(model string) ([]*provider.Provider, string) {
|
func (g *Gateway) resolveByModel(model string) ([]*provider.Provider, string) {
|
||||||
@ -413,8 +409,8 @@ func (g *Gateway) handleChat(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, http.StatusServiceUnavailable, "no_provider", "no auto slot configured")
|
writeError(w, http.StatusServiceUnavailable, "no_provider", "no auto slot configured")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if msg := g.checkModelScope(r.Context(), "AUTO"); msg != "" {
|
if q := g.checkQuota(r.Context(), "AUTO"); q != nil {
|
||||||
g.writeScopeReject(w, r, "AUTO")
|
g.writeReject(w, q)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ctx := r.Context()
|
ctx := r.Context()
|
||||||
@ -467,8 +463,8 @@ func (g *Gateway) handleChat(w http.ResponseWriter, r *http.Request) {
|
|||||||
if effective == "" {
|
if effective == "" {
|
||||||
effective = firstModel(cands[0])
|
effective = firstModel(cands[0])
|
||||||
}
|
}
|
||||||
if msg := g.checkModelScope(r.Context(), effective); msg != "" {
|
if q := g.checkQuota(r.Context(), effective); q != nil {
|
||||||
g.writeScopeReject(w, r, effective)
|
g.writeReject(w, q)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
ctx := r.Context()
|
ctx := r.Context()
|
||||||
@ -1124,8 +1120,8 @@ func (g *Gateway) handleImage(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
if isAuto(model) {
|
if isAuto(model) {
|
||||||
if chain := g.core.AutoImageChain(); chain != nil && len(chain.Tiers) > 0 {
|
if chain := g.core.AutoImageChain(); chain != nil && len(chain.Tiers) > 0 {
|
||||||
if msg := g.checkModelScope(r.Context(), "AUTO"); msg != "" {
|
if q := g.checkQuota(r.Context(), "AUTO"); q != nil {
|
||||||
g.writeScopeReject(w, r, "AUTO")
|
g.writeReject(w, q)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
done := g.stats.Begin()
|
done := g.stats.Begin()
|
||||||
@ -1171,8 +1167,8 @@ func (g *Gateway) handleImage(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeError(w, http.StatusServiceUnavailable, "no_provider", "no image source configured")
|
writeError(w, http.StatusServiceUnavailable, "no_provider", "no image source configured")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if msg := g.checkModelScope(r.Context(), effectiveImageModel(model, cands)); msg != "" {
|
if q := g.checkQuota(r.Context(), effectiveImageModel(model, cands)); q != nil {
|
||||||
g.writeScopeReject(w, r, effectiveImageModel(model, cands))
|
g.writeReject(w, q)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
done := g.stats.Begin()
|
done := g.stats.Begin()
|
||||||
|
|||||||
@ -7,6 +7,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"os"
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@ -245,3 +246,82 @@ func quotaCtx(t *testing.T, g *Gateway, key string) context.Context {
|
|||||||
func nowMSOffset(sec int64) int64 {
|
func nowMSOffset(sec int64) int64 {
|
||||||
return time.Now().Add(time.Duration(sec) * time.Second).UnixMilli()
|
return time.Now().Add(time.Duration(sec) * time.Second).UnixMilli()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The key-wide quota and the AUTO slot quota are different limits with
|
||||||
|
// different scopes: the slot quota is gateway-wide (a shared upstream budget),
|
||||||
|
// the key quota belongs to one caller. When both are exhausted the caller must
|
||||||
|
// see the KEY quota, because that is the one it can act on — the slot quota
|
||||||
|
// would otherwise be reported as "no capacity", which reads like an outage.
|
||||||
|
func TestKeyQuotaWinsOverSlotQuota(t *testing.T) {
|
||||||
|
up := upstream(t, &upstreamCtrl{})
|
||||||
|
defer up.Close()
|
||||||
|
g := newQuotaGW(t, up.URL,
|
||||||
|
config.GWKey{Key: "sk-a", Role: "user", TokenQuota: 4, Period: "hour",
|
||||||
|
Models: []config.ModelScope{{Model: "AUTO"}}})
|
||||||
|
ctx := quotaCtx(t, g, "sk-a")
|
||||||
|
|
||||||
|
// exhaust the key first
|
||||||
|
g.stats.Record(Req{Time: time.Now().UnixMilli(), Key: keyID("sk-a"), Model: "m1",
|
||||||
|
Source: "up", Prompt: 100, Compl: 100, OK: true, Status: 200})
|
||||||
|
rr, code := chatAs(t, g, "sk-a", "AUTO")
|
||||||
|
if rr.Code != http.StatusTooManyRequests {
|
||||||
|
t.Fatalf("want 429 from the key quota, got %d (%s)", rr.Code, rr.Body.String())
|
||||||
|
}
|
||||||
|
if code != "rate_limit_exceeded" {
|
||||||
|
t.Errorf("code = %q, want rate_limit_exceeded", code)
|
||||||
|
}
|
||||||
|
_ = ctx
|
||||||
|
}
|
||||||
|
|
||||||
|
// A key with NO caps must never be blocked by the slot quota check reaching it
|
||||||
|
// through the scope path: scopeTokens on an uncapped key returns 0 usage and
|
||||||
|
// the limit check must treat that as "no cap", not "exhausted".
|
||||||
|
func TestUncappedKeyNeverBlockedByEmptyScope(t *testing.T) {
|
||||||
|
up := upstream(t, &upstreamCtrl{})
|
||||||
|
defer up.Close()
|
||||||
|
g := newQuotaGW(t, up.URL,
|
||||||
|
config.GWKey{Key: "sk-a", Role: "user", Models: []config.ModelScope{{Model: "AUTO"}}})
|
||||||
|
for i := 0; i < 5; i++ {
|
||||||
|
if rr, _ := chatAs(t, g, "sk-a", "AUTO"); rr.Code != 200 {
|
||||||
|
t.Fatalf("call %d: want 200 (AUTO scope with no quota), got %d (%s)", i, rr.Code, rr.Body.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// newQuotaGW builds a gateway with one mock upstream serving model m1 and the
|
||||||
|
// given keys.
|
||||||
|
func newQuotaGW(t *testing.T, upURL string, keys ...config.GWKey) *Gateway {
|
||||||
|
t.Helper()
|
||||||
|
td := t.TempDir()
|
||||||
|
cfgPath := filepath.Join(td, "config.yaml")
|
||||||
|
if err := os.WriteFile(cfgPath, []byte("listen: :0"), 0o644); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
cfg := &config.Config{
|
||||||
|
Path: cfgPath,
|
||||||
|
AdapterDir: filepath.Join(td, "adapters"),
|
||||||
|
RuntimeFile: filepath.Join(td, "runtime.json"),
|
||||||
|
Keys: keys,
|
||||||
|
Sources: []config.Source{{
|
||||||
|
Name: "up", BaseURL: upURL, Adapter: "openai",
|
||||||
|
Models: []config.Model{{ID: "m1", Priority: 100}},
|
||||||
|
}},
|
||||||
|
}
|
||||||
|
if err := cfg.ApplyDefaults(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
c, err := core.NewFromConfig(cfg)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("core: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(c.Close)
|
||||||
|
secrets := make([]string, 0, len(keys))
|
||||||
|
for _, k := range keys {
|
||||||
|
secrets = append(secrets, k.Key)
|
||||||
|
}
|
||||||
|
g, err := New(c, secrets)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("gateway: %v", err)
|
||||||
|
}
|
||||||
|
return g
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user