diff --git a/backend/internal/domain/channel_monitor_quota.go b/backend/internal/domain/channel_monitor_quota.go index 656163dddf..fed21dedb7 100644 --- a/backend/internal/domain/channel_monitor_quota.go +++ b/backend/internal/domain/channel_monitor_quota.go @@ -20,11 +20,17 @@ import "time" // - "7d-sonnet" Claude 7 天 Sonnet 独立额度 // - "7d-fable" Claude 7 天 Fable 独立额度 // - "weekly" 周窗口(Kimi/Zhipu coding plan) +// - "daily" 日窗口(Gemini 日配额 / Grok 日请求) // - "30d" 30 天窗口(Grok 月度) -// - "total" 无窗口语义的总量额度(Gemini/Antigravity 等) +// - "total" 无窗口语义的总量额度(Antigravity per-model 等) +// +// 同一 Window 可能出现多条(Gemini 多档日配额、Antigravity per-model、 +// Grok requests/tokens),用 Label 区分:Label 是机器 token(requests/tokens/ +// shared/pro/flash 或模型名),前端已知 token 走 i18n,未知原样展示。 type MonitorQuotaTier struct { Window string `json:"window"` - UsedPercent float64 `json:"used_percent"` // 0-100;仅有绝对值时按 used/limit 计算 + Label string `json:"label,omitempty"` + UsedPercent float64 `json:"used_percent"` // 0-100+;仅有绝对值时按 used/limit 计算 Used float64 `json:"used,omitempty"` Limit float64 `json:"limit,omitempty"` ResetAt string `json:"reset_at,omitempty"` // RFC3339;未知时留空 @@ -44,8 +50,11 @@ type MonitorQuotaSnapshot struct { Balances []MonitorBalance `json:"balances,omitempty"` // 多币种余额(如 DeepSeek CNY+USD) Currency string `json:"currency,omitempty"` // 主余额币种 PlanLevel string `json:"plan_level,omitempty"` // 套餐等级(如智谱 level) - Error string `json:"error,omitempty"` // Success=false 时的错误摘要 - FetchedAt time.Time `json:"fetched_at"` + // CredentialInvalid 上游 401/403 鉴权失败(区别于网络/解析错误), + // 检测状态据此推导 failed 而非 error。 + CredentialInvalid bool `json:"credential_invalid,omitempty"` + Error string `json:"error,omitempty"` // Success=false 时的错误摘要 + FetchedAt time.Time `json:"fetched_at"` } // MonitorBalance 单币种余额条目。 diff --git a/backend/internal/repository/channel_monitor_repo.go b/backend/internal/repository/channel_monitor_repo.go index aa8b12ff8b..c4586bc38c 100644 --- a/backend/internal/repository/channel_monitor_repo.go +++ b/backend/internal/repository/channel_monitor_repo.go @@ -3,6 +3,7 @@ package repository import ( "context" "database/sql" + "encoding/json" "fmt" "strings" "time" @@ -10,6 +11,7 @@ import ( dbent "github.com/Wei-Shaw/sub2api/ent" "github.com/Wei-Shaw/sub2api/ent/channelmonitor" "github.com/Wei-Shaw/sub2api/ent/channelmonitorhistory" + "github.com/Wei-Shaw/sub2api/internal/domain" "github.com/Wei-Shaw/sub2api/internal/service" "github.com/lib/pq" @@ -51,10 +53,14 @@ func (r *channelMonitorRepository) Create(ctx context.Context, m *service.Channe SetJitterSeconds(m.JitterSeconds). SetCreatedBy(m.CreatedBy). SetExtraHeaders(channelMonitorHeadersForPersistence(m)). - SetBodyOverrideMode(defaultBodyModeRepo(m.BodyOverrideMode)) + SetBodyOverrideMode(defaultBodyModeRepo(m.BodyOverrideMode)). + SetCheckMode(defaultCheckModeRepo(m.CheckMode)) if m.TemplateID != nil { builder = builder.SetTemplateID(*m.TemplateID) } + if m.AccountID != nil { + builder = builder.SetAccountID(*m.AccountID) + } if m.BodyOverride != nil { builder = builder.SetBodyOverride(m.BodyOverride) } @@ -118,12 +124,18 @@ func (r *channelMonitorRepository) Update(ctx context.Context, m *service.Channe SetIntervalSeconds(m.IntervalSeconds). SetJitterSeconds(m.JitterSeconds). SetExtraHeaders(channelMonitorHeadersForPersistence(m)). - SetBodyOverrideMode(defaultBodyModeRepo(m.BodyOverrideMode)) + SetBodyOverrideMode(defaultBodyModeRepo(m.BodyOverrideMode)). + SetCheckMode(defaultCheckModeRepo(m.CheckMode)) if m.TemplateID != nil { updater = updater.SetTemplateID(*m.TemplateID) } else { updater = updater.ClearTemplateID() } + if m.AccountID != nil { + updater = updater.SetAccountID(*m.AccountID) + } else { + updater = updater.ClearAccountID() + } if m.BodyOverride != nil { updater = updater.SetBodyOverride(m.BodyOverride) } else { @@ -237,6 +249,9 @@ func (r *channelMonitorRepository) InsertHistoryBatch(ctx context.Context, rows if row.PingLatencyMs != nil { c = c.SetPingLatencyMs(*row.PingLatencyMs) } + if row.Quota != nil { + c = c.SetQuota(row.Quota) + } bulk = append(bulk, c) } if _, err := client.ChannelMonitorHistory.CreateBulk(bulk...).Save(ctx); err != nil { @@ -276,6 +291,7 @@ func (r *channelMonitorRepository) ListHistory(ctx context.Context, monitorID in PingLatencyMs: row.PingLatencyMs, Message: row.Message, CheckedAt: row.CheckedAt, + Quota: row.Quota, } out = append(out, entry) } @@ -324,6 +340,20 @@ func assignNullInt(dst **int, n sql.NullInt64) { *dst = &v } +// scanMonitorQuota 把裸 SQL 读出的 JSONB quota 列解包为配额快照。 +// NULL(探活模式旧行)返回 nil;解析失败也返回 nil 并由调用方日志感知, +// 不阻断列表渲染(与聚合层"失败仅日志"的原则一致)。 +func scanMonitorQuota(data []byte) *domain.MonitorQuotaSnapshot { + if len(data) == 0 { + return nil + } + snapshot := &domain.MonitorQuotaSnapshot{} + if err := json.Unmarshal(data, snapshot); err != nil { + return nil + } + return snapshot +} + // ComputeAvailability 计算指定窗口内每个模型的可用率与平均延迟。 // "可用" = status IN (operational, degraded)。 // @@ -396,7 +426,7 @@ func (r *channelMonitorRepository) ListLatestForMonitorIDs(ctx context.Context, } const q = ` SELECT DISTINCT ON (monitor_id, model) - monitor_id, model, status, latency_ms, ping_latency_ms, checked_at + monitor_id, model, status, latency_ms, ping_latency_ms, checked_at, quota FROM channel_monitor_histories WHERE monitor_id = ANY($1) ORDER BY monitor_id, model, checked_at DESC @@ -411,11 +441,13 @@ func (r *channelMonitorRepository) ListLatestForMonitorIDs(ctx context.Context, var monitorID int64 l := &service.ChannelMonitorLatest{} var latency, ping sql.NullInt64 - if err := rows.Scan(&monitorID, &l.Model, &l.Status, &latency, &ping, &l.CheckedAt); err != nil { + var quota []byte + if err := rows.Scan(&monitorID, &l.Model, &l.Status, &latency, &ping, &l.CheckedAt, "a); err != nil { return nil, fmt.Errorf("scan latest batch row: %w", err) } assignNullInt(&l.LatencyMs, latency) assignNullInt(&l.PingLatencyMs, ping) + l.Quota = scanMonitorQuota(quota) out[monitorID] = append(out[monitorID], l) } if err := rows.Err(); err != nil { @@ -757,12 +789,17 @@ func entToServiceMonitor(row *dbent.ChannelMonitor) *service.ChannelMonitor { ExtraHeaders: headers, BodyOverrideMode: row.BodyOverrideMode, BodyOverride: row.BodyOverride, + CheckMode: defaultCheckModeRepo(row.CheckMode), DuplicateOperationID: duplicateOperationID, } if row.TemplateID != nil { id := *row.TemplateID out.TemplateID = &id } + if row.AccountID != nil { + id := *row.AccountID + out.AccountID = &id + } return out } @@ -807,6 +844,14 @@ func defaultAPIModeRepo(apiMode string) string { return apiMode } +// defaultCheckModeRepo 空串归一为 probe(存量行有列默认值,这里兜底防御)。 +func defaultCheckModeRepo(checkMode string) string { + if checkMode == "" { + return "probe" + } + return checkMode +} + func emptySliceIfNil(in []string) []string { if in == nil { return []string{} diff --git a/backend/internal/service/channel_monitor_aggregator.go b/backend/internal/service/channel_monitor_aggregator.go index 09020f5fa4..b6f632bf4b 100644 --- a/backend/internal/service/channel_monitor_aggregator.go +++ b/backend/internal/service/channel_monitor_aggregator.go @@ -204,6 +204,8 @@ func buildStatusSummary( if l, ok := latestByModel[primary]; ok { summary.PrimaryStatus = l.Status summary.PrimaryLatencyMs = l.LatencyMs + // 配额快照只挂主模型行(quota 模式唯一行 / quota_probe 的主行)。 + summary.LatestQuota = l.Quota } if a, ok := availByModel[primary]; ok { summary.Availability7d = a.AvailabilityPct @@ -242,6 +244,7 @@ func buildUserViewFromSummary( } if primaryLatest != nil { view.PrimaryPingLatencyMs = primaryLatest.PingLatencyMs + view.LatestQuota = primaryLatest.Quota } return view } diff --git a/backend/internal/service/channel_monitor_checker.go b/backend/internal/service/channel_monitor_checker.go index 910a6e80d9..bc8aa70645 100644 --- a/backend/internal/service/channel_monitor_checker.go +++ b/backend/internal/service/channel_monitor_checker.go @@ -170,6 +170,11 @@ type providerAdapter struct { var providerAdapters = map[string]providerAdapter{ MonitorProviderOpenAI: providerOpenAIChatAdapter, MonitorProviderGrok: providerGrokChatAdapter, + // 国产 3 家(配额模式引入):均为 OpenAI 兼容 Chat Completions, + // 仅智谱路径前缀不同(/api/paas/v4/chat/completions)。 + MonitorProviderKimi: providerKimiChatAdapter, + MonitorProviderZhipu: providerZhipuChatAdapter, + MonitorProviderDeepseek: providerDeepseekChatAdapter, MonitorProviderAnthropic: { buildPath: func(string) string { return providerAnthropicPath }, buildBody: func(model, prompt string) ([]byte, error) { @@ -212,6 +217,15 @@ var providerOpenAIChatAdapter = newOpenAICompatibleChatAdapter(providerOpenAIPat //nolint:gochecknoglobals // 适配器表是只读静态数据,初始化后不变更。 var providerGrokChatAdapter = newOpenAICompatibleChatAdapter(providerGrokPath) +//nolint:gochecknoglobals // 适配器表是只读静态数据,初始化后不变更。 +var providerKimiChatAdapter = newOpenAICompatibleChatAdapter(providerOpenAIPath) + +//nolint:gochecknoglobals // 适配器表是只读静态数据,初始化后不变更。 +var providerZhipuChatAdapter = newOpenAICompatibleChatAdapter(providerZhipuPath) + +//nolint:gochecknoglobals // 适配器表是只读静态数据,初始化后不变更。 +var providerDeepseekChatAdapter = newOpenAICompatibleChatAdapter(providerOpenAIPath) + func newOpenAICompatibleChatAdapter(path string) providerAdapter { return providerAdapter{ buildPath: func(string) string { return path }, @@ -257,8 +271,8 @@ func providerAdapterFor(provider, apiMode string) (providerAdapter, string, bool return adapter, MonitorAPIModeChatCompletions, ok } -// isSupportedProvider 校验 provider 字符串是否在 adapter 表中。 -// 供 validate.go 的 validateProvider 复用,避免两份 switch 漂移。 +// isSupportedProvider 校验 provider 字符串是否在探活 adapter 表中。 +// 注意:完整的 provider 合法性(含仅配额模式的 antigravity)走 validate.go 的 monitorProviders。 func isSupportedProvider(p string) bool { _, ok := providerAdapters[p] return ok @@ -447,6 +461,10 @@ var bodyMergeKeyDenyList = map[string]map[string]bool{ MonitorProviderGrok: {"model": true, "messages": true, "stream": true}, MonitorProviderAnthropic: {"model": true, "messages": true}, MonitorProviderGemini: {"contents": true}, + // 国产 3 家与 OpenAI Chat Completions 同构。 + MonitorProviderKimi: {"model": true, "messages": true, "stream": true}, + MonitorProviderZhipu: {"model": true, "messages": true, "stream": true}, + MonitorProviderDeepseek: {"model": true, "messages": true, "stream": true}, } func checkAPIMode(opts *CheckOptions) string { @@ -463,8 +481,20 @@ func bodyMergeDenyKey(provider, apiMode string) string { return provider } +// isOpenAICompatibleChatProvider 该 provider 的探活请求是否为 OpenAI Chat +// Completions 同构(replace 模式的 body 校验按 messages 必填处理)。 +func isOpenAICompatibleChatProvider(provider string) bool { + switch provider { + case MonitorProviderOpenAI, MonitorProviderGrok, + MonitorProviderKimi, MonitorProviderZhipu, MonitorProviderDeepseek: + return true + default: + return false + } +} + func validateReplaceRequestBody(provider, apiMode string, body map[string]any) error { - if provider != MonitorProviderOpenAI && provider != MonitorProviderGrok { + if !isOpenAICompatibleChatProvider(provider) { return nil } switch defaultAPIMode(apiMode) { diff --git a/backend/internal/service/channel_monitor_const.go b/backend/internal/service/channel_monitor_const.go index 6a41add298..4ed6a026ea 100644 --- a/backend/internal/service/channel_monitor_const.go +++ b/backend/internal/service/channel_monitor_const.go @@ -45,10 +45,12 @@ const ( monitorChallengeMin = 1 monitorChallengeMax = 50 - // providerOpenAIPath OpenAI Chat Completions 路径。 + // providerOpenAIPath OpenAI Chat Completions 路径(Kimi / DeepSeek 同为 OpenAI 兼容)。 providerOpenAIPath = "/v1/chat/completions" // providerGrokPath Grok OpenAI-compatible Chat Completions 路径。 providerGrokPath = "/v1/chat/completions" + // providerZhipuPath 智谱 OpenAI 兼容 Chat Completions 路径(前缀与官方不同)。 + providerZhipuPath = "/api/paas/v4/chat/completions" // providerOpenAIResponsesPath OpenAI Responses API 路径。 providerOpenAIResponsesPath = "/v1/responses" // providerAnthropicPath Anthropic Messages 路径。 @@ -56,11 +58,36 @@ const ( // providerGeminiPathTemplate Gemini generateContent 路径模板(含 model 占位)。 providerGeminiPathTemplate = "/v1beta/models/%s:generateContent" - // MonitorProviderOpenAI / Anthropic / Gemini / Grok provider 字符串常量(也是 ent enum 的实际值)。 - MonitorProviderOpenAI = "openai" - MonitorProviderAnthropic = "anthropic" - MonitorProviderGemini = "gemini" - MonitorProviderGrok = "grok" + // MonitorProviderOpenAI 等 provider 字符串常量(也是 ent enum 的实际值)。 + // 后 4 个 provider(antigravity/kimi/zhipu/deepseek)为配额模式引入: + // antigravity 无探活 adapter(仅配额),其余 3 个复用 OpenAI 兼容探活。 + MonitorProviderOpenAI = "openai" + MonitorProviderAnthropic = "anthropic" + MonitorProviderGemini = "gemini" + MonitorProviderGrok = "grok" + MonitorProviderAntigravity = "antigravity" + MonitorProviderKimi = "kimi" + MonitorProviderZhipu = "zhipu" + MonitorProviderDeepseek = "deepseek" + + // MonitorCheckMode 检测模式(channel_monitors.check_mode)。 + // probe - LLM 探活(默认,原有行为) + // quota - 仅查关联账号用量/余额,零 LLM 成本 + // quota_probe - 探活 + 配额并存(配额快照挂到主模型历史行) + MonitorCheckModeProbe = "probe" + MonitorCheckModeQuota = "quota" + MonitorCheckModeQuotaProbe = "quota_probe" + + // MonitorDefaultQuotaModel 是 quota 模式监控未显式指定模型时占位的虚拟模型名 + // (primary_model 列 NotEmpty,用 "quota" 让历史行/时间线机制无需特判)。 + MonitorDefaultQuotaModel = "quota" + + // monitorQuotaFetchCacheTTL 配额快照缓存时长。多个监控可能关联同一账号, + // 而 interval 最小 15s 且国产配额服务无缓存,TTL 防止打爆上游配额端点。 + monitorQuotaFetchCacheTTL = 5 * time.Minute + // monitorQuotaDegradedUsedPercent 任一用量窗口使用率超过该阈值时, + // 配额检查状态记为 degraded(对齐账号页展示阈值)。 + monitorQuotaDegradedUsedPercent = 90.0 // MonitorDefaultGrokModel 是新增 Grok 监控未显式指定模型时使用的轻量测活模型。 MonitorDefaultGrokModel = "grok-4.5" @@ -118,7 +145,16 @@ var ( "CHANNEL_MONITOR_NOT_FOUND", "channel monitor not found", ) ErrChannelMonitorInvalidProvider = infraerrors.BadRequest( - "CHANNEL_MONITOR_INVALID_PROVIDER", "provider must be one of openai/anthropic/gemini/grok", + "CHANNEL_MONITOR_INVALID_PROVIDER", "provider must be one of openai/anthropic/gemini/grok/antigravity/kimi/zhipu/deepseek", + ) + ErrChannelMonitorInvalidCheckMode = infraerrors.BadRequest( + "CHANNEL_MONITOR_INVALID_CHECK_MODE", "check_mode must be one of probe/quota/quota_probe; antigravity only supports quota", + ) + ErrChannelMonitorAccountRequired = infraerrors.BadRequest( + "CHANNEL_MONITOR_ACCOUNT_REQUIRED", "account_id is required for quota-based check_mode", + ) + ErrChannelMonitorProviderIncompatible = infraerrors.BadRequest( + "CHANNEL_MONITOR_PROVIDER_INCOMPATIBLE", "monitor provider must match the linked account platform", ) ErrChannelMonitorInvalidAPIMode = infraerrors.BadRequest( "CHANNEL_MONITOR_INVALID_API_MODE", "api_mode must be chat_completions or responses; responses is only supported for openai", diff --git a/backend/internal/service/channel_monitor_quota_fetcher.go b/backend/internal/service/channel_monitor_quota_fetcher.go new file mode 100644 index 0000000000..83d8e0f455 --- /dev/null +++ b/backend/internal/service/channel_monitor_quota_fetcher.go @@ -0,0 +1,429 @@ +package service + +import ( + "context" + "fmt" + "log/slog" + "sort" + "strconv" + "strings" + "sync" + "time" + + "github.com/Wei-Shaw/sub2api/internal/domain" + "github.com/Wei-Shaw/sub2api/internal/pkg/xai" +) + +// 渠道监控「配额模式」的配额抓取器。 +// +// 不直接对接上游,而是把账号侧现成的用量服务归一成 domain.MonitorQuotaSnapshot: +// - 海外 5 家(anthropic/openai/gemini/antigravity/grok)→ AccountUsageService.GetUsage +// - 国产 coding plan(kimi/zhipu/deepseek)→ CNProviderQuotaService.QueryUsage +// - 国产 payg(kimi/deepseek)→ CNProviderBalanceService.QueryBalance +// (zhipu payg 无公开余额端点,QueryBalance 会返回该错误,原样透出) +// +// Fetch 永不返回 error:所有失败都降级为 Success=false 的快照照常入库, +// 由 deriveQuotaCheckResult 推导为 failed/error 状态。 +// +// 多个监控可能关联同一账号,而 interval 最小 15s 且国产配额服务自身无缓存, +// 所以成功快照统一带 monitorQuotaFetchCacheTTL 缓存,防止打爆上游配额端点。 + +// monitorUsageSource 海外平台账号用量查询(AccountUsageService 天然满足)。 +type monitorUsageSource interface { + GetUsage(ctx context.Context, accountID int64, force ...bool) (*UsageInfo, error) +} + +// monitorCNQuotaSource 国产 coding plan 滚动窗口额度探测(CNProviderQuotaService 天然满足)。 +type monitorCNQuotaSource interface { + QueryUsage(ctx context.Context, accountID int64) (*CNProviderQuotaProbeResult, error) +} + +// monitorCNBalanceSource 国产 payg 余额探测(CNProviderBalanceService 天然满足)。 +type monitorCNBalanceSource interface { + QueryBalance(ctx context.Context, accountID int64) (*CNProviderBalanceResult, error) +} + +// monitorAccountSource 账号加载(AccountRepository 天然满足)。 +type monitorAccountSource interface { + GetByID(ctx context.Context, id int64) (*Account, error) +} + +// ChannelMonitorQuotaFetcher 配额抓取器(带成功快照 TTL 缓存)。 +type ChannelMonitorQuotaFetcher struct { + usage monitorUsageSource + cnQuota monitorCNQuotaSource + cnBalance monitorCNBalanceSource + accounts monitorAccountSource + + mu sync.Mutex + cache map[int64]monitorQuotaCacheEntry +} + +type monitorQuotaCacheEntry struct { + snapshot *domain.MonitorQuotaSnapshot + expiry time.Time +} + +// NewChannelMonitorQuotaFetcher 构造配额抓取器。 +func NewChannelMonitorQuotaFetcher( + usage monitorUsageSource, + cnQuota monitorCNQuotaSource, + cnBalance monitorCNBalanceSource, + accounts monitorAccountSource, +) *ChannelMonitorQuotaFetcher { + return &ChannelMonitorQuotaFetcher{ + usage: usage, + cnQuota: cnQuota, + cnBalance: cnBalance, + accounts: accounts, + cache: make(map[int64]monitorQuotaCacheEntry), + } +} + +// LoadAccount 加载账号(不走缓存)。供 Create/Update 时校验 +// provider 与 account.platform 一致;账号不存在时返回错误。 +func (f *ChannelMonitorQuotaFetcher) LoadAccount(ctx context.Context, id int64) (*Account, error) { + if f == nil || f.accounts == nil { + return nil, fmt.Errorf("quota fetcher is not configured") + } + return f.accounts.GetByID(ctx, id) +} + +// Fetch 抓取账号的最新配额快照。永不返回 error:失败降级为 +// Success=false 快照(Error 带摘要),保证检测历史的时间线连续。 +func (f *ChannelMonitorQuotaFetcher) Fetch(ctx context.Context, accountID int64) *domain.MonitorQuotaSnapshot { + now := time.Now() + + if cached, ok := f.cachedSnapshot(accountID, now); ok { + return cached + } + + snapshot := f.fetchUncached(ctx, accountID, now) + if snapshot.Success { + f.storeSnapshot(accountID, snapshot, now.Add(monitorQuotaFetchCacheTTL)) + } + return snapshot +} + +func (f *ChannelMonitorQuotaFetcher) cachedSnapshot(accountID int64, now time.Time) (*domain.MonitorQuotaSnapshot, bool) { + f.mu.Lock() + defer f.mu.Unlock() + entry, ok := f.cache[accountID] + if !ok || now.After(entry.expiry) { + return nil, false + } + return entry.snapshot, true +} + +func (f *ChannelMonitorQuotaFetcher) storeSnapshot(accountID int64, snapshot *domain.MonitorQuotaSnapshot, expiry time.Time) { + f.mu.Lock() + defer f.mu.Unlock() + f.cache[accountID] = monitorQuotaCacheEntry{snapshot: snapshot, expiry: expiry} +} + +func (f *ChannelMonitorQuotaFetcher) fetchUncached(ctx context.Context, accountID int64, now time.Time) *domain.MonitorQuotaSnapshot { + if f == nil { + return quotaErrorSnapshot("usage", "quota fetcher is not configured", now) + } + + account, err := f.LoadAccount(ctx, accountID) + if err != nil || account == nil { + // FK ON DELETE SET NULL 后 account_id 可能为空/失效;显式报「账号未关联」, + // 推导为 degraded(配置问题,不是渠道故障)。 + slog.Warn("channel_monitor: load linked account failed", + "account_id", accountID, "error", err) + return quotaErrorSnapshot("usage", "linked account not found", now) + } + + switch account.Platform { + case domain.PlatformKimi, domain.PlatformZhipu, domain.PlatformDeepseek: + if account.IsCodingPlan() { + return f.fetchCNQuota(ctx, accountID, now) + } + return f.fetchCNBalance(ctx, accountID, now) + default: + return f.fetchUsage(ctx, accountID, now) + } +} + +// fetchUsage 海外平台:AccountUsageService.GetUsage → 快照。 +func (f *ChannelMonitorQuotaFetcher) fetchUsage(ctx context.Context, accountID int64, now time.Time) *domain.MonitorQuotaSnapshot { + if f.usage == nil { + return quotaErrorSnapshot("usage", "usage service is not configured", now) + } + usage, err := f.usage.GetUsage(ctx, accountID) + if err != nil { + msg := truncateMessage(sanitizeErrorMessage(err.Error())) + return &domain.MonitorQuotaSnapshot{ + Source: "usage", + Success: false, + CredentialInvalid: isCredentialErrorMessage(msg), + Error: msg, + FetchedAt: now, + } + } + snapshot := &domain.MonitorQuotaSnapshot{ + Source: "usage", + Success: true, + PlanLevel: usage.SubscriptionTier, + Tiers: usageQuotaTiers(usage), + FetchedAt: now, + } + if snapshot.PlanLevel == "" { + snapshot.PlanLevel = usage.SubscriptionTierRaw + } + return snapshot +} + +// usageQuotaTiers 把 UsageInfo 的各平台窗口归一为 tier 列表(无数据的窗口跳过)。 +func usageQuotaTiers(usage *UsageInfo) []domain.MonitorQuotaTier { + if usage == nil { + return nil + } + tiers := make([]domain.MonitorQuotaTier, 0, 8) + appendProgressTier(&tiers, "5h", "", usage.FiveHour) + appendProgressTier(&tiers, "7d", "", usage.SevenDay) + appendProgressTier(&tiers, "7d-sonnet", "", usage.SevenDaySonnet) + appendProgressTier(&tiers, "7d-fable", "", usage.SevenDayFable) + appendProgressTier(&tiers, "30d", "", usage.ThirtyDay) + // Gemini 多档日配额:同 Window 不同 Label。 + appendProgressTier(&tiers, "daily", "shared", usage.GeminiSharedDaily) + appendProgressTier(&tiers, "daily", "pro", usage.GeminiProDaily) + appendProgressTier(&tiers, "daily", "flash", usage.GeminiFlashDaily) + // Grok requests/tokens 两个日窗口 + 月度计费窗口。 + appendQuotaWindowTier(&tiers, "daily", "requests", usage.GrokRequestQuota) + appendQuotaWindowTier(&tiers, "daily", "tokens", usage.GrokTokenQuota) + // Antigravity per-model 总量额度,Label = 模型名(按名排序保证输出稳定)。 + for _, model := range sortedQuotaModelNames(usage.AntigravityQuota) { + q := usage.AntigravityQuota[model] + if q == nil { + continue + } + tiers = append(tiers, domain.MonitorQuotaTier{ + Window: "total", + Label: model, + UsedPercent: float64(q.Utilization), + ResetAt: q.ResetTime, + }) + } + if len(tiers) == 0 { + return nil + } + return tiers +} + +func appendProgressTier(tiers *[]domain.MonitorQuotaTier, window, label string, p *UsageProgress) { + if p == nil { + return + } + tier := domain.MonitorQuotaTier{ + Window: window, + Label: label, + UsedPercent: p.Utilization, + } + if p.ResetsAt != nil { + tier.ResetAt = p.ResetsAt.UTC().Format(time.RFC3339) + } + if p.LimitRequests > 0 { + tier.Used = float64(p.UsedRequests) + tier.Limit = float64(p.LimitRequests) + } + *tiers = append(*tiers, tier) +} + +func appendQuotaWindowTier(tiers *[]domain.MonitorQuotaTier, window, label string, q *xai.QuotaWindow) { + if q == nil || q.Limit == nil || *q.Limit <= 0 { + return + } + used := *q.Limit + if q.Remaining != nil { + used = *q.Limit - *q.Remaining + if used < 0 { + used = 0 + } + } + tier := domain.MonitorQuotaTier{ + Window: window, + Label: label, + Used: float64(used), + Limit: float64(*q.Limit), + UsedPercent: float64(used) / float64(*q.Limit) * 100, + } + if q.ResetAt != "" { + tier.ResetAt = q.ResetAt + } else if q.ResetUnix != nil && *q.ResetUnix > 0 { + tier.ResetAt = time.Unix(*q.ResetUnix, 0).UTC().Format(time.RFC3339) + } + *tiers = append(*tiers, tier) +} + +func sortedQuotaModelNames(quotas map[string]*AntigravityModelQuota) []string { + names := make([]string, 0, len(quotas)) + for name := range quotas { + names = append(names, name) + } + sort.Strings(names) + return names +} + +// fetchCNQuota 国产 coding plan:CNProviderQuotaService.QueryUsage → 快照。 +func (f *ChannelMonitorQuotaFetcher) fetchCNQuota(ctx context.Context, accountID int64, now time.Time) *domain.MonitorQuotaSnapshot { + if f.cnQuota == nil { + return quotaErrorSnapshot("cn_quota", "cn quota service is not configured", now) + } + result, err := f.cnQuota.QueryUsage(ctx, accountID) + if err != nil { + msg := truncateMessage(sanitizeErrorMessage(err.Error())) + return &domain.MonitorQuotaSnapshot{ + Source: "cn_quota", + Success: false, + CredentialInvalid: isCredentialErrorMessage(msg), + Error: msg, + FetchedAt: now, + } + } + snapshot := &domain.MonitorQuotaSnapshot{ + Source: "cn_quota", + Success: result.Success, + PlanLevel: result.PlanLevel, + Error: result.Error, + FetchedAt: now, + } + if !result.Success && !result.CredentialValid { + snapshot.CredentialInvalid = true + } + if len(result.Tiers) > 0 { + snapshot.Tiers = make([]domain.MonitorQuotaTier, 0, len(result.Tiers)) + for _, t := range result.Tiers { + snapshot.Tiers = append(snapshot.Tiers, domain.MonitorQuotaTier{ + Window: t.Window, + UsedPercent: t.UsedPercent, + ResetAt: t.ResetAt, + }) + } + } + if !snapshot.Success { + snapshot.Error = firstNonEmpty(snapshot.Error, "cn quota probe failed") + } + return snapshot +} + +// fetchCNBalance 国产 payg:CNProviderBalanceService.QueryBalance → 快照。 +func (f *ChannelMonitorQuotaFetcher) fetchCNBalance(ctx context.Context, accountID int64, now time.Time) *domain.MonitorQuotaSnapshot { + if f.cnBalance == nil { + return quotaErrorSnapshot("cn_balance", "cn balance service is not configured", now) + } + result, err := f.cnBalance.QueryBalance(ctx, accountID) + if err != nil { + msg := truncateMessage(sanitizeErrorMessage(err.Error())) + return &domain.MonitorQuotaSnapshot{ + Source: "cn_balance", + Success: false, + CredentialInvalid: isCredentialErrorMessage(msg), + Error: msg, + FetchedAt: now, + } + } + snapshot := &domain.MonitorQuotaSnapshot{ + Source: "cn_balance", + Success: result.Success, + Currency: result.Currency, + Error: result.Error, + FetchedAt: now, + } + if result.Success { + balance := result.Balance + snapshot.Balance = &balance + } else if result.StatusCode == 401 || result.StatusCode == 403 { + snapshot.CredentialInvalid = true + } + if len(result.Balances) > 0 { + snapshot.Balances = make([]domain.MonitorBalance, 0, len(result.Balances)) + for _, b := range result.Balances { + snapshot.Balances = append(snapshot.Balances, domain.MonitorBalance{ + Currency: b.Currency, + Balance: b.Balance, + }) + } + } + if !snapshot.Success { + snapshot.Error = firstNonEmpty(snapshot.Error, "cn balance probe failed") + } + return snapshot +} + +// quotaErrorSnapshot 构造统一错误快照。 +func quotaErrorSnapshot(source, message string, now time.Time) *domain.MonitorQuotaSnapshot { + return &domain.MonitorQuotaSnapshot{ + Source: source, + Success: false, + Error: truncateMessage(sanitizeErrorMessage(message)), + FetchedAt: now, + } +} + +// isCredentialErrorMessage 上游 401/403 鉴权失败的启发式识别 +// (海外 GetUsage 的错误没有结构化状态码,只能看文本)。 +func isCredentialErrorMessage(msg string) bool { + msg = strings.ToLower(msg) + return strings.Contains(msg, "401") || + strings.Contains(msg, "403") || + strings.Contains(msg, "unauthorized") || + strings.Contains(msg, "forbidden") || + strings.Contains(msg, "invalid_api_key") || + strings.Contains(msg, "authentication") +} + +// deriveQuotaCheckResult 把配额快照推导为检测状态(复用既有 status 枚举, +// 时间线/可用率机制自动生效): +// - 查询成功且无告警 → operational +// - 任一窗口使用率 >= 阈值或余额耗尽 → degraded +// - 账号未关联(配置问题) → degraded +// - 凭据失效(401/403) → failed +// - 网络/解析等其他错误 → error +func deriveQuotaCheckResult(snapshot *domain.MonitorQuotaSnapshot, model string, checkedAt time.Time) *CheckResult { + res := &CheckResult{Model: model, CheckedAt: checkedAt} + if snapshot == nil { + res.Status = MonitorStatusError + res.Message = "quota snapshot missing" + return res + } + + switch { + case !snapshot.Success && snapshot.CredentialInvalid: + res.Status = MonitorStatusFailed + res.Message = snapshot.Error + case !snapshot.Success && strings.Contains(snapshot.Error, "linked account not found"): + res.Status = MonitorStatusDegraded + res.Message = snapshot.Error + case !snapshot.Success: + res.Status = MonitorStatusError + res.Message = snapshot.Error + default: + if hint := quotaDegradedHint(snapshot); hint != "" { + res.Status = MonitorStatusDegraded + res.Message = hint + } else { + res.Status = MonitorStatusOperational + } + } + return res +} + +// quotaDegradedHint 生成 degraded 的 message(指出触发告警的窗口/余额); +// 空串表示无告警。 +func quotaDegradedHint(snapshot *domain.MonitorQuotaSnapshot) string { + for _, tier := range snapshot.Tiers { + if tier.UsedPercent >= monitorQuotaDegradedUsedPercent { + name := tier.Window + if tier.Label != "" { + name = tier.Label + "/" + tier.Window + } + 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, "?")) + } + return "" +} diff --git a/backend/internal/service/channel_monitor_service.go b/backend/internal/service/channel_monitor_service.go index 248b72133d..9b0bb07639 100644 --- a/backend/internal/service/channel_monitor_service.go +++ b/backend/internal/service/channel_monitor_service.go @@ -11,6 +11,7 @@ import ( "sync" "time" + "github.com/Wei-Shaw/sub2api/internal/domain" "golang.org/x/sync/errgroup" ) @@ -77,6 +78,10 @@ type ChannelMonitorService struct { // scheduler 由 wire 通过 SetScheduler 注入;CRUD 后调用对应钩子即时同步任务。 // 测试或未注入场景下保持 nil,所有钩子调用变为 no-op。 scheduler MonitorScheduler + // quotaFetcher 由 wire 通过 SetQuotaFetcher 注入(accountUsage/CN 服务在本服务 + // 之后构造,构造参数注入会破坏既有依赖顺序)。nil 时 fail-closed: + // 配额模式的检测产出「未配置」错误快照,Create/Update 关联账号直接报错。 + quotaFetcher *ChannelMonitorQuotaFetcher } const maxChannelMonitorNameRunes = 100 @@ -150,6 +155,10 @@ func (s *ChannelMonitorService) Create(ctx context.Context, p ChannelMonitorCrea if err := validateExtraHeaders(p.ExtraHeaders); err != nil { return nil, err } + if err := s.validateLinkedAccount(ctx, p.Provider, p.AccountID); err != nil { + return nil, err + } + checkMode := defaultCheckMode(p.CheckMode) encrypted, err := s.encryptor.Encrypt(p.APIKey) if err != nil { return nil, fmt.Errorf("encrypt api key: %w", err) @@ -160,7 +169,7 @@ func (s *ChannelMonitorService) Create(ctx context.Context, p ChannelMonitorCrea APIMode: defaultAPIMode(p.APIMode), Endpoint: normalizeEndpoint(p.Endpoint), APIKey: encrypted, // 注意:传入 repository 时该字段为密文 - PrimaryModel: normalizeMonitorPrimaryModel(p.Provider, p.PrimaryModel), + PrimaryModel: normalizeMonitorPrimaryModel(p.Provider, checkMode, p.PrimaryModel), ExtraModels: normalizeModels(p.ExtraModels), GroupName: strings.TrimSpace(p.GroupName), Enabled: p.Enabled, @@ -171,6 +180,8 @@ func (s *ChannelMonitorService) Create(ctx context.Context, p ChannelMonitorCrea ExtraHeaders: emptyHeadersIfNil(p.ExtraHeaders), BodyOverrideMode: defaultBodyMode(p.BodyOverrideMode), BodyOverride: p.BodyOverride, + CheckMode: checkMode, + AccountID: cloneInt64Pointer(p.AccountID), } if err := s.repo.Create(ctx, m); err != nil { return nil, fmt.Errorf("create channel monitor: %w", err) @@ -236,6 +247,8 @@ func (s *ChannelMonitorService) Duplicate( ExtraHeaders: cloneChannelMonitorHeaders(source.ExtraHeaders), BodyOverrideMode: source.BodyOverrideMode, BodyOverride: bodyOverride, + CheckMode: defaultCheckMode(source.CheckMode), + AccountID: cloneInt64Pointer(source.AccountID), DuplicateOperationID: operationID, } if err := s.repo.Create(ctx, duplicate); err != nil { @@ -290,11 +303,22 @@ func (s *ChannelMonitorService) decryptAPIKeyForDuplicate(source *ChannelMonitor return "", ErrChannelMonitorAPIKeyDecryptFailed } plain, err := s.encryptor.Decrypt(source.APIKey) - if err != nil || strings.TrimSpace(plain) == "" { + if err != nil { slog.Warn("channel_monitor: decrypt api key for duplicate failed", "monitor_id", source.ID, "error", err) return "", ErrChannelMonitorAPIKeyDecryptFailed } + // quota 模式明文为空串是合法状态(api_key_encrypted 存的是加密空串): + // 重加密空串即可。若在此报错,克隆出的配额监控会被 runner 当作 + // 解密失败而 Unschedule,静默停摆。 + if strings.TrimSpace(plain) == "" { + if monitorCheckModeUsesQuota(defaultCheckMode(source.CheckMode)) { + return "", nil + } + slog.Warn("channel_monitor: decrypted api key for duplicate is empty", + "monitor_id", source.ID) + return "", ErrChannelMonitorAPIKeyDecryptFailed + } return plain, nil } @@ -343,10 +367,16 @@ func cloneChannelMonitorJSONMap(source map[string]any) (map[string]any, error) { } // validateCreateParams 把 Create 入参的所有校验聚拢为一个函数,避免 Create 主体超过 30 行。 +// 按 check_mode 分支:probe 沿用 endpoint+api_key 必填;quota 只需关联账号; +// quota_probe 两者皆需。 func validateCreateParams(p ChannelMonitorCreateParams) error { if err := validateProvider(p.Provider); err != nil { return err } + checkMode := defaultCheckMode(p.CheckMode) + if err := validateCheckMode(p.Provider, checkMode); err != nil { + return err + } if err := validateAPIMode(p.Provider, p.APIMode); err != nil { return err } @@ -356,18 +386,45 @@ func validateCreateParams(p ChannelMonitorCreateParams) error { if err := validateJitter(p.JitterSeconds, p.IntervalSeconds); err != nil { return err } - if err := validateEndpoint(p.Endpoint); err != nil { - return err + usesQuota := monitorCheckModeUsesQuota(checkMode) + // probe 分支(含 quota_probe 的探活部分)仍需 endpoint + api_key; + // quota 模式 endpoint/api_key 留空,避免要求用户填无意义的占位值。 + if checkMode != MonitorCheckModeQuota { + if err := validateEndpoint(p.Endpoint); err != nil { + return err + } + if strings.TrimSpace(p.APIKey) == "" { + return ErrChannelMonitorMissingAPIKey + } } - if strings.TrimSpace(p.APIKey) == "" { - return ErrChannelMonitorMissingAPIKey + if usesQuota && (p.AccountID == nil || *p.AccountID <= 0) { + return ErrChannelMonitorAccountRequired } - if normalizeMonitorPrimaryModel(p.Provider, p.PrimaryModel) == "" { + if normalizeMonitorPrimaryModel(p.Provider, checkMode, p.PrimaryModel) == "" { return ErrChannelMonitorMissingPrimaryModel } return nil } +// validateLinkedAccount 校验关联账号存在且平台与监控 provider 一致。 +// fetcher 未注入时 fail-closed(拒绝创建配额监控,而不是创建后静默坏)。 +func (s *ChannelMonitorService) validateLinkedAccount(ctx context.Context, provider string, accountID *int64) error { + if accountID == nil || *accountID <= 0 { + return nil + } + if s.quotaFetcher == nil { + return ErrChannelMonitorAccountRequired + } + account, err := s.quotaFetcher.LoadAccount(ctx, *accountID) + if err != nil || account == nil { + return ErrChannelMonitorAccountRequired + } + if account.Platform != provider { + return ErrChannelMonitorProviderIncompatible + } + return nil +} + // Update 更新监控。APIKey 字段:nil 或空字符串 = 不修改;非空 = 加密后覆盖。 func (s *ChannelMonitorService) Update(ctx context.Context, id int64, p ChannelMonitorUpdateParams) (*ChannelMonitor, error) { existing, err := s.repo.GetByID(ctx, id) @@ -382,6 +439,14 @@ func (s *ChannelMonitorService) Update(ctx context.Context, id int64, p ChannelM if err != nil { return nil, err } + if err := s.validateProbeAPIKey(existing, newPlainAPIKey); err != nil { + return nil, err + } + if p.Provider != nil || p.CheckMode != nil || p.AccountID != nil { + if err := s.revalidateLinkedAccount(ctx, existing); err != nil { + return nil, err + } + } if err := s.repo.Update(ctx, existing); err != nil { return nil, fmt.Errorf("update channel monitor: %w", err) @@ -401,6 +466,76 @@ func (s *ChannelMonitorService) Update(ctx context.Context, id int64, p ChannelM return existing, nil } +// validateMonitorModeFields 校验 check_mode 与其它字段的组合约束 +// (在 provider/check_mode/account_id/endpoint 全部应用后调用): +// - quota / quota_probe 必须关联账号 +// - probe / quota_probe 必须持有 endpoint(探活目标) +func validateMonitorModeFields(m *ChannelMonitor) error { + checkMode := defaultCheckMode(m.CheckMode) + if monitorCheckModeUsesQuota(checkMode) && m.AccountID == nil { + return ErrChannelMonitorAccountRequired + } + if checkMode != MonitorCheckModeQuota && strings.TrimSpace(m.Endpoint) == "" { + return ErrChannelMonitorInvalidEndpoint + } + return nil +} + +// validateProbeAPIKey 探活模式(probe / quota_probe)必须持有可用明文 key: +// 存量密文解密为空串(quota 监控切回探活但未重填 key)时拒绝。 +// 密文损坏的情况交给既有 APIKeyDecryptFailed 链路(Get/RunCheck 会显式报错)。 +func (s *ChannelMonitorService) validateProbeAPIKey(m *ChannelMonitor, newPlainKey string) error { + if defaultCheckMode(m.CheckMode) == MonitorCheckModeQuota { + return nil + } + if strings.TrimSpace(newPlainKey) != "" { + return nil + } + if strings.TrimSpace(m.APIKey) == "" { + return ErrChannelMonitorMissingAPIKey + } + plain, err := s.encryptor.Decrypt(m.APIKey) + if err != nil { + return nil + } + if strings.TrimSpace(plain) == "" { + return ErrChannelMonitorMissingAPIKey + } + return nil +} + +// revalidateLinkedAccount 在 provider/check_mode/account_id 任一变化后复核关联账号: +// - 账号已被删除或平台失配:probe 模式自动解绑(静默修复), +// quota 模式显式报错(配额监控必须有可用数据源) +func (s *ChannelMonitorService) revalidateLinkedAccount(ctx context.Context, m *ChannelMonitor) error { + usesQuota := monitorCheckModeUsesQuota(defaultCheckMode(m.CheckMode)) + if m.AccountID == nil { + if usesQuota { + return ErrChannelMonitorAccountRequired + } + return nil + } + if s.quotaFetcher == nil { + return ErrChannelMonitorAccountRequired + } + account, err := s.quotaFetcher.LoadAccount(ctx, *m.AccountID) + if err != nil || account == nil { + if usesQuota { + return ErrChannelMonitorAccountRequired + } + m.AccountID = nil + return nil + } + if account.Platform != m.Provider { + if usesQuota { + return ErrChannelMonitorProviderIncompatible + } + m.AccountID = nil + return nil + } + return nil +} + // applyAPIKeyUpdate 处理 Update 中的 APIKey 字段: // - 入参 raw 为 nil 或空白:不修改 existing.APIKey(仍为密文),返回 updated=false // - 非空:加密后写入 existing.APIKey;同时把明文返回给调用方, @@ -454,6 +589,9 @@ func (s *ChannelMonitorService) ListHistory(ctx context.Context, id int64, model // 写历史记录并更新 last_checked_at。返回每个模型的检测结果。 // 仅当 channel_monitor_enabled=true 且 channel_monitor_mode=v1 时真正探测; // mode=v2 时返回 ErrChannelMonitorActiveProbesRetired,不产生上游流量。 +// +// 按 check_mode 分派:probe(默认,现状探活)/ quota(仅查关联账号配额, +// 零 LLM 成本)/ quota_probe(探活 + 配额快照挂主模型行)。 func (s *ChannelMonitorService) RunCheck(ctx context.Context, id int64) ([]*CheckResult, error) { rt := s.probeRuntime(ctx) if !rt.Enabled { @@ -466,14 +604,59 @@ func (s *ChannelMonitorService) RunCheck(ctx context.Context, id int64) ([]*Chec if err != nil { return nil, err } - if m.APIKeyDecryptFailed { + checkMode := defaultCheckMode(m.CheckMode) + if checkMode != MonitorCheckModeQuota && m.APIKeyDecryptFailed { return nil, ErrChannelMonitorAPIKeyDecryptFailed } - results := s.runChecksConcurrent(ctx, m) + + var results []*CheckResult + switch checkMode { + case MonitorCheckModeQuota: + results = s.runQuotaOnlyCheck(ctx, m) + case MonitorCheckModeQuotaProbe: + results = s.runChecksConcurrent(ctx, m) + attachQuotaSnapshot(results, s.fetchQuotaSnapshot(ctx, m)) + default: + results = s.runChecksConcurrent(ctx, m) + } s.persistCheckResults(ctx, m, results) return results, nil } +// runQuotaOnlyCheck quota 模式:一次配额抓取 → 单条 CheckResult +// (Model=PrimaryModel,默认 "quota";无 ping/latency,状态由快照推导)。 +func (s *ChannelMonitorService) runQuotaOnlyCheck(ctx context.Context, m *ChannelMonitor) []*CheckResult { + snapshot := s.fetchQuotaSnapshot(ctx, m) + res := deriveQuotaCheckResult(snapshot, m.PrimaryModel, time.Now()) + res.Quota = snapshot + return []*CheckResult{res} +} + +// fetchQuotaSnapshot 抓取关联账号配额。未关联账号 / fetcher 未注入时返回 +// 显式错误快照(不返回 error,保证检测周期与历史时间线连续)。 +func (s *ChannelMonitorService) fetchQuotaSnapshot(ctx context.Context, m *ChannelMonitor) *domain.MonitorQuotaSnapshot { + if m.AccountID == nil { + return quotaErrorSnapshot("usage", "linked account not found", time.Now()) + } + if s.quotaFetcher == nil { + return quotaErrorSnapshot("usage", "quota fetcher is not configured", time.Now()) + } + return s.quotaFetcher.Fetch(ctx, *m.AccountID) +} + +// attachQuotaSnapshot quota_probe:把配额快照挂到主模型行(results[0])。 +// 配额失败不改变探活状态,仅在探活 message 为空时附注失败原因。 +func attachQuotaSnapshot(results []*CheckResult, snapshot *domain.MonitorQuotaSnapshot) { + if len(results) == 0 || snapshot == nil { + return + } + primary := results[0] + primary.Quota = snapshot + if !snapshot.Success && strings.TrimSpace(primary.Message) == "" { + primary.Message = truncateMessage("quota fetch failed: " + snapshot.Error) + } +} + // persistCheckResults 写入本次检测的历史记录并更新 last_checked_at。 // 任一写库失败都只记日志,不影响调用方拿到 results(与 MVP 期望一致:宁可漏记历史也要先返回结果)。 func (s *ChannelMonitorService) persistCheckResults(ctx context.Context, m *ChannelMonitor, results []*CheckResult) { @@ -487,6 +670,7 @@ func (s *ChannelMonitorService) persistCheckResults(ctx context.Context, m *Chan PingLatencyMs: r.PingLatencyMs, Message: r.Message, CheckedAt: r.CheckedAt, + Quota: r.Quota, }) } if err := s.repo.InsertHistoryBatch(ctx, rows); err != nil { @@ -541,6 +725,14 @@ func (s *ChannelMonitorService) SetScheduler(sched MonitorScheduler) { s.scheduler = sched } +// SetQuotaFetcher 由 wire 注入配额抓取器(账号侧用量服务聚合)。 +func (s *ChannelMonitorService) SetQuotaFetcher(fetcher *ChannelMonitorQuotaFetcher) { + if s == nil { + return + } + s.quotaFetcher = fetcher +} + // ListEnabledMonitors 返回所有 enabled=true 的监控(解密后),供 runner 启动时建立任务表。 func (s *ChannelMonitorService) ListEnabledMonitors(ctx context.Context) ([]*ChannelMonitor, error) { all, err := s.repo.ListEnabled(ctx) @@ -693,14 +885,36 @@ func applyMonitorUpdate(existing *ChannelMonitor, p ChannelMonitorUpdateParams) providerChanged = existing.Provider != *p.Provider existing.Provider = *p.Provider } - if p.Endpoint != nil { - if err := validateEndpoint(*p.Endpoint); err != nil { + if p.CheckMode != nil { + mode := defaultCheckMode(*p.CheckMode) + if err := validateCheckMode(existing.Provider, mode); err != nil { return err } + existing.CheckMode = mode + } + if p.AccountID != nil { + if *p.AccountID > 0 { + id := *p.AccountID + existing.AccountID = &id + } else { + existing.AccountID = nil // 0 = 清空关联 + } + } + if p.Endpoint != nil { + // quota 模式允许清空 endpoint(校验由 validateMonitorModeFields 兜底)。 + if strings.TrimSpace(*p.Endpoint) != "" { + if err := validateEndpoint(*p.Endpoint); err != nil { + return err + } + } existing.Endpoint = normalizeEndpoint(*p.Endpoint) } + // 模式与字段的组合校验(provider/check_mode/account_id/endpoint 全部应用后)。 + if err := validateMonitorModeFields(existing); err != nil { + return err + } if p.PrimaryModel != nil { - primaryModel := normalizeMonitorPrimaryModel(existing.Provider, *p.PrimaryModel) + primaryModel := normalizeMonitorPrimaryModel(existing.Provider, defaultCheckMode(existing.CheckMode), *p.PrimaryModel) if primaryModel == "" { return ErrChannelMonitorMissingPrimaryModel } diff --git a/backend/internal/service/channel_monitor_types.go b/backend/internal/service/channel_monitor_types.go index 20ab1da935..298e037ee1 100644 --- a/backend/internal/service/channel_monitor_types.go +++ b/backend/internal/service/channel_monitor_types.go @@ -1,6 +1,10 @@ package service -import "time" +import ( + "time" + + "github.com/Wei-Shaw/sub2api/internal/domain" +) // MonitorBodyOverrideMode 自定义请求体处理模式。 // @@ -45,6 +49,11 @@ type ChannelMonitor struct { CreatedAt time.Time UpdatedAt time.Time + // 配额模式(check_mode = quota / quota_probe): + // 关联已有账号复用账号侧用量服务,Endpoint/APIKey 可为空(quota 模式)。 + CheckMode string // probe(默认)/ quota / quota_probe;空串按 probe 处理 + AccountID *int64 // 关联账号 ID;账号删除后被 DB 置空(监控保留并报「账号未关联」) + // 请求自定义快照(来自模板拷贝 or 用户手填,运行时直接读取) TemplateID *int64 // 仅用于 UI 分组 + 一键应用,运行时不用 ExtraHeaders map[string]string // 与 adapter 默认 headers 合并,用户优先 @@ -89,6 +98,10 @@ type ChannelMonitorCreateParams struct { ExtraHeaders map[string]string BodyOverrideMode string BodyOverride map[string]any + + // 配额模式:CheckMode 空串默认 probe;quota/quota_probe 必须关联账号。 + CheckMode string + AccountID *int64 } // ChannelMonitorUpdateParams 更新参数(指针字段表示"未提供则不更新")。 @@ -112,6 +125,11 @@ type ChannelMonitorUpdateParams struct { ExtraHeaders *map[string]string BodyOverrideMode *string BodyOverride *map[string]any + + // 配额模式:CheckMode nil = 不更新;AccountID nil = 不更新, + // 指向 0 = 清空关联(退回 probe 模式时由 CheckMode 分支兜底)。 + CheckMode *string + AccountID *int64 } // CheckResult 单个模型一次检测的结果。 @@ -122,6 +140,8 @@ type CheckResult struct { PingLatencyMs *int Message string CheckedAt time.Time + // Quota 配额模式附带快照(quota 模式唯一数据;quota_probe 挂在主模型行)。 + Quota *domain.MonitorQuotaSnapshot } // UserMonitorView 用户只读视图:监控概览(含主模型最近状态 + 7d 可用率 + 附加模型最近状态)。 @@ -137,6 +157,9 @@ type UserMonitorView struct { Availability7d float64 // 0-100 ExtraModels []ExtraModelStatus Timeline []UserMonitorTimelinePoint // 主模型最近 N 个历史点(按 checked_at DESC,最新在前) + // LatestQuota 主模型最近一次配额快照;channel_monitor_show_quota=false + // 时由 handler 服务端剥离。 + LatestQuota *domain.MonitorQuotaSnapshot } // UserMonitorTimelinePoint 用户视图 timeline 单点数据(去除 message 以减小响应体)。 @@ -183,6 +206,7 @@ type ChannelMonitorHistoryRow struct { PingLatencyMs *int Message string CheckedAt time.Time + Quota *domain.MonitorQuotaSnapshot } // ChannelMonitorHistoryEntry 历史记录查询返回行(含 ent 主键 ID)。 @@ -194,6 +218,7 @@ type ChannelMonitorHistoryEntry struct { PingLatencyMs *int Message string CheckedAt time.Time + Quota *domain.MonitorQuotaSnapshot } // ChannelMonitorLatest 最近一次检测的简明信息(用于 UserMonitorView 聚合)。 @@ -203,6 +228,7 @@ type ChannelMonitorLatest struct { LatencyMs *int PingLatencyMs *int CheckedAt time.Time + Quota *domain.MonitorQuotaSnapshot } // ChannelMonitorAvailability 单个模型在某窗口内的可用率与平均延迟(用于 UserMonitorDetail 聚合)。 @@ -223,4 +249,5 @@ type MonitorStatusSummary struct { PrimaryLatencyMs *int Availability7d float64 // 0-100,无历史时为 0 ExtraModels []ExtraModelStatus + LatestQuota *domain.MonitorQuotaSnapshot // 主模型最近配额快照(配额模式) } diff --git a/backend/internal/service/channel_monitor_validate.go b/backend/internal/service/channel_monitor_validate.go index 7740dc83b7..faf7247fbc 100644 --- a/backend/internal/service/channel_monitor_validate.go +++ b/backend/internal/service/channel_monitor_validate.go @@ -9,15 +9,81 @@ import ( // 渠道监控参数校验与归一化辅助函数。 // 校验失败一律返回 channel_monitor_const.go 中预定义的 Err* 错误,错误信息不含具体 IP/hostname,避免泄露内网拓扑。 +// monitorProviders 渠道监控支持的全部 provider(与迁移 226 的 CHECK 约束一致)。 +// 不再以 adapter 表为唯一来源:antigravity 没有探活 adapter,但支持配额模式。 +// +//nolint:gochecknoglobals // 静态查表,初始化后不变。 +var monitorProviders = map[string]struct{}{ + MonitorProviderOpenAI: {}, + MonitorProviderAnthropic: {}, + MonitorProviderGemini: {}, + MonitorProviderGrok: {}, + MonitorProviderAntigravity: {}, + MonitorProviderKimi: {}, + MonitorProviderZhipu: {}, + MonitorProviderDeepseek: {}, +} + +// probeCapableProviders 支持探活(probe / quota_probe)的 provider。 +// antigravity 上游无 Chat/Responses 可打(仅 IDE 代理形态),只允许配额模式。 +// +//nolint:gochecknoglobals // 静态查表,初始化后不变。 +var probeCapableProviders = map[string]struct{}{ + MonitorProviderOpenAI: {}, + MonitorProviderAnthropic: {}, + MonitorProviderGemini: {}, + MonitorProviderGrok: {}, + MonitorProviderKimi: {}, + MonitorProviderZhipu: {}, + MonitorProviderDeepseek: {}, +} + // validateProvider 校验 provider 字符串。 -// 唯一来源于 providerAdapters:新增 provider 只需要在 channel_monitor_checker.go 注册 adapter。 func validateProvider(p string) error { - if !isSupportedProvider(p) { + if _, ok := monitorProviders[p]; !ok { return ErrChannelMonitorInvalidProvider } return nil } +// providerSupportsProbe 该 provider 是否注册了探活 adapter(antigravity 为 false)。 +func providerSupportsProbe(p string) bool { + _, ok := probeCapableProviders[p] + return ok +} + +// defaultCheckMode 空串归一为 probe,保证存量数据与旧客户端兼容。 +func defaultCheckMode(checkMode string) string { + if strings.TrimSpace(checkMode) == "" { + return MonitorCheckModeProbe + } + return strings.TrimSpace(checkMode) +} + +// monitorCheckModeUsesQuota 该模式是否需要关联账号查配额。 +func monitorCheckModeUsesQuota(checkMode string) bool { + return checkMode == MonitorCheckModeQuota || checkMode == MonitorCheckModeQuotaProbe +} + +// validateCheckMode 校验 check_mode 与 provider 的组合矩阵: +// +// provider | probe | quota | quota_probe +// ------------------------+-------+-------+------------ +// openai/anthropic/... | Y | Y | Y +// antigravity(无 adapter)| N | Y | N +func validateCheckMode(provider, checkMode string) error { + checkMode = defaultCheckMode(checkMode) + switch checkMode { + case MonitorCheckModeProbe, MonitorCheckModeQuota, MonitorCheckModeQuotaProbe: + default: + return ErrChannelMonitorInvalidCheckMode + } + if checkMode != MonitorCheckModeQuota && !providerSupportsProbe(provider) { + return ErrChannelMonitorInvalidCheckMode + } + return nil +} + // validateAPIMode 校验 provider 与 api_mode 的组合。 // responses 只对 OpenAI 有意义;其它 provider 使用 chat_completions 作为默认占位。 func validateAPIMode(provider, apiMode string) error { @@ -124,13 +190,18 @@ func normalizeModels(in []string) []string { return out } -// normalizeMonitorPrimaryModel applies the Grok health-check default while -// preserving the existing required-model behavior for every other provider. -func normalizeMonitorPrimaryModel(provider, model string) string { +// normalizeMonitorPrimaryModel applies provider/check_mode defaults while +// preserving the existing required-model behavior: +// - Grok 探活默认轻量测活模型 +// - quota 模式占位 "quota"(primary_model 列 NotEmpty;历史行/时间线机制无需特判) +func normalizeMonitorPrimaryModel(provider, checkMode, model string) string { model = strings.TrimSpace(model) if model == "" && provider == MonitorProviderGrok { return MonitorDefaultGrokModel } + if model == "" && monitorCheckModeUsesQuota(defaultCheckMode(checkMode)) { + return MonitorDefaultQuotaModel + } return model }