File size: 15,266 Bytes
6d60378 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 | // 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。
|