fix(monitor): align quota-fetcher credential/balance semantics with scheduler
P2-1/P2-3 from review: - fetchCNQuota: credential-invalid now judged by StatusCode 401/403 (aligned with fetchCNBalance) instead of `!Success && !CredentialValid` — CN quota service only sets CredentialValid=true on the success path, so 500/429/zhipu business errors were all misclassified as failed instead of error. - fetchCNBalance: snapshot carries new BalanceLow flag computed with the scheduler's exact criterion (`!Available || allCNBalancesBelowThreshold`) against Gateway.CNProviders.BalanceThreshold (ctor now takes cfg; wire regenerated). quotaDegradedHint reports "balance low" instead of the old `<=0` check, so an account already paused by the scheduler (balance 5 / threshold 10) no longer shows green in the monitor. - threshold helper falls back to viper default 0.5 for nil/<=0 config to avoid a zero-threshold regression where balance=0 stops alerting. Tests: CN quota status-code matrix (rewrites the test that cemented the old behavior), balance-low matrix (below-threshold / unavailable / multi-currency healthy), threshold-from-config; PayG stubs now set Available explicitly (zero-value trap).
This commit is contained in:
@@ -337,7 +337,7 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) {
|
||||
batchImageWorkerRuntime := service.ProvideBatchImageWorkerRuntime(batchImageRepository, accountRepository, batchImageQueue, usageBillingRepository, usageLogRepository, batchImageModelPricingResolver, apiKeyAuthCacheInvalidator, configConfig)
|
||||
scheduledTestRunnerService := service.ProvideScheduledTestRunnerService(scheduledTestPlanRepository, scheduledTestService, accountTestService, rateLimitService, configConfig)
|
||||
paymentOrderExpiryService := service.ProvidePaymentOrderExpiryService(paymentService, leaderLockCache, db)
|
||||
channelMonitorQuotaFetcher := service.NewChannelMonitorQuotaFetcher(accountUsageService, cnProviderQuotaService, cnProviderBalanceService, accountRepository)
|
||||
channelMonitorQuotaFetcher := service.NewChannelMonitorQuotaFetcher(accountUsageService, cnProviderQuotaService, cnProviderBalanceService, accountRepository, configConfig)
|
||||
channelMonitorRunner := service.ProvideChannelMonitorRunner(channelMonitorService, settingService, channelMonitorQuotaFetcher)
|
||||
channelMonitorV2Aggregator := service.ProvideChannelMonitorV2Aggregator(channelMonitorV2Repository, db, settingService)
|
||||
userPlatformQuotaUsageFlusher := service.ProvideUserPlatformQuotaUsageFlusher(configConfig, billingCache, serviceUserPlatformQuotaRepository, timingWheelService)
|
||||
|
||||
@@ -50,6 +50,10 @@ type MonitorQuotaSnapshot struct {
|
||||
Balances []MonitorBalance `json:"balances,omitempty"` // 多币种余额(如 DeepSeek CNY+USD)
|
||||
Currency string `json:"currency,omitempty"` // 主余额币种
|
||||
PlanLevel string `json:"plan_level,omitempty"` // 套餐等级(如智谱 level)
|
||||
// BalanceLow 余额低于阈值或账号被上游标记不可用(仅 cn_balance 来源)。
|
||||
// 抓取器按 Gateway.CNProviders.BalanceThreshold 判定,口径与账号停调
|
||||
// (CNProviderBalanceCheckService.checkOne)一致:任一币种达标即健康。
|
||||
BalanceLow bool `json:"balance_low,omitempty"`
|
||||
// CredentialInvalid 上游 401/403 鉴权失败(区别于网络/解析错误),
|
||||
// 检测状态据此推导 failed 而非 error。
|
||||
CredentialInvalid bool `json:"credential_invalid,omitempty"`
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/Wei-Shaw/sub2api/internal/config"
|
||||
"github.com/Wei-Shaw/sub2api/internal/domain"
|
||||
"github.com/Wei-Shaw/sub2api/internal/pkg/xai"
|
||||
"golang.org/x/sync/singleflight"
|
||||
@@ -58,6 +59,8 @@ type ChannelMonitorQuotaFetcher struct {
|
||||
cnQuota monitorCNQuotaSource
|
||||
cnBalance monitorCNBalanceSource
|
||||
accounts monitorAccountSource
|
||||
// balanceThreshold cn_balance 余额告警阈值(与账号停调共用配置,见 monitorBalanceThreshold)。
|
||||
balanceThreshold float64
|
||||
|
||||
mu sync.Mutex
|
||||
cache map[int64]monitorQuotaCacheEntry
|
||||
@@ -76,8 +79,12 @@ func NewChannelMonitorQuotaFetcher(
|
||||
cnQuota *CNProviderQuotaService,
|
||||
cnBalance *CNProviderBalanceService,
|
||||
accounts AccountRepository,
|
||||
cfg *config.Config,
|
||||
) *ChannelMonitorQuotaFetcher {
|
||||
f := &ChannelMonitorQuotaFetcher{cache: make(map[int64]monitorQuotaCacheEntry)}
|
||||
f := &ChannelMonitorQuotaFetcher{
|
||||
cache: make(map[int64]monitorQuotaCacheEntry),
|
||||
balanceThreshold: monitorBalanceThreshold(cfg),
|
||||
}
|
||||
if usage != nil {
|
||||
f.usage = usage
|
||||
}
|
||||
@@ -93,6 +100,17 @@ func NewChannelMonitorQuotaFetcher(
|
||||
return f
|
||||
}
|
||||
|
||||
// monitorBalanceThreshold 余额告警阈值,与账号停调(CNProviderBalanceCheckService)
|
||||
// 共用 gateway.cn_providers.balance_threshold,保证监控 degraded 与调度器停调
|
||||
// 口径一致(任一币种达标即健康)。未配置/非正值时回退 viper 默认 0.5(config.go),
|
||||
// 避免 0 阈值下「余额=0 也不告警」相对旧 `<=0` 判定的回归。
|
||||
func monitorBalanceThreshold(cfg *config.Config) float64 {
|
||||
if cfg != nil && cfg.Gateway.CNProviders.BalanceThreshold > 0 {
|
||||
return cfg.Gateway.CNProviders.BalanceThreshold
|
||||
}
|
||||
return 0.5
|
||||
}
|
||||
|
||||
// LoadAccount 加载账号(不走缓存)。供 Create/Update 时校验
|
||||
// provider 与 account.platform 一致;账号不存在时返回错误。
|
||||
func (f *ChannelMonitorQuotaFetcher) LoadAccount(ctx context.Context, id int64) (*Account, error) {
|
||||
@@ -341,7 +359,10 @@ func (f *ChannelMonitorQuotaFetcher) fetchCNQuota(ctx context.Context, accountID
|
||||
Error: result.Error,
|
||||
FetchedAt: now,
|
||||
}
|
||||
if !result.Success && !result.CredentialValid {
|
||||
// 只有 401/403 判凭据失效(与 fetchCNBalance 口径一致):CN quota 服务的
|
||||
// CredentialValid 仅在成功路径置 true,若按 `!Success && !CredentialValid`
|
||||
// 推导,500/429/智谱业务错误全会被误判为 failed。
|
||||
if !result.Success && (result.StatusCode == 401 || result.StatusCode == 403) {
|
||||
snapshot.CredentialInvalid = true
|
||||
}
|
||||
if len(result.Tiers) > 0 {
|
||||
@@ -386,6 +407,10 @@ func (f *ChannelMonitorQuotaFetcher) fetchCNBalance(ctx context.Context, account
|
||||
if result.Success {
|
||||
balance := result.Balance
|
||||
snapshot.Balance = &balance
|
||||
// 与账号停调(checkOne)同口径:上游标记不可用或全部币种低于阈值
|
||||
// 才告警,任一币种达标即健康(余额 5 元/阈值 10 元的账号调度器已
|
||||
// 停调,监控不能仍绿灯)。
|
||||
snapshot.BalanceLow = !result.Available || allCNBalancesBelowThreshold(result, f.balanceThreshold)
|
||||
} else if result.StatusCode == 401 || result.StatusCode == 403 {
|
||||
snapshot.CredentialInvalid = true
|
||||
}
|
||||
@@ -454,7 +479,7 @@ func usageFailureInfo(usage *UsageInfo) (failed, credentialInvalid bool, msg str
|
||||
// deriveQuotaCheckResult 把配额快照推导为检测状态(复用既有 status 枚举,
|
||||
// 时间线/可用率机制自动生效):
|
||||
// - 查询成功且无告警 → operational
|
||||
// - 任一窗口使用率 >= 阈值或余额耗尽 → degraded
|
||||
// - 任一窗口使用率 >= 阈值或余额低于阈值/不可用 → degraded
|
||||
// - 账号未关联(配置问题) → degraded
|
||||
// - 凭据失效(401/403) → failed
|
||||
// - 网络/解析等其他错误 → error
|
||||
@@ -499,8 +524,11 @@ func quotaDegradedHint(snapshot *domain.MonitorQuotaSnapshot) string {
|
||||
return fmt.Sprintf("quota high: %s at %s%%", name, strconv.FormatFloat(tier.UsedPercent, 'f', 1, 64))
|
||||
}
|
||||
}
|
||||
if snapshot.Balance != nil && *snapshot.Balance <= 0 {
|
||||
return fmt.Sprintf("balance depleted (%s)", firstNonEmpty(snapshot.Currency, "?"))
|
||||
if snapshot.BalanceLow {
|
||||
if snapshot.Balance != nil {
|
||||
return fmt.Sprintf("balance low: %s %s", strconv.FormatFloat(*snapshot.Balance, 'f', -1, 64), firstNonEmpty(snapshot.Currency, "?"))
|
||||
}
|
||||
return fmt.Sprintf("balance low (%s)", firstNonEmpty(snapshot.Currency, "?"))
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Wei-Shaw/sub2api/internal/config"
|
||||
"github.com/Wei-Shaw/sub2api/internal/domain"
|
||||
"github.com/Wei-Shaw/sub2api/internal/pkg/xai"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -87,11 +88,12 @@ func newQuotaFetcherTestSetup(t *testing.T) (*ChannelMonitorQuotaFetcher, *stubM
|
||||
cnBalance := &stubMonitorCNBalanceSource{}
|
||||
accounts := &stubMonitorAccountSource{accounts: make(map[int64]*Account)}
|
||||
fetcher := &ChannelMonitorQuotaFetcher{
|
||||
usage: usage,
|
||||
cnQuota: cnQuota,
|
||||
cnBalance: cnBalance,
|
||||
accounts: accounts,
|
||||
cache: make(map[int64]monitorQuotaCacheEntry),
|
||||
usage: usage,
|
||||
cnQuota: cnQuota,
|
||||
cnBalance: cnBalance,
|
||||
accounts: accounts,
|
||||
balanceThreshold: monitorBalanceThreshold(nil),
|
||||
cache: make(map[int64]monitorQuotaCacheEntry),
|
||||
}
|
||||
return fetcher, usage, cnQuota, cnBalance, accounts
|
||||
}
|
||||
@@ -167,9 +169,10 @@ func TestQuotaFetcher_PayGAccountUsesCNBalance(t *testing.T) {
|
||||
Credentials: map[string]any{"account_mode": AccountModePayG},
|
||||
}
|
||||
cnBalance.result = &CNProviderBalanceResult{
|
||||
Success: true,
|
||||
Balance: 12.34,
|
||||
Currency: "CNY",
|
||||
Success: true,
|
||||
Available: true,
|
||||
Balance: 12.34,
|
||||
Currency: "CNY",
|
||||
Balances: []CNProviderBalanceEntry{
|
||||
{Currency: "CNY", Balance: 12.34},
|
||||
{Currency: "USD", Balance: 1.5},
|
||||
@@ -185,6 +188,7 @@ func TestQuotaFetcher_PayGAccountUsesCNBalance(t *testing.T) {
|
||||
require.Equal(t, "CNY", snapshot.Currency)
|
||||
require.Len(t, snapshot.Balances, 2)
|
||||
require.Equal(t, "USD", snapshot.Balances[1].Currency)
|
||||
require.False(t, snapshot.BalanceLow)
|
||||
require.Empty(t, snapshot.Error)
|
||||
}
|
||||
|
||||
@@ -278,20 +282,44 @@ func TestUsageFailureInfo_ClassificationMatrix(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestQuotaFetcher_CNQuotaCredentialInvalidFlagPropagates(t *testing.T) {
|
||||
fetcher, _, cnQuota, _, accounts := newQuotaFetcherTestSetup(t)
|
||||
accounts.accounts[5] = &Account{
|
||||
ID: 5,
|
||||
Platform: domain.PlatformZhipu,
|
||||
Credentials: map[string]any{"account_mode": AccountModeCoding},
|
||||
// 凭据失效只认 401/403(与 fetchCNBalance 口径一致):CN quota 服务的
|
||||
// CredentialValid 仅成功路径置 true,500/429/智谱业务错误须推导为 error 而非 failed。
|
||||
func TestQuotaFetcher_CNQuotaCredentialInvalidByStatusCode(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
accountID int64
|
||||
statusCode int
|
||||
credentialBad bool
|
||||
expectedStatus string
|
||||
}{
|
||||
{name: "401 unauthorized", accountID: 5, statusCode: 401, credentialBad: true, expectedStatus: MonitorStatusFailed},
|
||||
{name: "403 forbidden", accountID: 15, statusCode: 403, credentialBad: true, expectedStatus: MonitorStatusFailed},
|
||||
{name: "500 server error", accountID: 16, statusCode: 500, expectedStatus: MonitorStatusError},
|
||||
{name: "429 rate limited", accountID: 17, statusCode: 429, expectedStatus: MonitorStatusError},
|
||||
// 智谱 2xx 但业务级失败:StatusCode=200,非凭据问题。
|
||||
{name: "200 business error", accountID: 18, statusCode: 200, expectedStatus: MonitorStatusError},
|
||||
}
|
||||
cnQuota.result = &CNProviderQuotaProbeResult{Success: false, CredentialValid: false, Error: "api key expired"}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fetcher, _, cnQuota, _, accounts := newQuotaFetcherTestSetup(t)
|
||||
accounts.accounts[tc.accountID] = &Account{
|
||||
ID: tc.accountID,
|
||||
Platform: domain.PlatformZhipu,
|
||||
Credentials: map[string]any{"account_mode": AccountModeCoding},
|
||||
}
|
||||
cnQuota.result = &CNProviderQuotaProbeResult{
|
||||
Success: false,
|
||||
StatusCode: tc.statusCode,
|
||||
Error: "api key expired",
|
||||
}
|
||||
|
||||
snapshot := fetcher.Fetch(context.Background(), 5)
|
||||
snapshot := fetcher.Fetch(context.Background(), tc.accountID)
|
||||
|
||||
require.False(t, snapshot.Success)
|
||||
require.True(t, snapshot.CredentialInvalid)
|
||||
require.Equal(t, "api key expired", snapshot.Error)
|
||||
require.False(t, snapshot.Success)
|
||||
require.Equal(t, tc.credentialBad, snapshot.CredentialInvalid)
|
||||
require.Equal(t, tc.expectedStatus, deriveQuotaCheckResult(snapshot, "quota", time.Now()).Status)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestQuotaFetcher_CNBalanceHTTP403MarksCredentialInvalid(t *testing.T) {
|
||||
@@ -305,6 +333,86 @@ func TestQuotaFetcher_CNBalanceHTTP403MarksCredentialInvalid(t *testing.T) {
|
||||
require.True(t, snapshot.CredentialInvalid)
|
||||
}
|
||||
|
||||
// 余额告警口径与账号停调(CNProviderBalanceCheckService.checkOne)一致:
|
||||
// 上游标记不可用或全部币种低于阈值 → BalanceLow → degraded;任一币种达标即健康。
|
||||
func TestQuotaFetcher_CNBalanceLowMarksDegraded(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
accountID int64
|
||||
result *CNProviderBalanceResult
|
||||
balanceLow bool
|
||||
wantStatus string
|
||||
wantMessage string
|
||||
}{
|
||||
{
|
||||
// 审查例:余额 5/阈值 10 的账号调度器已停调,监控不能仍绿灯。
|
||||
name: "balance below threshold",
|
||||
accountID: 21,
|
||||
result: &CNProviderBalanceResult{Success: true, Available: true, Balance: 5, Currency: "CNY"},
|
||||
balanceLow: true,
|
||||
wantStatus: MonitorStatusDegraded, wantMessage: "balance low: 5 CNY",
|
||||
},
|
||||
{
|
||||
name: "upstream marked unavailable",
|
||||
accountID: 22,
|
||||
result: &CNProviderBalanceResult{Success: true, Available: false, Balance: 20, Currency: "CNY"},
|
||||
balanceLow: true,
|
||||
wantStatus: MonitorStatusDegraded, wantMessage: "balance low: 20 CNY",
|
||||
},
|
||||
{
|
||||
// deepseek 双币种:任一币种(USD 20)达标即健康。
|
||||
name: "any currency above threshold is healthy",
|
||||
accountID: 23,
|
||||
result: &CNProviderBalanceResult{
|
||||
Success: true, Available: true, Balance: 5, Currency: "CNY",
|
||||
Balances: []CNProviderBalanceEntry{{Currency: "CNY", Balance: 5}, {Currency: "USD", Balance: 20}},
|
||||
},
|
||||
wantStatus: MonitorStatusOperational,
|
||||
},
|
||||
{
|
||||
name: "single currency above threshold",
|
||||
accountID: 24,
|
||||
result: &CNProviderBalanceResult{Success: true, Available: true, Balance: 20, Currency: "CNY"},
|
||||
wantStatus: MonitorStatusOperational,
|
||||
},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
fetcher, _, _, cnBalance, accounts := newQuotaFetcherTestSetup(t)
|
||||
fetcher.balanceThreshold = 10
|
||||
accounts.accounts[tc.accountID] = &Account{
|
||||
ID: tc.accountID,
|
||||
Platform: domain.PlatformKimi,
|
||||
Credentials: map[string]any{"account_mode": AccountModePayG},
|
||||
}
|
||||
cnBalance.result = tc.result
|
||||
|
||||
snapshot := fetcher.Fetch(context.Background(), tc.accountID)
|
||||
|
||||
require.True(t, snapshot.Success)
|
||||
require.Equal(t, tc.balanceLow, snapshot.BalanceLow)
|
||||
res := deriveQuotaCheckResult(snapshot, "quota", time.Now())
|
||||
require.Equal(t, tc.wantStatus, res.Status)
|
||||
if tc.wantMessage != "" {
|
||||
require.Contains(t, res.Message, tc.wantMessage)
|
||||
} else {
|
||||
require.Empty(t, res.Message)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewChannelMonitorQuotaFetcher_ThresholdFromConfig(t *testing.T) {
|
||||
require.InDelta(t, 0.5, NewChannelMonitorQuotaFetcher(nil, nil, nil, nil, nil).balanceThreshold, 0.0001)
|
||||
|
||||
cfg10 := &config.Config{Gateway: config.GatewayConfig{CNProviders: config.GatewayCNProvidersConfig{BalanceThreshold: 10}}}
|
||||
require.InDelta(t, 10, NewChannelMonitorQuotaFetcher(nil, nil, nil, nil, cfg10).balanceThreshold, 0.0001)
|
||||
|
||||
// 非正值(含显式 0)回退默认,避免 0 阈值下「余额=0 也不告警」。
|
||||
cfg0 := &config.Config{Gateway: config.GatewayConfig{CNProviders: config.GatewayCNProvidersConfig{BalanceThreshold: 0}}}
|
||||
require.InDelta(t, 0.5, NewChannelMonitorQuotaFetcher(nil, nil, nil, nil, cfg0).balanceThreshold, 0.0001)
|
||||
}
|
||||
|
||||
func TestQuotaFetcher_NilDependenciesProduceErrorSnapshots(t *testing.T) {
|
||||
// fetcher 本体为 nil:直接降级为错误快照,不 panic。
|
||||
var nilFetcher *ChannelMonitorQuotaFetcher
|
||||
@@ -495,10 +603,10 @@ func TestDeriveQuotaCheckResult_StatusMatrix(t *testing.T) {
|
||||
require.Contains(t, res.Message, "95.0%")
|
||||
|
||||
balance := -0.5
|
||||
depleted := &domain.MonitorQuotaSnapshot{Success: true, Balance: &balance, Currency: "CNY"}
|
||||
res = deriveQuotaCheckResult(depleted, "quota", now)
|
||||
lowBalance := &domain.MonitorQuotaSnapshot{Success: true, BalanceLow: true, Balance: &balance, Currency: "CNY"}
|
||||
res = deriveQuotaCheckResult(lowBalance, "quota", now)
|
||||
require.Equal(t, MonitorStatusDegraded, res.Status)
|
||||
require.Contains(t, res.Message, "balance depleted")
|
||||
require.Contains(t, res.Message, "balance low")
|
||||
|
||||
invalid := &domain.MonitorQuotaSnapshot{Success: false, CredentialInvalid: true, Error: "401 unauthorized"}
|
||||
res = deriveQuotaCheckResult(invalid, "quota", now)
|
||||
|
||||
Reference in New Issue
Block a user