fix(channel-monitor): 配额快照识别值通道失败并加 60s 负缓存与 singleflight

- fetchUsage 显式识别 UsageInfo.Error/ErrorCode/NeedsReauth/IsBanned/IsForbidden
  (antigravity/grok 等平台失败不走 Go error 通道,此前会被误判为 operational)
- 凭据失效语义(401/403)标记 CredentialInvalid → failed 状态;
  grok quota_unknown 已知未知态豁免不判失败
- 失败快照进 60s 负缓存,避免故障/凭据失效期间以最小 15s 间隔打上游
- 同账号并发抓取由 singleflight 合并(脱离调用方 ctx,45s 总超时兜底)
This commit is contained in:
Randark
2026-08-18 09:14:20 +00:00
parent ebcae03afd
commit 22df600d09
3 changed files with 214 additions and 18 deletions
@@ -85,6 +85,12 @@ const (
// monitorQuotaFetchCacheTTL 配额快照缓存时长。多个监控可能关联同一账号,
// 而 interval 最小 15s 且国产配额服务无缓存,TTL 防止打爆上游配额端点。
monitorQuotaFetchCacheTTL = 5 * time.Minute
// monitorQuotaErrorCacheTTL 失败快照的负缓存时长:失败也短缓存,避免
// 故障/凭据失效期间每次调度(最小 15s)都带真实凭据打上游;到期自动重试。
monitorQuotaErrorCacheTTL = 60 * time.Second
// monitorQuotaFetchTimeout singleflight 内单次配额抓取的总超时
// (脱离调用方 ctx,防止某个监控的取消波及共享同一账号的其他监控)。
monitorQuotaFetchTimeout = 45 * time.Second
// monitorQuotaDegradedUsedPercent 任一用量窗口使用率超过该阈值时,
// 配额检查状态记为 degraded(对齐账号页展示阈值)。
monitorQuotaDegradedUsedPercent = 90.0
@@ -12,6 +12,7 @@ import (
"github.com/Wei-Shaw/sub2api/internal/domain"
"github.com/Wei-Shaw/sub2api/internal/pkg/xai"
"golang.org/x/sync/singleflight"
)
// 渠道监控「配额模式」的配额抓取器。
@@ -26,7 +27,9 @@ import (
// 由 deriveQuotaCheckResult 推导为 failed/error 状态。
//
// 多个监控可能关联同一账号,而 interval 最小 15s 且国产配额服务自身无缓存,
// 所以成功快照统一带 monitorQuotaFetchCacheTTL 缓存,防止打爆上游配额端点。
// 所以快照统一带 TTL 缓存(成功 monitorQuotaFetchCacheTTL、失败
// monitorQuotaErrorCacheTTL 负缓存),防止打爆上游配额端点;同账号的并发
// 抓取由 singleflight 合并为一次上游查询。
// monitorUsageSource 海外平台账号用量查询(AccountUsageService 天然满足)。
type monitorUsageSource interface {
@@ -48,15 +51,17 @@ type monitorAccountSource interface {
GetByID(ctx context.Context, id int64) (*Account, error)
}
// ChannelMonitorQuotaFetcher 配额抓取器(带成功快照 TTL 缓存)。
// ChannelMonitorQuotaFetcher 配额抓取器(成功/失败快照均带 TTL 缓存,
// 同账号并发抓取由 singleflight 合并)。
type ChannelMonitorQuotaFetcher struct {
usage monitorUsageSource
cnQuota monitorCNQuotaSource
cnBalance monitorCNBalanceSource
accounts monitorAccountSource
mu sync.Mutex
cache map[int64]monitorQuotaCacheEntry
mu sync.Mutex
cache map[int64]monitorQuotaCacheEntry
flight singleflight.Group
}
type monitorQuotaCacheEntry struct {
@@ -111,11 +116,31 @@ func (f *ChannelMonitorQuotaFetcher) Fetch(ctx context.Context, accountID int64)
return cached
}
snapshot := f.fetchUncached(ctx, accountID, now)
if snapshot.Success {
f.storeSnapshot(accountID, snapshot, now.Add(monitorQuotaFetchCacheTTL))
// singleflight 合并同账号并发抓取;脱离调用方 ctx(仿 CN 配额服务),
// 避免某个监控的取消波及共享同一账号的其他监控。
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
})
select {
case <-ctx.Done():
return quotaErrorSnapshot("usage", "context canceled", now)
case res := <-ch:
snapshot, ok := res.Val.(*domain.MonitorQuotaSnapshot)
if res.Err != nil || !ok || snapshot == nil {
return quotaErrorSnapshot("usage", "quota fetch failed", now)
}
return snapshot
}
return snapshot
}
func (f *ChannelMonitorQuotaFetcher) cachedSnapshot(accountID int64, now time.Time) (*domain.MonitorQuotaSnapshot, bool) {
@@ -175,6 +200,20 @@ func (f *ChannelMonitorQuotaFetcher) fetchUsage(ctx context.Context, accountID i
FetchedAt: now,
}
}
if usage == nil {
return quotaErrorSnapshot("usage", "usage service returned no data", now)
}
// openai/gemini/antigravity/grok 的失败多走「值通道」(err==nil 但错误
// 降级在 UsageInfo 字段里),必须显式识别,否则会被误判为 operational。
if failed, credInvalid, msg := usageFailureInfo(usage); failed {
return &domain.MonitorQuotaSnapshot{
Source: "usage",
Success: false,
CredentialInvalid: credInvalid,
Error: truncateMessage(sanitizeErrorMessage(msg)),
FetchedAt: now,
}
}
snapshot := &domain.MonitorQuotaSnapshot{
Source: "usage",
Success: true,
@@ -387,6 +426,31 @@ func isCredentialErrorMessage(msg string) bool {
strings.Contains(msg, "authentication")
}
// usageFailureInfo 识别 GetUsage 经「值通道」返回的失败:antigravity/grok
// 等平台 err==nil 但把错误降级在 UsageInfo 字段里(Error/ErrorCode/状态标记)。
// 返回 failed=false 表示可用;credentialInvalid 表示凭据失效(401/403 语义,
// 推导为 failed 状态);msg 为失败摘要。
//
// grok 的 ErrorCode=quota_unknown 是「尚未观测到计费快照/限流头」的已知未知态,
// 不是失败(严格按 ErrorCode 判会把健康 grok 账号永久判 error),显式豁免。
func usageFailureInfo(usage *UsageInfo) (failed, credentialInvalid bool, msg string) {
if usage == nil {
return false, false, ""
}
if usage.ErrorCode == "quota_unknown" {
return false, false, ""
}
failed = usage.Error != "" || usage.NeedsReauth || usage.IsBanned ||
usage.IsForbidden || usage.ErrorCode != ""
if !failed {
return false, false, ""
}
credentialInvalid = usage.NeedsReauth || usage.IsBanned || usage.IsForbidden ||
usage.ErrorCode == errorCodeUnauthenticated || usage.ErrorCode == errorCodeForbidden
msg = firstNonEmpty(usage.Error, usage.ForbiddenReason, usage.ErrorCode, "usage fetch failed")
return failed, credentialInvalid, msg
}
// deriveQuotaCheckResult 把配额快照推导为检测状态(复用既有 status 枚举,
// 时间线/可用率机制自动生效):
// - 查询成功且无告警 → operational
@@ -5,6 +5,7 @@ package service
import (
"context"
"errors"
"sync"
"testing"
"time"
@@ -16,18 +17,33 @@ import (
// --- fetcher 依赖 stub ---
type stubMonitorUsageSource struct {
usage *UsageInfo
err error
usage *UsageInfo
err error
// block 非 nil 时 GetUsage 阻塞在该 channel 上,用于并发/singleflight 测试。
block chan struct{}
mu sync.Mutex
calls int
lastCtx context.Context
}
func (s *stubMonitorUsageSource) GetUsage(ctx context.Context, accountID int64, force ...bool) (*UsageInfo, error) {
s.mu.Lock()
s.calls++
s.lastCtx = ctx
s.mu.Unlock()
if s.block != nil {
<-s.block
}
return s.usage, s.err
}
func (s *stubMonitorUsageSource) getCalls() int {
s.mu.Lock()
defer s.mu.Unlock()
return s.calls
}
type stubMonitorCNQuotaSource struct {
result *CNProviderQuotaProbeResult
err error
@@ -110,7 +126,7 @@ func TestQuotaFetcher_OverseasAccountUsesUsageService(t *testing.T) {
require.NotEmpty(t, fiveHour.ResetAt)
require.Equal(t, "7d", snapshot.Tiers[1].Window)
require.Equal(t, 1, usage.calls)
require.Equal(t, 1, usage.getCalls())
require.Equal(t, 0, cnQuota.calls)
}
@@ -182,7 +198,7 @@ func TestQuotaFetcher_AccountMissingYieldsLinkedAccountSnapshot(t *testing.T) {
require.False(t, snapshot.Success)
require.Equal(t, "linked account not found", snapshot.Error)
require.Equal(t, 0, usage.calls) // 未走到数据源
require.Equal(t, 0, usage.getCalls()) // 未走到数据源
}
func TestQuotaFetcher_UsageAuthErrorMarksCredentialInvalid(t *testing.T) {
@@ -197,6 +213,71 @@ func TestQuotaFetcher_UsageAuthErrorMarksCredentialInvalid(t *testing.T) {
require.Contains(t, snapshot.Error, "401")
}
// 值通道失败:antigravity/grok 等平台 err==nil 但错误降级在 UsageInfo 字段里,
// 必须识别为失败快照,否则会被误判为 operational。
func TestQuotaFetcher_UsageValueChannelFailureYieldsFailureSnapshot(t *testing.T) {
fetcher, usage, _, _, accounts := newQuotaFetcherTestSetup(t)
// 凭据失效(401 语义)→ failed。
accounts.accounts[3] = &Account{ID: 3, Platform: domain.PlatformAnthropic}
usage.usage = &UsageInfo{Error: "usage API error: HTTP 401", ErrorCode: errorCodeUnauthenticated, NeedsReauth: true}
snapshot := fetcher.Fetch(context.Background(), 3)
require.False(t, snapshot.Success)
require.True(t, snapshot.CredentialInvalid)
require.Contains(t, snapshot.Error, "401")
require.Equal(t, MonitorStatusFailed, deriveQuotaCheckResult(snapshot, "quota", time.Now()).Status)
// 限流等非凭据失败 → error(而非 operational)。
accounts.accounts[13] = &Account{ID: 13, Platform: domain.PlatformAnthropic}
usage.usage = &UsageInfo{Error: "usage API error: HTTP 429", ErrorCode: errorCodeRateLimited}
snapshot = fetcher.Fetch(context.Background(), 13)
require.False(t, snapshot.Success)
require.False(t, snapshot.CredentialInvalid)
require.Contains(t, snapshot.Error, "429")
require.Equal(t, MonitorStatusError, deriveQuotaCheckResult(snapshot, "quota", time.Now()).Status)
// grok 已知未知态(尚未观测到计费/限流头)不算失败。
accounts.accounts[14] = &Account{ID: 14, Platform: domain.PlatformGrok}
usage.usage = &UsageInfo{ErrorCode: "quota_unknown", Error: "Grok quota is unknown until billing is probed"}
snapshot = fetcher.Fetch(context.Background(), 14)
require.True(t, snapshot.Success)
require.Empty(t, snapshot.Error)
require.Empty(t, snapshot.Tiers)
require.Equal(t, MonitorStatusOperational, deriveQuotaCheckResult(snapshot, "quota", time.Now()).Status)
}
func TestUsageFailureInfo_ClassificationMatrix(t *testing.T) {
cases := []struct {
name string
usage *UsageInfo
failed bool
credentialInvalid bool
msg string
}{
{name: "nil usage", usage: nil},
{name: "healthy empty", usage: &UsageInfo{}},
{name: "error text only", usage: &UsageInfo{Error: "boom"}, failed: true, msg: "boom"},
{name: "needs reauth", usage: &UsageInfo{NeedsReauth: true}, failed: true, credentialInvalid: true, msg: "usage fetch failed"},
{name: "banned", usage: &UsageInfo{IsBanned: true}, failed: true, credentialInvalid: true, msg: "usage fetch failed"},
{name: "forbidden with reason", usage: &UsageInfo{IsForbidden: true, ForbiddenReason: "usage limited"}, failed: true, credentialInvalid: true, msg: "usage limited"},
{name: "error code unauthenticated", usage: &UsageInfo{ErrorCode: errorCodeUnauthenticated}, failed: true, credentialInvalid: true, msg: errorCodeUnauthenticated},
{name: "error code forbidden", usage: &UsageInfo{ErrorCode: errorCodeForbidden}, failed: true, credentialInvalid: true, msg: errorCodeForbidden},
{name: "error code rate limited", usage: &UsageInfo{ErrorCode: errorCodeRateLimited}, failed: true, msg: errorCodeRateLimited},
{name: "error code network error", usage: &UsageInfo{ErrorCode: errorCodeNetworkError}, failed: true, msg: errorCodeNetworkError},
{name: "grok quota unknown exempted", usage: &UsageInfo{ErrorCode: "quota_unknown", Error: "Grok quota is unknown until billing is probed"}},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
failed, credentialInvalid, msg := usageFailureInfo(tc.usage)
require.Equal(t, tc.failed, failed)
require.Equal(t, tc.credentialInvalid, credentialInvalid)
if tc.msg != "" {
require.Equal(t, tc.msg, msg)
}
})
}
}
func TestQuotaFetcher_CNQuotaCredentialInvalidFlagPropagates(t *testing.T) {
fetcher, _, cnQuota, _, accounts := newQuotaFetcherTestSetup(t)
accounts.accounts[5] = &Account{
@@ -251,7 +332,7 @@ func TestQuotaFetcher_CachesSuccessSnapshotPerAccount(t *testing.T) {
snapshot := fetcher.Fetch(context.Background(), 8)
require.True(t, snapshot.Success)
}
require.Equal(t, 1, usage.calls, "success snapshots should be served from cache")
require.Equal(t, 1, usage.getCalls(), "success snapshots should be served from cache")
// 缓存过期后重新拉取。
fetcher.mu.Lock()
@@ -261,18 +342,63 @@ func TestQuotaFetcher_CachesSuccessSnapshotPerAccount(t *testing.T) {
fetcher.mu.Unlock()
_ = fetcher.Fetch(context.Background(), 8)
require.Equal(t, 2, usage.calls)
require.Equal(t, 2, usage.getCalls())
}
func TestQuotaFetcher_DoesNotCacheFailures(t *testing.T) {
func TestQuotaFetcher_CachesFailureSnapshotWithShortTTL(t *testing.T) {
fetcher, usage, _, _, accounts := newQuotaFetcherTestSetup(t)
accounts.accounts[4] = &Account{ID: 4, Platform: domain.PlatformOpenAI}
usage.err = errors.New("boom")
_ = fetcher.Fetch(context.Background(), 4)
_ = fetcher.Fetch(context.Background(), 4)
for i := 0; i < 2; i++ {
snapshot := fetcher.Fetch(context.Background(), 4)
require.False(t, snapshot.Success)
}
require.Equal(t, 1, usage.getCalls(), "failure snapshots should be served from the short negative cache")
require.Equal(t, 2, usage.calls, "failed snapshots must not be cached")
// 失败快照的 TTL 是负缓存时长(而非成功 TTL)。
fetcher.mu.Lock()
entry := fetcher.cache[4]
require.WithinDuration(t, entry.snapshot.FetchedAt.Add(monitorQuotaErrorCacheTTL), entry.expiry, time.Second)
entry.expiry = time.Now().Add(-time.Second)
fetcher.cache[4] = entry
fetcher.mu.Unlock()
_ = fetcher.Fetch(context.Background(), 4)
require.Equal(t, 2, usage.getCalls(), "expired negative cache should refetch")
}
func TestQuotaFetcher_ConcurrentFetchesShareSingleFlight(t *testing.T) {
fetcher, usage, _, _, accounts := newQuotaFetcherTestSetup(t)
accounts.accounts[12] = &Account{ID: 12, Platform: domain.PlatformOpenAI}
usage.usage = &UsageInfo{FiveHour: &UsageProgress{Utilization: 10}}
usage.block = make(chan struct{})
var wg sync.WaitGroup
snapshots := make([]*domain.MonitorQuotaSnapshot, 5)
for i := 0; i < 5; i++ {
wg.Add(1)
go func(idx int) {
defer wg.Done()
snapshots[idx] = fetcher.Fetch(context.Background(), 12)
}(i)
}
// 上游被 block 卡住时,5 个并发 Fetch 应只产生 1 次真实查询。
require.Eventually(t, func() bool { return usage.getCalls() == 1 },
5*time.Second, 10*time.Millisecond)
close(usage.block)
wg.Wait()
for _, snapshot := range snapshots {
require.NotNil(t, snapshot)
require.True(t, snapshot.Success)
}
require.Equal(t, 1, usage.getCalls())
// 成功快照已缓存:再取一次仍不打上游。
_ = fetcher.Fetch(context.Background(), 12)
require.Equal(t, 1, usage.getCalls())
}
// --- UsageInfo → tiers 归一 ---