feat(channel-monitor): quota mode service layer (fetcher + dispatch + repo)

- 常量:MonitorProvider{Antigravity,Kimi,Zhipu,Deepseek}、MonitorCheckMode 三态、
  配额缓存/告警阈值参数;新增 check_mode 相关业务错误
- 校验:monitorProviders(8)/probeCapableProviders(7) 集合替代 adapter 表判定,
  validateCheckMode 矩阵(antigravity 仅 quota);quota 模式 primary_model
  默认占位 "quota",endpoint/api_key 条件必填
- checker:kimi/deepseek 复用 /v1/chat/completions、zhipu 走
  /api/paas/v4/chat/completions;merge 黑名单与 replace 校验同步扩容
- ChannelMonitorQuotaFetcher:窄接口聚合账号侧三个用量服务 + 账号加载,
  归一为 domain.MonitorQuotaSnapshot;成功快照 5min TTL 缓存防配额端点风暴;
  Fetch 永不返回 error(失败降级 Success=false 快照);状态推导
  operational/degraded(≥90% 或余额耗尽/账号未关联)/failed(401/403)/error
- RunCheck 按 check_mode 分派:quota 单行结果、quota_probe 探活+配额挂主行、
  probe 不变;Update 组合校验 + 关联账号复核(probe 自动解绑失效账号);
  Duplicate 拷贝新字段且 quota 模式空明文重加密(防 runner 误停摆)
- repo:Create/Update/entToServiceMonitor 透传 check_mode/account_id;
  InsertHistoryBatch 写 quota;ListLatestForMonitorIDs 裸 SQL 加 quota 列 +
  scanMonitorQuota JSONB 解包(NULL 安全);ListHistory 透传 quota
This commit is contained in:
Randark
2026-08-18 04:03:13 +00:00
parent 615e6901ea
commit c44711ac98
9 changed files with 900 additions and 36 deletions
@@ -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 单币种余额条目。
@@ -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, &quota); 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{}
@@ -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
}
@@ -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) {
@@ -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",
@@ -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 ""
}
@@ -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
}
@@ -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 // 主模型最近配额快照(配额模式)
}
@@ -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
}