gcli2api / internal /pool /pool.go
a3216's picture
chore: 同步到上游 1.12.0-panel + 凭证同步/独立启动器/保活
6d60378 verified
Raw History Blame Contribute Delete
15.3 kB
// Pool 账号池核心:结构定义、构造(New/Set* 注入)、在途租约(Acquire/Release)
// 与账号增删(Add/SyncToDir/upsertLocked)。选号/冷却/状态/持久化见同包其他文件。
package pool
import (
"strings"
"sync"
"sync/atomic"
"time"
"github.com/linguo2625469/workbuddy2api-panel/internal/auth"
)
type Pool struct {
mu sync.RWMutex
byUID map[string]*entry
stateFp string
dirty atomic.Bool // 内存有变更待落盘
// store 池状态快照镜像(redisstore.Store);nil = 无需镜像(未配置 Redis / Noop 之外也可能 nil)。
// SaveState/LoadState 经它接线,与本地 state.json 并存作启动恢复备份。
store StoreSnapshotter
// 熔断器调优(SetBreaker 注入;默认值见 defaultBreaker*)。
breakerThreshold int
breakerCooldown time.Duration
breakerCooldownMax time.Duration
// softRateMax 软冷却指数退避的封顶(SetSoftRateMax 注入;默认 defaultSoftRateMax)。
softRateMax time.Duration
// costExploreInterval costTier 条件探索窗口(issue #136 方案 a′,SetCostExploreInterval
// 注入;默认 defaultCostExploreInterval 30m)。tier 0 垄断 + tier 1 存在且距上次
// 探索 ≥ 窗口时,本次 pick 生效层切 tier 1-only(探索=搭车改道,零新增上游请求)。
// 0 = 关停(完全回到现状行为)。
costExploreInterval time.Duration
// exploreLast 各 (realm, 模型) 的上次探索时刻,键 = realm + "\x1f" + model。
// 运行态(不持久化,同 lastUsed/usedSeq 口径):重启归零 → 每个仍冻结的
// (域, 模型) 多至 1 次即时重探;已毕业号经 ModelCosts 恢复 tier,学费不重付。
// 只在探索事件时写入(tier 1 枯竭期间停走,陈旧无害);不做对称清理。
exploreLast map[string]time.Time
// costExploreEvents 累计探索事件数(/status 透出;pick 写锁内 ++,无需 atomic)。
costExploreEvents int64
// degradeThreshold / degradeCooldown / degradeCooldownMax 连败降权参数
// (SetDegrade 注入;默认值见 defaultDegrade*,issue #114)。
degradeThreshold int
degradeCooldown time.Duration
degradeCooldownMax time.Duration
// creditFloor 积分保底(SetCreditFloor 注入;0 = 关闭,缺省即现状)。
// 账号 credits < floor 时不再参与选号——防止收费请求把余额打穿、连免费模型都
// 402 冷却到次日签到。签到回血(SetCreditsDetailed)越过 floor 即自动恢复。
// 全池触底且无免费模型可接时选号返回 nil(硬语义:宁 503 不打穿)。
// 收费与否的判据见 floorBlockedForModel:本地实测台账优先,缺失时用上游目录
// 倍率(modelRateOf 注入)兜底,避免「无观测的高价新模型」绕过保底。
creditFloor int64
// modelRateOf 按 (realm, 模型) 查上游目录积分倍率("0.79" / "" = 未知)。
// 由 main 用 upstream.Client.ModelRate 注入——pool 不依赖 upstream 包(避免
// 循环依赖与分层破坏),nil 时倍率兜底不生效(退化为仅本地台账判定)。
// 仅在持 p.mu 时由 floorBlockedForModel 调用;回调不得反向调用 Pool 方法。
modelRateOf func(realm, model string) string
// 加权路由的闲置补偿调优(SetWeights 注入;默认值见 defaultIdle*)。
idleWeightPerHour float64
idleWeightMax float64
// preferExpiring 最早到期优先路由开关(默认 true)。开启且快过期窗口内存在有效
// 批次时,选号在成本层内先按最早到期排序;关闭后只使用普通加权路由。
preferExpiring bool
// maxInFlight 单账号最大在途请求数;0 = 不限(租约关闭)。
maxInFlight int
// maxInFlightGlobal global 域单账号在途上限分档(WAF 403 修复 P1-1:global 域
// WAF 风控更紧,压低并发);0 = 未设置,回落 maxInFlight(不分档,零回归)。
maxInFlightGlobal int
// randInt64N 仅供测试注入确定性随机源;nil 时用 math/rand/v2 全局源。
// 生产代码不应设置此字段。
randInt64N func(n int64) int64
// persistFails 本地 state.json 连续落盘失败计数(仅 saveLocked 在持锁下读写,无需 atomic)。
// 用于落盘失败的日志节流:首败/每 N 次提醒/恢复各打一条,避免磁盘满时刷屏。
persistFails int
// pickSeq 单调递增的选号序号:每次 pick 选中账号时自增并记到 entry.usedSeq,
// 为 LRU 兜底/防惊群提供与 time.Now() 精度无关的严格全序(Windows ~0.5ms 精度下
// lastUsed 墙钟会全等)。仅 pick 写锁路径读写,无需 atomic。
pickSeq uint64
// stopCh 关闭信号:Close 关闭它使 startFlusher 的后台 goroutine 退出。
// nil = 未启动 flusher(stateFp 为空时 New 不起 flusher)。
stopCh chan struct{}
// closeOnce 保证 Close 幂等(多次调用不重复 close channel)。
closeOnce sync.Once
}
// defaultBreaker* 熔断器默认参数(FreeBuff2API 参考口径)。
func New(stateFp string) *Pool {
p := &Pool{
byUID: map[string]*entry{},
stateFp: stateFp,
breakerThreshold: defaultBreakerThreshold,
breakerCooldown: defaultBreakerCooldown,
breakerCooldownMax: defaultBreakerCooldownMax,
idleWeightPerHour: defaultIdleWeightPerHour,
idleWeightMax: defaultIdleWeightMax,
preferExpiring: true,
degradeThreshold: defaultDegradeThreshold,
degradeCooldown: defaultDegradeCooldown,
degradeCooldownMax: defaultDegradeCooldownMax,
// 探索缺省 30m:tier 0 垄断下的 tier 1 探索窗口(issue #136)。用户经
// config 显式 "0" 关停(SetCostExploreInterval(0))。
costExploreInterval: defaultCostExploreInterval,
exploreLast: map[string]time.Time{},
}
if stateFp != "" {
p.load()
p.startFlusher()
}
return p
}
// Close 停止后台落盘 goroutine 并做最后一次落盘(幂等)。
// 进程退出前调用,消除 startFlusher 的 goroutine 泄漏;不调用也不影响正确性
// (进程退出即回收),仅是生命周期卫生。
func (p *Pool) Close() {
if p.stopCh == nil {
return
}
p.closeOnce.Do(func() {
close(p.stopCh)
})
p.Flush()
}
// SetBreaker 注入熔断器参数(main 从 config 解析后调用)。非正值保留原值(用默认)。
func (p *Pool) SetBreaker(threshold int, cooldown, cooldownMax time.Duration) {
p.mu.Lock()
defer p.mu.Unlock()
if threshold > 0 {
p.breakerThreshold = threshold
}
if cooldown > 0 {
p.breakerCooldown = cooldown
}
if cooldownMax > 0 {
p.breakerCooldownMax = cooldownMax
}
}
// SetSoftRateMax 注入软冷却指数退避的封顶时长(main 从 config 解析后调用)。
// 非正值保留原值(用默认 2h),风格同 SetBreaker。
func (p *Pool) SetSoftRateMax(d time.Duration) {
p.mu.Lock()
defer p.mu.Unlock()
if d > 0 {
p.softRateMax = d
}
}
// SetCostExploreInterval 注入 costTier 条件探索窗口(main 从 config 解析后调用,
// issue #136)。0 = 关停(完全回到现状行为);正值覆盖默认 30m。
// 注意:与 SetSoftRateMax「非正值保留默认」不同,0 在这里是**合法值**(关停开关,
// 与 config 的 "0" 关停语义对齐)——不设 0 语义就无法关停探索。
func (p *Pool) SetCostExploreInterval(d time.Duration) {
p.mu.Lock()
defer p.mu.Unlock()
if d < 0 {
return // 负值非法,保留现值
}
p.costExploreInterval = d
}
// CostExploreStatus 透出探索台账(/status 用):累计探索事件数 + 各 (域, 模型)
// 的最近探索时刻(键内 \x1f 分隔符输出为 "|",与 model_costs 行对照即可读出
// 「探索→毕业」全链路)。RLock 只读遍历;map 大小受「服务过的 (域, 模型)」集合
// 约束(与 modelCost 同界,天然有界)。
func (p *Pool) CostExploreStatus() (events int64, last map[string]time.Time) {
p.mu.RLock()
defer p.mu.RUnlock()
last = make(map[string]time.Time, len(p.exploreLast))
for k, ts := range p.exploreLast {
// 键 realm+"\x1f"+model → 输出 "|"(JSON 安全可读;\x1f 不可打印)。
last[strings.ReplaceAll(k, "\x1f", "|")] = ts
}
return p.costExploreEvents, last
}
// SetWeights 注入加权路由的闲置补偿参数。非正值保留原值(用默认)。
func (p *Pool) SetWeights(idlePerHour, idleMax float64) {
p.mu.Lock()
defer p.mu.Unlock()
if idlePerHour > 0 {
p.idleWeightPerHour = idlePerHour
}
if idleMax > 0 {
p.idleWeightMax = idleMax
}
}
// SetPreferExpiring 注入最早到期优先路由开关(main 从 config 解析后调用)。
func (p *Pool) SetPreferExpiring(enabled bool) {
p.mu.Lock()
defer p.mu.Unlock()
p.preferExpiring = enabled
}
// SetDegrade 注入连败降权参数(main 从 config 解析后调用,issue #114)。
// 非正值保留原值(用默认,见 defaultDegrade*),风格同 SetBreaker/SetSoftRateMax。
func (p *Pool) SetDegrade(threshold int, cooldown, cooldownMax time.Duration) {
p.mu.Lock()
defer p.mu.Unlock()
if threshold > 0 {
p.degradeThreshold = threshold
}
if cooldown > 0 {
p.degradeCooldown = cooldown
}
if cooldownMax > 0 {
p.degradeCooldownMax = cooldownMax
}
}
// CreditFloor 透出生效的积分保底值(/status 用)。0 = 关闭。
func (p *Pool) CreditFloor() int64 {
p.mu.RLock()
defer p.mu.RUnlock()
return p.creditFloor
}
// SetCreditFloor 注入积分保底线(main 从 config 解析后调用)。
// 0 = 关闭(缺省即现状,零回归);负值非法保留原值(0)。
// 语义见 Pool.creditFloor 字段注释。
func (p *Pool) SetCreditFloor(n int64) {
p.mu.Lock()
defer p.mu.Unlock()
if n >= 0 {
p.creditFloor = n
}
}
// SetModelRateOf 注入上游目录积分倍率查表(main 用 upstream.Client.ModelRate 装配)。
// 供积分保底兜底判定「未实测过的模型是否收费」——本地台账无观测时,不能因为
// 「没学过」就放行,否则高价新模型会把触底号一次性打穿(kimi-k3-1 实案:
// 全池无观测 → 保底全部放行 → 两笔扣 111 分打穿到 0 并硬冷却到次日 04:00)。
// fn 可为 nil(清注入);回调只在持 p.mu 时被调用,不得反向调用 Pool 方法。
func (p *Pool) SetModelRateOf(fn func(realm, model string) string) {
p.mu.Lock()
defer p.mu.Unlock()
p.modelRateOf = fn
}
// SetMaxInFlight 注入单账号最大在途请求数;0 = 不限。负值保留原值。
func (p *Pool) SetMaxInFlight(n int) {
p.mu.Lock()
defer p.mu.Unlock()
if n >= 0 {
p.maxInFlight = n
}
}
// SetMaxInFlightGlobal 注入 global 域单账号在途上限(WAF 403 修复 P1-1 分档);
// 0 = 未设置,global 账号回落 maxInFlight(不分档)。负值保留原值。
func (p *Pool) SetMaxInFlightGlobal(n int) {
p.mu.Lock()
defer p.mu.Unlock()
if n >= 0 {
p.maxInFlightGlobal = n
}
}
// inFlightLimit 报告账号的生效在途上限(global 分档优先,回落 maxInFlight);
// 0 = 不限。调用方需已持 p.mu(或快照过 limit,见 Acquire)。
func (p *Pool) inFlightLimit(e *entry) int {
if p.maxInFlightGlobal > 0 && e.a.Realm() == "global" {
return p.maxInFlightGlobal
}
return p.maxInFlight
}
// SetStore 注入池状态快照镜像(redisstore.Store)。nil 表示不镜像(纯本地恢复)。
// 必须在 SyncToDir 之前调用,使"择新恢复"发生在账号对齐之前。
func (p *Pool) SetStore(s StoreSnapshotter) {
p.mu.Lock()
defer p.mu.Unlock()
p.store = s
}
// RestoreFromSnapshot 择新恢复:比较本地 state.json 与 Redis 快照,采用较新者。
// 无快照、快照无 savedAt、或本地不存在/不可读时,都会被判定为"本地优先/跳过快照",
// 同时打一条恢复来源日志。必须在 SyncToDir 之前调用(SyncToDir 只增删不入值)。
// Acquire 为 uid 占一个在途名额(会话粘性命中后调用);池上限内返回 true。
// 名额用 entry.inFlight 原子自增,满额返回 false。上限按账号 realm 分档
// (global 档 maxInFlightGlobal,P1-1;未设置回落 maxInFlight)。
func (p *Pool) Acquire(uid string) bool {
p.mu.RLock()
e, ok := p.byUID[uid]
if !ok {
p.mu.RUnlock()
return false
}
limit := p.inFlightLimit(e)
p.mu.RUnlock()
if limit <= 0 {
// 不限:计数仍累加(供状态观测),但永不拒绝。
e.inFlight.Add(1)
return true
}
for {
cur := e.inFlight.Load()
if cur >= int64(limit) {
return false
}
if e.inFlight.CompareAndSwap(cur, cur+1) {
return true
}
}
}
// Release 释放一个在途名额。幂等减到 0 为止(防重复释放扣成负数)。
func (p *Pool) Release(uid string) {
p.mu.RLock()
e, ok := p.byUID[uid]
p.mu.RUnlock()
if !ok {
return
}
for {
cur := e.inFlight.Load()
if cur <= 0 {
return
}
if e.inFlight.CompareAndSwap(cur, cur-1) {
return
}
}
}
// SetRandomSource 仅供测试注入确定性随机源;生产代码不应调用。
// 注入源取 n∈[0,n) 后,pickWeighted 的抽签结果完全可预测。
func (p *Pool) SetRandomSource(fn func(n int64) int64) {
p.mu.Lock()
defer p.mu.Unlock()
p.randInt64N = fn
}
// Add 加入账号;已存在则保留原状态、更新凭证(upsert 单账号,不影响其他账号)。
func (p *Pool) Add(a *auth.Auth) {
p.mu.Lock()
defer p.mu.Unlock()
p.upsertLocked(a)
}
// SyncToDir 用最新扫描结果对齐池:新账号加入、消失的账号剔除(状态保留)。
// 剔除结果持久化回 state.json,避免已删账号在下次启动时被 load() 复活。
func (p *Pool) SyncToDir(auths []*auth.Auth) {
p.mu.Lock()
defer p.mu.Unlock()
seen := make(map[string]bool, len(auths))
for _, a := range auths {
seen[a.UID] = true
p.upsertLocked(a)
}
changed := false
for uid := range p.byUID {
if !seen[uid] {
delete(p.byUID, uid)
changed = true
}
}
if changed {
p.saveLocked()
}
}
// Remove 从池中移除账号并立即落盘(管理面板用)。返回被移除账号的凭证
// (含 FilePath,供调用方删除 auth 文件);uid 不存在返回 nil。
// 在途请求的 Release 对已删条目是 no-op,无需等待。
func (p *Pool) Remove(uid string) *auth.Auth {
p.mu.Lock()
defer p.mu.Unlock()
e, ok := p.byUID[uid]
if !ok {
return nil
}
delete(p.byUID, uid)
p.dirty.Store(true)
p.saveLocked()
return e.a
}
// upsertLocked 更新或插入单个账号;已存在则只换凭证、保留 credits/cooling 状态。
// 调用方必须已持有 p.mu;Add 与 SyncToDir 共用此 upsert 逻辑。
func (p *Pool) upsertLocked(a *auth.Auth) {
if e, ok := p.byUID[a.UID]; ok {
e.a = a // 保留 credits/cooling 状态
return
}
p.byUID[a.UID] = &entry{a: a}
}
// Pick 返回 healthy 中积分最高的账号;无可用返回 nil。