diff --git a/backend/internal/config/config.go b/backend/internal/config/config.go index 2422b2cab0..2fefe820c9 100644 --- a/backend/internal/config/config.go +++ b/backend/internal/config/config.go @@ -957,6 +957,9 @@ type GatewayOpenAIWSSchedulerScoreWeights struct { Queue float64 `mapstructure:"queue"` ErrorRate float64 `mapstructure:"error_rate"` TTFT float64 `mapstructure:"ttft"` + // Reset 倾向「会话窗口最早重置」的账号(use-it-or-lose-it)。 + // >0 时,剩余重置时间越短的账号得分越高,从而被优先用尽。默认 0(关闭,不改变原有行为)。 + Reset float64 `mapstructure:"reset"` } // GatewayOpenAISchedulerConfig OpenAI 高级调度器配置。 @@ -1055,6 +1058,11 @@ type GatewaySchedulingConfig struct { // 兜底层账户选择策略: "last_used"(按最后使用时间排序,默认) 或 "random"(随机) FallbackSelectionMode string `mapstructure:"fallback_selection_mode"` + // PreferSoonestReset 开启后,负载感知选择会优先选用「会话窗口最早重置」的账号 + // (use-it-or-lose-it:先用尽即将重置的账号,保留重置时间还很久的账号)。 + // 默认 false,保持原有「优先级 → 负载率 → LRU」行为不变。 + PreferSoonestReset bool `mapstructure:"prefer_soonest_reset"` + // 负载计算 LoadBatchEnabled bool `mapstructure:"load_batch_enabled"` LoadBatchCacheTTLMS int `mapstructure:"load_batch_cache_ttl_ms"` @@ -1870,6 +1878,7 @@ func setDefaults() { viper.SetDefault("gateway.openai_ws.scheduler_score_weights.queue", 0.7) viper.SetDefault("gateway.openai_ws.scheduler_score_weights.error_rate", 0.8) viper.SetDefault("gateway.openai_ws.scheduler_score_weights.ttft", 0.5) + viper.SetDefault("gateway.openai_ws.scheduler_score_weights.reset", 0.0) // OpenAI HTTP upstream protocol strategy viper.SetDefault("gateway.openai_http2.enabled", true) viper.SetDefault("gateway.openai_http2.allow_proxy_fallback_to_http1", true) @@ -1906,6 +1915,7 @@ func setDefaults() { viper.SetDefault("gateway.scheduling.fallback_wait_timeout", 30*time.Second) viper.SetDefault("gateway.scheduling.fallback_max_waiting", 100) viper.SetDefault("gateway.scheduling.fallback_selection_mode", "last_used") + viper.SetDefault("gateway.scheduling.prefer_soonest_reset", false) viper.SetDefault("gateway.scheduling.load_batch_enabled", true) viper.SetDefault("gateway.scheduling.load_batch_cache_ttl_ms", 200) viper.SetDefault("gateway.scheduling.snapshot_mget_chunk_size", 128) diff --git a/backend/internal/service/gateway_service.go b/backend/internal/service/gateway_service.go index 5a9a33880e..209c9615b7 100644 --- a/backend/internal/service/gateway_service.go +++ b/backend/internal/service/gateway_service.go @@ -2218,13 +2218,17 @@ func (s *GatewayService) SelectAccountWithLoadAwareness(ctx context.Context, gro } } - // 分层过滤选择:优先级 → 负载率 → LRU + // 分层过滤选择:优先级 →(可选)最早重置 → 负载率 → LRU for len(available) > 0 { // 1. 取优先级最小的集合 candidates := filterByMinPriority(available) - // 2. 取负载率最低的集合 + // 2. (可选)use-it-or-lose-it:优先选用会话窗口最早重置的账号 + if cfg.PreferSoonestReset { + candidates = filterBySoonestReset(candidates) + } + // 3. 取负载率最低的集合 candidates = filterByMinLoadRate(candidates) - // 3. LRU 选择最久未用的账号 + // 4. LRU 选择最久未用的账号 selected := selectByLRU(candidates, preferOAuth) if selected == nil { break @@ -2983,6 +2987,39 @@ func filterByMinLoadRate(accounts []accountWithLoad) []accountWithLoad { return result } +// filterBySoonestReset 过滤出「会话窗口最早重置」的账号集合(use-it-or-lose-it)。 +// 仅保留拥有未来重置时间(SessionWindowEnd 在当前时间之后)且最早的账号; +// 窗口为空或已过期的账号视为无活跃窗口、优先级最低。 +// 当所有账号都没有活跃窗口时,返回原集合(不改变后续 LRU 选择)。 +func filterBySoonestReset(accounts []accountWithLoad) []accountWithLoad { + if len(accounts) <= 1 { + return accounts + } + now := time.Now() + var minEnd *time.Time + for _, acc := range accounts { + end := acc.account.SessionWindowEnd + if end == nil || !now.Before(*end) { + continue + } + if minEnd == nil || end.Before(*minEnd) { + minEnd = end + } + } + if minEnd == nil { + // 没有任何账号拥有活跃窗口,保持原集合 + return accounts + } + result := make([]accountWithLoad, 0, len(accounts)) + for _, acc := range accounts { + end := acc.account.SessionWindowEnd + if end != nil && now.Before(*end) && end.Equal(*minEnd) { + result = append(result, acc) + } + } + return result +} + // selectByLRU 从集合中选择最久未用的账号 // 如果有多个账号具有相同的最小 LastUsedAt,则随机选择一个 func selectByLRU(accounts []accountWithLoad, preferOAuth bool) *accountWithLoad { diff --git a/backend/internal/service/gateway_soonest_reset_test.go b/backend/internal/service/gateway_soonest_reset_test.go new file mode 100644 index 0000000000..953d4d1d0e --- /dev/null +++ b/backend/internal/service/gateway_soonest_reset_test.go @@ -0,0 +1,81 @@ +//go:build unit + +package service + +import ( + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func accWithWindowEnd(id int64, end *time.Time) accountWithLoad { + return accountWithLoad{ + account: &Account{ + ID: id, + Schedulable: true, + Status: StatusActive, + SessionWindowEnd: end, + }, + loadInfo: &AccountLoadInfo{AccountID: id}, + } +} + +func TestFilterBySoonestReset_PicksSoonestFutureWindow(t *testing.T) { + now := time.Now() + soon := now.Add(1 * time.Hour) + later := now.Add(24 * time.Hour) + accounts := []accountWithLoad{ + accWithWindowEnd(1, testTimePtr(later)), + accWithWindowEnd(2, testTimePtr(soon)), + accWithWindowEnd(3, testTimePtr(later)), + } + got := filterBySoonestReset(accounts) + require.Len(t, got, 1) + require.Equal(t, int64(2), got[0].account.ID, "重置时间最早的账号被选中") +} + +func TestFilterBySoonestReset_IgnoresNilAndExpiredWindows(t *testing.T) { + now := time.Now() + expired := now.Add(-1 * time.Hour) + active := now.Add(2 * time.Hour) + accounts := []accountWithLoad{ + accWithWindowEnd(1, nil), // 无活跃窗口 + accWithWindowEnd(2, testTimePtr(expired)), // 已过期,视为无活跃窗口 + accWithWindowEnd(3, testTimePtr(active)), // 唯一活跃窗口 + } + got := filterBySoonestReset(accounts) + require.Len(t, got, 1) + require.Equal(t, int64(3), got[0].account.ID, "仅保留拥有未来重置时间的账号") +} + +func TestFilterBySoonestReset_NoActiveWindowReturnsAll(t *testing.T) { + now := time.Now() + expired := now.Add(-30 * time.Minute) + accounts := []accountWithLoad{ + accWithWindowEnd(1, nil), + accWithWindowEnd(2, testTimePtr(expired)), + } + got := filterBySoonestReset(accounts) + require.Len(t, got, 2, "没有任何账号拥有活跃窗口时,返回原集合不做过滤") +} + +func TestFilterBySoonestReset_TiedSoonestKeepsAll(t *testing.T) { + now := time.Now() + end := now.Add(90 * time.Minute) + accounts := []accountWithLoad{ + accWithWindowEnd(1, testTimePtr(end)), + accWithWindowEnd(2, testTimePtr(end)), + accWithWindowEnd(3, testTimePtr(now.Add(5*time.Hour))), + } + got := filterBySoonestReset(accounts) + require.Len(t, got, 2, "并列最早重置的账号都保留,交由后续 LRU 决定") + ids := map[int64]bool{got[0].account.ID: true, got[1].account.ID: true} + require.True(t, ids[1] && ids[2]) +} + +func TestFilterBySoonestReset_SingleOrEmptyUnchanged(t *testing.T) { + require.Empty(t, filterBySoonestReset(nil)) + single := []accountWithLoad{accWithWindowEnd(1, nil)} + require.Len(t, filterBySoonestReset(single), 1) +} diff --git a/backend/internal/service/openai_account_scheduler.go b/backend/internal/service/openai_account_scheduler.go index ece3716b57..59da55d3fe 100644 --- a/backend/internal/service/openai_account_scheduler.go +++ b/backend/internal/service/openai_account_scheduler.go @@ -745,6 +745,35 @@ func (s *defaultOpenAIAccountScheduler) buildOpenAIAccountLoadPlan( plan.loadSkew = calcLoadSkewByMoments(loadRateSum, loadRateSumSquares, len(candidates)) weights := s.service.openAIWSSchedulerWeights() + + // Reset 因子(use-it-or-lose-it):在拥有「未来会话窗口结束时间」的账号中, + // 剩余时间越短 → 因子越接近 1(越早重置越优先用尽)。无活跃窗口的账号因子为 0。 + // 仅在 weights.Reset > 0 时计算,默认关闭不影响原有行为。 + minResetRemaining, maxResetRemaining := 0.0, 0.0 + hasResetSample := false + if weights.Reset > 0 { + now := time.Now() + for _, candidate := range candidates { + end := candidate.account.SessionWindowEnd + if end == nil || !now.Before(*end) { + continue + } + remaining := end.Sub(now).Seconds() + if !hasResetSample { + minResetRemaining, maxResetRemaining = remaining, remaining + hasResetSample = true + continue + } + if remaining < minResetRemaining { + minResetRemaining = remaining + } + if remaining > maxResetRemaining { + maxResetRemaining = remaining + } + } + } + + now := time.Now() for i := range candidates { item := &candidates[i] priorityFactor := 1.0 @@ -758,12 +787,24 @@ func (s *defaultOpenAIAccountScheduler) buildOpenAIAccountLoadPlan( if item.hasTTFT && hasTTFTSample && maxTTFT > minTTFT { ttftFactor = 1 - clamp01((item.ttft-minTTFT)/(maxTTFT-minTTFT)) } + resetFactor := 0.0 + if weights.Reset > 0 && hasResetSample { + if end := item.account.SessionWindowEnd; end != nil && now.Before(*end) { + if maxResetRemaining > minResetRemaining { + resetFactor = 1 - clamp01((end.Sub(now).Seconds()-minResetRemaining)/(maxResetRemaining-minResetRemaining)) + } else { + // 所有有窗口的账号剩余时间相同:一律给满分,让其优于无窗口账号。 + resetFactor = 1 + } + } + } item.score = weights.Priority*priorityFactor + weights.Load*loadFactor + weights.Queue*queueFactor + weights.ErrorRate*errorFactor + - weights.TTFT*ttftFactor + weights.TTFT*ttftFactor + + weights.Reset*resetFactor } plan.candidates = candidates @@ -1415,6 +1456,7 @@ func (s *OpenAIGatewayService) openAIWSSchedulerWeights() GatewayOpenAIWSSchedul Queue: s.cfg.Gateway.OpenAIWS.SchedulerScoreWeights.Queue, ErrorRate: s.cfg.Gateway.OpenAIWS.SchedulerScoreWeights.ErrorRate, TTFT: s.cfg.Gateway.OpenAIWS.SchedulerScoreWeights.TTFT, + Reset: s.cfg.Gateway.OpenAIWS.SchedulerScoreWeights.Reset, } } return GatewayOpenAIWSSchedulerScoreWeightsView{ @@ -1423,6 +1465,7 @@ func (s *OpenAIGatewayService) openAIWSSchedulerWeights() GatewayOpenAIWSSchedul Queue: 0.7, ErrorRate: 0.8, TTFT: 0.5, + Reset: 0.0, } } @@ -1432,6 +1475,8 @@ type GatewayOpenAIWSSchedulerScoreWeightsView struct { Queue float64 ErrorRate float64 TTFT float64 + // Reset 倾向「会话窗口最早重置」的账号;0 表示关闭(默认)。 + Reset float64 } func clamp01(value float64) float64 { diff --git a/backend/internal/service/openai_account_scheduler_reset_test.go b/backend/internal/service/openai_account_scheduler_reset_test.go new file mode 100644 index 0000000000..1ba6cb570c --- /dev/null +++ b/backend/internal/service/openai_account_scheduler_reset_test.go @@ -0,0 +1,77 @@ +package service + +import ( + "testing" + "time" + + "github.com/Wei-Shaw/sub2api/internal/config" + "github.com/stretchr/testify/require" +) + +func openAIResetTestScheduler(reset float64) *defaultOpenAIAccountScheduler { + cfg := &config.Config{} + cfg.Gateway.OpenAIWS.SchedulerScoreWeights = config.GatewayOpenAIWSSchedulerScoreWeights{ + Priority: 1.0, + Load: 1.0, + Queue: 0.7, + ErrorRate: 0.8, + TTFT: 0.5, + Reset: reset, + } + return &defaultOpenAIAccountScheduler{service: &OpenAIGatewayService{cfg: cfg}} +} + +func openAIPlanScores(plan openAIAccountLoadPlan) map[int64]float64 { + scores := make(map[int64]float64, len(plan.candidates)) + for _, c := range plan.candidates { + scores[c.account.ID] = c.score + } + return scores +} + +// Reset 权重 > 0 时,会话窗口最早重置的账号应获得更高分。 +func TestBuildOpenAIAccountLoadPlan_ResetWeightPrefersSoonestReset(t *testing.T) { + now := time.Now() + soon := now.Add(1 * time.Hour) + later := now.Add(20 * time.Hour) + filtered := []*Account{ + {ID: 1, Priority: 0, SessionWindowEnd: &later}, + {ID: 2, Priority: 0, SessionWindowEnd: &soon}, + } + sched := openAIResetTestScheduler(5.0) + + plan := sched.buildOpenAIAccountLoadPlan(OpenAIAccountScheduleRequest{}, filtered, map[int64]*AccountLoadInfo{}) + scores := openAIPlanScores(plan) + require.Greater(t, scores[2], scores[1], "重置时间最早的账号(ID=2)得分更高") +} + +// Reset 权重为 0(默认)时,窗口重置时间不应影响打分,保持原有行为。 +func TestBuildOpenAIAccountLoadPlan_ResetWeightZeroNoEffect(t *testing.T) { + now := time.Now() + soon := now.Add(1 * time.Hour) + later := now.Add(20 * time.Hour) + filtered := []*Account{ + {ID: 1, Priority: 0, SessionWindowEnd: &later}, + {ID: 2, Priority: 0, SessionWindowEnd: &soon}, + } + sched := openAIResetTestScheduler(0.0) + + plan := sched.buildOpenAIAccountLoadPlan(OpenAIAccountScheduleRequest{}, filtered, map[int64]*AccountLoadInfo{}) + scores := openAIPlanScores(plan) + require.Equal(t, scores[1], scores[2], "Reset 权重为 0 时两账号得分相同") +} + +// 无活跃窗口的账号 reset 因子为 0,应低于拥有未来窗口的账号。 +func TestBuildOpenAIAccountLoadPlan_ResetWeightIgnoresNilWindow(t *testing.T) { + now := time.Now() + soon := now.Add(2 * time.Hour) + filtered := []*Account{ + {ID: 1, Priority: 0, SessionWindowEnd: nil}, + {ID: 2, Priority: 0, SessionWindowEnd: &soon}, + } + sched := openAIResetTestScheduler(5.0) + + plan := sched.buildOpenAIAccountLoadPlan(OpenAIAccountScheduleRequest{}, filtered, map[int64]*AccountLoadInfo{}) + scores := openAIPlanScores(plan) + require.Greater(t, scores[2], scores[1], "拥有活跃窗口的账号得分高于无窗口账号") +} diff --git a/deploy/config.example.yaml b/deploy/config.example.yaml index 35c7696423..2cb65d8fec 100644 --- a/deploy/config.example.yaml +++ b/deploy/config.example.yaml @@ -320,6 +320,9 @@ gateway: queue: 0.7 error_rate: 0.8 ttft: 0.5 + # use-it-or-lose-it:倾向「会话窗口最早重置」的账号,剩余重置时间越短得分越高。 + # 0 表示关闭(默认,不改变原有行为);调大可让即将重置的账号被优先用尽。 + reset: 0.0 # OpenAI 高级调度器补充配置 openai_scheduler: # 是否允许 session_hash sticky 在账号健康度恶化时临时逃逸;false 可一键回退旧行为 @@ -421,6 +424,10 @@ gateway: # Fallback max waiting queue size # 兜底最大排队长度 fallback_max_waiting: 100 + # Prefer the account whose session window resets soonest (use-it-or-lose-it). + # 负载感知选择时优先用尽「会话窗口最早重置」的账号;false 保持 + # 原有「优先级 → 负载率 → LRU」行为(默认)。 + prefer_soonest_reset: false # Enable batch load calculation for scheduling # 启用调度批量负载计算 load_batch_enabled: true