mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 14:58:23 +08:00
@@ -141,16 +141,7 @@ func (f *ChannelMonitorQuotaFetcher) Fetch(ctx context.Context, accountID int64)
|
||||
// 避免某个监控的取消波及共享同一账号的其他监控。
|
||||
key := "monitor-quota:" + strconv.FormatInt(accountID, 10)
|
||||
ch := f.flight.DoChan(key, func() (any, error) {
|
||||
fetchCtx, cancel := context.WithTimeout(context.Background(), monitorQuotaFetchTimeout)
|
||||
defer cancel()
|
||||
snapshot := f.fetchUncached(fetchCtx, accountID, time.Now())
|
||||
// 失败也进短 TTL 负缓存:凭据失效/故障期间不必每次调度都打上游。
|
||||
ttl := monitorQuotaFetchCacheTTL
|
||||
if !snapshot.Success {
|
||||
ttl = monitorQuotaErrorCacheTTL
|
||||
}
|
||||
f.storeSnapshot(accountID, snapshot, time.Now().Add(ttl))
|
||||
return snapshot, nil
|
||||
return f.fetchShared(accountID), nil
|
||||
})
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
@@ -164,6 +155,33 @@ func (f *ChannelMonitorQuotaFetcher) Fetch(ctx context.Context, accountID int64)
|
||||
}
|
||||
}
|
||||
|
||||
// fetchShared 是 singleflight 的执行体:抓取一次并写入缓存,结果由同一 key 上
|
||||
// 所有等待者共享。
|
||||
//
|
||||
// 开头必须重查缓存。Fetch 顶部的 cachedSnapshot 与下面的 flight.DoChan 之间有一个
|
||||
// 窗口:期间另一个 goroutine 的 flight 可能已经跑完、写好缓存,并且它的 singleflight
|
||||
// key 也已被摘掉,于是本 goroutine 不会并入那次飞行,而是另起一个新的、对同一账号
|
||||
// 再打一次上游——正是 singleflight 要消除的那种重复查询。
|
||||
//
|
||||
// 这次重查一定命中,所以合并是确定的而不是尽力而为:storeSnapshot 发生在本函数
|
||||
// 返回之前,而 singleflight 删 key 发生在返回之后,因此「能新起一次飞行」必然蕴含
|
||||
// 「上一次的快照已经可见」。
|
||||
func (f *ChannelMonitorQuotaFetcher) fetchShared(accountID int64) *domain.MonitorQuotaSnapshot {
|
||||
if cached, ok := f.cachedSnapshot(accountID, time.Now()); ok {
|
||||
return cached
|
||||
}
|
||||
fetchCtx, cancel := context.WithTimeout(context.Background(), monitorQuotaFetchTimeout)
|
||||
defer cancel()
|
||||
snapshot := f.fetchUncached(fetchCtx, accountID, time.Now())
|
||||
// 失败也进短 TTL 负缓存:凭据失效/故障期间不必每次调度都打上游。
|
||||
ttl := monitorQuotaFetchCacheTTL
|
||||
if !snapshot.Success {
|
||||
ttl = monitorQuotaErrorCacheTTL
|
||||
}
|
||||
f.storeSnapshot(accountID, snapshot, time.Now().Add(ttl))
|
||||
return snapshot
|
||||
}
|
||||
|
||||
func (f *ChannelMonitorQuotaFetcher) cachedSnapshot(accountID int64, now time.Time) (*domain.MonitorQuotaSnapshot, bool) {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
@@ -564,6 +564,48 @@ func TestQuotaFetcher_ConcurrentFetchesShareSingleFlight(t *testing.T) {
|
||||
require.Equal(t, 1, usage.getCalls())
|
||||
}
|
||||
|
||||
// 上面那条并发用例只能以极低概率撞上真正的缺陷窗口(实测约 300 次一次),
|
||||
// 所以这里直接钉住窗口本身。
|
||||
//
|
||||
// 窗口在 Fetch 顶部的 cachedSnapshot 与 flight.DoChan 之间:一个 goroutine 判定
|
||||
// 缓存未命中之后、真正入队之前,另一个 goroutine 的 flight 可能已经跑完、写好缓存
|
||||
// 并被摘掉 key,于是前者不会并入那次飞行,而是另起一次新的,对同一账号重复打上游。
|
||||
// fetchShared 就是 flight 的执行体,直接调用它等价于「已经越过顶层缓存判定、拿到了
|
||||
// 属于自己的那次飞行」这个状态。
|
||||
func TestQuotaFetcher_SharedFetchCacheRecheckAvoidsDuplicateUpstream(t *testing.T) {
|
||||
fetcher, usage, _, _, accounts := newQuotaFetcherTestSetup(t)
|
||||
accounts.accounts[13] = &Account{ID: 13, Platform: domain.PlatformOpenAI}
|
||||
usage.usage = &UsageInfo{FiveHour: &UsageProgress{Utilization: 10}}
|
||||
|
||||
first := fetcher.Fetch(context.Background(), 13)
|
||||
require.True(t, first.Success)
|
||||
require.Equal(t, 1, usage.getCalls())
|
||||
|
||||
shared := fetcher.fetchShared(13)
|
||||
require.Equal(t, 1, usage.getCalls(), "重查缓存后不得再打一次上游")
|
||||
require.Same(t, first, shared, "应原样返回已缓存的快照")
|
||||
}
|
||||
|
||||
// 重查不能变成无条件短路:缓存过期后同一个执行体必须照常回源,
|
||||
// 否则快照会永远停在第一次的值上。
|
||||
func TestQuotaFetcher_SharedFetchStillRefetchesAfterCacheExpiry(t *testing.T) {
|
||||
fetcher, usage, _, _, accounts := newQuotaFetcherTestSetup(t)
|
||||
accounts.accounts[14] = &Account{ID: 14, Platform: domain.PlatformOpenAI}
|
||||
usage.usage = &UsageInfo{FiveHour: &UsageProgress{Utilization: 10}}
|
||||
|
||||
require.True(t, fetcher.Fetch(context.Background(), 14).Success)
|
||||
require.Equal(t, 1, usage.getCalls())
|
||||
|
||||
fetcher.mu.Lock()
|
||||
entry := fetcher.cache[14]
|
||||
entry.expiry = time.Now().Add(-time.Second)
|
||||
fetcher.cache[14] = entry
|
||||
fetcher.mu.Unlock()
|
||||
|
||||
require.True(t, fetcher.fetchShared(14).Success)
|
||||
require.Equal(t, 2, usage.getCalls(), "缓存过期后必须回源")
|
||||
}
|
||||
|
||||
// --- UsageInfo → tiers 归一 ---
|
||||
|
||||
func TestUsageQuotaTiers_MapsAllWindowKinds(t *testing.T) {
|
||||
|
||||
Reference in New Issue
Block a user