From e652f6e20a594a3dd3430de6a49064376eb95a77 Mon Sep 17 00:00:00 2001 From: li Date: Fri, 28 Aug 2026 16:49:49 +0800 Subject: [PATCH] =?UTF-8?q?fix(monitor):=20=E9=85=8D=E9=A2=9D=E6=8A=93?= =?UTF-8?q?=E5=8F=96=E5=9C=A8=20singleflight=20=E6=89=A7=E8=A1=8C=E4=BD=93?= =?UTF-8?q?=E5=86=85=E9=87=8D=E6=9F=A5=E7=BC=93=E5=AD=98=EF=BC=8C=E6=B6=88?= =?UTF-8?q?=E9=99=A4=E9=87=8D=E5=A4=8D=E4=B8=8A=E6=B8=B8=E6=9F=A5=E8=AF=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ChannelMonitorQuotaFetcher.Fetch 顶部的 cachedSnapshot 与随后的 flight.DoChan 之间存在一个窗口:goroutine A 判定缓存未命中之后、真正入队之前,goroutine B 的 flight 可能已经跑完、写好缓存,并且它的 singleflight key 也已被摘掉。A 于是不会 并入 B 的那次飞行,而是另起一个新 flight,对同一账号再打一次上游——正是 singleflight 想要消除的重复查询。 多个监控可以关联同一账号,interval 最小 15s,国产配额服务自身无缓存, 这条路径的全部意义就是防止打爆上游配额端点(见文件头注释)。 修复:把 flight 执行体提为具名方法 fetchShared,开头重查一次缓存。 这次重查必然命中,所以合并是确定的而非尽力而为:storeSnapshot 发生在执行体 返回之前,singleflight 删 key 发生在返回之后,因此「能新起一次飞行」必然蕴含 「上一次的快照已经可见」。 该窗口也是主线 CI 上 TestQuotaFetcher_ConcurrentFetchesShareSingleFlight 间歇性失败(expected: 1, actual: 2)的成因,本地 -count=900 复现 3 次(约 0.33%), 带修复后 -count=1700 零失败。 新增两条确定性用例:既有并发用例只能以约 1/300 的概率撞上该窗口,故直接钉住 窗口本身——fetchShared 就是 flight 的执行体,直接调用它等价于「已越过顶层缓存 判定、拿到属于自己的那次飞行」这一状态。另一条守住重查不得退化成无条件短路: 缓存过期后必须照常回源。 --- .../service/channel_monitor_quota_fetcher.go | 38 ++++++++++++----- .../channel_monitor_quota_fetcher_test.go | 42 +++++++++++++++++++ 2 files changed, 70 insertions(+), 10 deletions(-) diff --git a/backend/internal/service/channel_monitor_quota_fetcher.go b/backend/internal/service/channel_monitor_quota_fetcher.go index 479762ae8d..f8157ddd87 100644 --- a/backend/internal/service/channel_monitor_quota_fetcher.go +++ b/backend/internal/service/channel_monitor_quota_fetcher.go @@ -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() diff --git a/backend/internal/service/channel_monitor_quota_fetcher_test.go b/backend/internal/service/channel_monitor_quota_fetcher_test.go index 6c5b06559b..822009a17e 100644 --- a/backend/internal/service/channel_monitor_quota_fetcher_test.go +++ b/backend/internal/service/channel_monitor_quota_fetcher_test.go @@ -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) {