From 41344c20ff54e4a828a09fa0ca77d0167a66b7f4 Mon Sep 17 00:00:00 2001 From: Randark Date: Tue, 18 Aug 2026 04:14:15 +0000 Subject: [PATCH] feat(monitor): wire quota fetcher & expose check_mode in handlers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - handler DTO: create/update 接收 check_mode/account_id,provider oneof 扩至 8 家, endpoint/api_key 改为 omitempty(条件必填下沉 service 校验); monitor/checkResult/historyItem 响应透传 check_mode/account_id/quota - 用户端 latest_quota 由 channel_monitor_show_quota 控制,关闭时服务端剥离 - wire: NewChannelMonitorQuotaFetcher 以具体服务类型收参(窄接口包内保留供 stub),ProvideChannelMonitorRunner 注入后 SetQuotaFetcher --- backend/cmd/server/wire_gen.go | 3 +- .../handler/admin/channel_monitor_handler.go | 65 ++++++++++++++----- .../handler/channel_monitor_user_handler.go | 24 ++++++- .../service/channel_monitor_quota_fetcher.go | 28 +++++--- backend/internal/service/wire.go | 10 ++- 5 files changed, 97 insertions(+), 33 deletions(-) diff --git a/backend/cmd/server/wire_gen.go b/backend/cmd/server/wire_gen.go index 2b0caf982c..e68a6bd7d5 100644 --- a/backend/cmd/server/wire_gen.go +++ b/backend/cmd/server/wire_gen.go @@ -337,7 +337,8 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) { batchImageWorkerRuntime := service.ProvideBatchImageWorkerRuntime(batchImageRepository, accountRepository, batchImageQueue, usageBillingRepository, usageLogRepository, batchImageModelPricingResolver, apiKeyAuthCacheInvalidator, configConfig) scheduledTestRunnerService := service.ProvideScheduledTestRunnerService(scheduledTestPlanRepository, scheduledTestService, accountTestService, rateLimitService, configConfig) paymentOrderExpiryService := service.ProvidePaymentOrderExpiryService(paymentService, leaderLockCache, db) - channelMonitorRunner := service.ProvideChannelMonitorRunner(channelMonitorService, settingService) + channelMonitorQuotaFetcher := service.NewChannelMonitorQuotaFetcher(accountUsageService, cnProviderQuotaService, cnProviderBalanceService, accountRepository) + channelMonitorRunner := service.ProvideChannelMonitorRunner(channelMonitorService, settingService, channelMonitorQuotaFetcher) channelMonitorV2Aggregator := service.ProvideChannelMonitorV2Aggregator(channelMonitorV2Repository, db, settingService) userPlatformQuotaUsageFlusher := service.ProvideUserPlatformQuotaUsageFlusher(configConfig, billingCache, serviceUserPlatformQuotaRepository, timingWheelService) v := provideCleanup(client, redisClient, opsMetricsCollector, opsAggregationService, opsAlertEvaluatorService, opsCleanupService, opsScheduledReportService, opsSystemLogSink, opsService, opsIngressRejectAggregator, apiKeyService, authCacheInvalidationWorker, schedulerSnapshotService, tokenRefreshService, accountExpiryService, cnProviderBalanceCheckService, openAICodexVersionSyncService, proxyExpiryService, subscriptionExpiryService, usageCleanupService, idempotencyCleanupService, batchImageCleanupService, batchImageWorkerRuntime, pricingService, emailQueueService, billingCacheService, usageRecordWorkerPool, subscriptionService, oAuthService, openAIOAuthService, geminiOAuthService, antigravityOAuthService, grokOAuthService, openAIGatewayService, scheduledTestRunnerService, backupService, paymentOrderExpiryService, channelMonitorRunner, channelMonitorV2Aggregator, userPlatformQuotaUsageFlusher, upstreamBillingProbeService, ollamaCloudUsageService, auditLogService, promptService) diff --git a/backend/internal/handler/admin/channel_monitor_handler.go b/backend/internal/handler/admin/channel_monitor_handler.go index ec5735f0ca..4832be2dec 100644 --- a/backend/internal/handler/admin/channel_monitor_handler.go +++ b/backend/internal/handler/admin/channel_monitor_handler.go @@ -7,6 +7,7 @@ import ( "strings" "time" + "github.com/Wei-Shaw/sub2api/internal/domain" "github.com/Wei-Shaw/sub2api/internal/handler/dto" infraerrors "github.com/Wei-Shaw/sub2api/internal/pkg/errors" "github.com/Wei-Shaw/sub2api/internal/pkg/response" @@ -39,10 +40,10 @@ func NewChannelMonitorHandler(monitorService *service.ChannelMonitorService) *Ch type channelMonitorCreateRequest struct { Name string `json:"name" binding:"required,max=100"` - Provider string `json:"provider" binding:"required,oneof=openai anthropic gemini grok"` + Provider string `json:"provider" binding:"required,oneof=openai anthropic gemini grok antigravity kimi zhipu deepseek"` APIMode string `json:"api_mode" binding:"omitempty,oneof=chat_completions responses"` - Endpoint string `json:"endpoint" binding:"required,max=500"` - APIKey string `json:"api_key" binding:"required,max=2000"` + Endpoint string `json:"endpoint" binding:"omitempty,max=500"` + APIKey string `json:"api_key" binding:"omitempty,max=2000"` PrimaryModel string `json:"primary_model" binding:"max=200"` ExtraModels []string `json:"extra_models"` GroupName string `json:"group_name" binding:"max=100"` @@ -53,11 +54,17 @@ type channelMonitorCreateRequest struct { ExtraHeaders map[string]string `json:"extra_headers"` BodyOverrideMode string `json:"body_override_mode" binding:"omitempty,oneof=off merge replace"` BodyOverride map[string]any `json:"body_override"` + + // CheckMode: probe(默认)/ quota / quota_probe。quota 模式 endpoint/api_key + // 可空(条件必填校验在 service 层按模式分支)。 + CheckMode string `json:"check_mode" binding:"omitempty,oneof=probe quota quota_probe"` + // AccountID: 配额模式关联的账号 ID。 + AccountID *int64 `json:"account_id"` } type channelMonitorUpdateRequest struct { Name *string `json:"name" binding:"omitempty,max=100"` - Provider *string `json:"provider" binding:"omitempty,oneof=openai anthropic gemini grok"` + Provider *string `json:"provider" binding:"omitempty,oneof=openai anthropic gemini grok antigravity kimi zhipu deepseek"` APIMode *string `json:"api_mode" binding:"omitempty,oneof=chat_completions responses"` Endpoint *string `json:"endpoint" binding:"omitempty,max=500"` APIKey *string `json:"api_key" binding:"omitempty,max=2000"` @@ -72,6 +79,10 @@ type channelMonitorUpdateRequest struct { ExtraHeaders *map[string]string `json:"extra_headers"` BodyOverrideMode *string `json:"body_override_mode" binding:"omitempty,oneof=off merge replace"` BodyOverride *map[string]any `json:"body_override"` + + // CheckMode/AccountID:nil = 不更新;AccountID 指向 0 = 清空关联。 + CheckMode *string `json:"check_mode" binding:"omitempty,oneof=probe quota quota_probe"` + AccountID *int64 `json:"account_id"` } type channelMonitorResponse struct { @@ -101,25 +112,33 @@ type channelMonitorResponse struct { ExtraHeaders map[string]string `json:"extra_headers"` BodyOverrideMode string `json:"body_override_mode"` BodyOverride map[string]any `json:"body_override"` + + // 配额模式:check_mode + 关联账号 + 主模型最近配额快照 + // (LatestQuota 由 List handler 批量聚合后填充;管理端不受 channel_monitor_show_quota 影响)。 + CheckMode string `json:"check_mode"` + AccountID *int64 `json:"account_id"` + LatestQuota *domain.MonitorQuotaSnapshot `json:"latest_quota,omitempty"` } type channelMonitorCheckResultResponse struct { - Model string `json:"model"` - Status string `json:"status"` - LatencyMs *int `json:"latency_ms"` - PingLatencyMs *int `json:"ping_latency_ms"` - Message string `json:"message"` - CheckedAt string `json:"checked_at"` + Model string `json:"model"` + Status string `json:"status"` + LatencyMs *int `json:"latency_ms"` + PingLatencyMs *int `json:"ping_latency_ms"` + Message string `json:"message"` + CheckedAt string `json:"checked_at"` + Quota *domain.MonitorQuotaSnapshot `json:"quota,omitempty"` } type channelMonitorHistoryItemResponse struct { - ID int64 `json:"id"` - Model string `json:"model"` - Status string `json:"status"` - LatencyMs *int `json:"latency_ms"` - PingLatencyMs *int `json:"ping_latency_ms"` - Message string `json:"message"` - CheckedAt string `json:"checked_at"` + ID int64 `json:"id"` + Model string `json:"model"` + Status string `json:"status"` + LatencyMs *int `json:"latency_ms"` + PingLatencyMs *int `json:"ping_latency_ms"` + Message string `json:"message"` + CheckedAt string `json:"checked_at"` + Quota *domain.MonitorQuotaSnapshot `json:"quota,omitempty"` } // maskAPIKey 对 API Key 明文做脱敏:前 4 字符 + "***",长度 ≤ 4 时只显示 "***"。 @@ -163,7 +182,10 @@ func channelMonitorToResponse(m *service.ChannelMonitor) *channelMonitorResponse ExtraHeaders: headers, BodyOverrideMode: m.BodyOverrideMode, BodyOverride: m.BodyOverride, - // PrimaryStatus / PrimaryLatencyMs / Availability7d 由 List handler 在批量聚合后填充。 + CheckMode: m.CheckMode, + AccountID: m.AccountID, + // PrimaryStatus / PrimaryLatencyMs / Availability7d / LatestQuota + // 由 List handler 在批量聚合后填充。 } if m.LastCheckedAt != nil { s := m.LastCheckedAt.UTC().Format(time.RFC3339) @@ -180,6 +202,7 @@ func checkResultToResponse(r *service.CheckResult) channelMonitorCheckResultResp PingLatencyMs: r.PingLatencyMs, Message: r.Message, CheckedAt: r.CheckedAt.UTC().Format(time.RFC3339), + Quota: r.Quota, } } @@ -192,6 +215,7 @@ func historyEntryToResponse(e *service.ChannelMonitorHistoryEntry) channelMonito PingLatencyMs: e.PingLatencyMs, Message: e.Message, CheckedAt: e.CheckedAt.UTC().Format(time.RFC3339), + Quota: e.Quota, } } @@ -270,6 +294,7 @@ func buildListItemResponse(m *service.ChannelMonitor, summary service.MonitorSta resp.PrimaryStatus = summary.PrimaryStatus resp.PrimaryLatencyMs = summary.PrimaryLatencyMs resp.Availability7d = summary.Availability7d + resp.LatestQuota = summary.LatestQuota resp.ExtraModelsStatus = make([]dto.ChannelMonitorExtraModelStatus, 0, len(summary.ExtraModels)) for _, e := range summary.ExtraModels { resp.ExtraModelsStatus = append(resp.ExtraModelsStatus, dto.ChannelMonitorExtraModelStatus{ @@ -327,6 +352,8 @@ func (h *ChannelMonitorHandler) Create(c *gin.Context) { ExtraHeaders: req.ExtraHeaders, BodyOverrideMode: req.BodyOverrideMode, BodyOverride: req.BodyOverride, + CheckMode: req.CheckMode, + AccountID: req.AccountID, }) if err != nil { response.ErrorFrom(c, err) @@ -421,6 +448,8 @@ func (h *ChannelMonitorHandler) Update(c *gin.Context) { ExtraHeaders: req.ExtraHeaders, BodyOverrideMode: req.BodyOverrideMode, BodyOverride: req.BodyOverride, + CheckMode: req.CheckMode, + AccountID: req.AccountID, }) if err != nil { response.ErrorFrom(c, err) diff --git a/backend/internal/handler/channel_monitor_user_handler.go b/backend/internal/handler/channel_monitor_user_handler.go index 7e42c5a74d..173ca56dfe 100644 --- a/backend/internal/handler/channel_monitor_user_handler.go +++ b/backend/internal/handler/channel_monitor_user_handler.go @@ -3,6 +3,7 @@ package handler import ( "time" + "github.com/Wei-Shaw/sub2api/internal/domain" "github.com/Wei-Shaw/sub2api/internal/handler/admin" "github.com/Wei-Shaw/sub2api/internal/handler/dto" "github.com/Wei-Shaw/sub2api/internal/pkg/response" @@ -39,6 +40,15 @@ func (h *ChannelMonitorUserHandler) featureEnabled(c *gin.Context) bool { return runtime.Enabled && runtime.Mode == service.ChannelMonitorModeV1 } +// quotaVisible 返回用户端是否展示配额/余额快照(channel_monitor_show_quota, +// fail-closed:未配置/非 "true" 一律视为关闭)。settingService 为 nil 时 fail-closed。 +func (h *ChannelMonitorUserHandler) quotaVisible(c *gin.Context) bool { + if h.settingService == nil { + return false + } + return h.settingService.GetChannelMonitorRuntime(c.Request.Context()).ShowQuota +} + // --- Response --- type channelMonitorUserListItem struct { @@ -53,6 +63,9 @@ type channelMonitorUserListItem struct { Availability7d float64 `json:"availability_7d"` ExtraModels []dto.ChannelMonitorExtraModelStatus `json:"extra_models"` Timeline []channelMonitorUserTimelinePoint `json:"timeline"` + // LatestQuota 主模型最近配额快照;channel_monitor_show_quota=false 时 + // 由 userMonitorViewToItem 的调用方传入 false 剥离(服务端脱敏,非仅前端隐藏)。 + LatestQuota *domain.MonitorQuotaSnapshot `json:"latest_quota,omitempty"` } // channelMonitorUserTimelinePoint 主模型最近一次检测的 timeline 点。 @@ -82,7 +95,7 @@ type channelMonitorUserModelStat struct { AvgLatency7dMs *int `json:"avg_latency_7d_ms"` } -func userMonitorViewToItem(v *service.UserMonitorView) channelMonitorUserListItem { +func userMonitorViewToItem(v *service.UserMonitorView, includeQuota bool) channelMonitorUserListItem { extras := make([]dto.ChannelMonitorExtraModelStatus, 0, len(v.ExtraModels)) for _, e := range v.ExtraModels { extras = append(extras, dto.ChannelMonitorExtraModelStatus{ @@ -100,7 +113,7 @@ func userMonitorViewToItem(v *service.UserMonitorView) channelMonitorUserListIte CheckedAt: p.CheckedAt.UTC().Format(time.RFC3339), }) } - return channelMonitorUserListItem{ + item := channelMonitorUserListItem{ ID: v.ID, Name: v.Name, Provider: v.Provider, @@ -113,6 +126,10 @@ func userMonitorViewToItem(v *service.UserMonitorView) channelMonitorUserListIte ExtraModels: extras, Timeline: timeline, } + if includeQuota { + item.LatestQuota = v.LatestQuota + } + return item } func userMonitorDetailToResponse(d *service.UserMonitorDetail) *channelMonitorUserDetailResponse { @@ -150,9 +167,10 @@ func (h *ChannelMonitorUserHandler) List(c *gin.Context) { response.ErrorFrom(c, err) return } + includeQuota := h.quotaVisible(c) items := make([]channelMonitorUserListItem, 0, len(views)) for _, v := range views { - items = append(items, userMonitorViewToItem(v)) + items = append(items, userMonitorViewToItem(v, includeQuota)) } response.Success(c, gin.H{"items": items}) } diff --git a/backend/internal/service/channel_monitor_quota_fetcher.go b/backend/internal/service/channel_monitor_quota_fetcher.go index 83d8e0f455..d208dadc1b 100644 --- a/backend/internal/service/channel_monitor_quota_fetcher.go +++ b/backend/internal/service/channel_monitor_quota_fetcher.go @@ -65,19 +65,27 @@ type monitorQuotaCacheEntry struct { } // NewChannelMonitorQuotaFetcher 构造配额抓取器。 +// 参数取具体服务类型以便 wire 直连;单元测试在同包内用 struct 字面量注入 stub。 func NewChannelMonitorQuotaFetcher( - usage monitorUsageSource, - cnQuota monitorCNQuotaSource, - cnBalance monitorCNBalanceSource, - accounts monitorAccountSource, + usage *AccountUsageService, + cnQuota *CNProviderQuotaService, + cnBalance *CNProviderBalanceService, + accounts AccountRepository, ) *ChannelMonitorQuotaFetcher { - return &ChannelMonitorQuotaFetcher{ - usage: usage, - cnQuota: cnQuota, - cnBalance: cnBalance, - accounts: accounts, - cache: make(map[int64]monitorQuotaCacheEntry), + f := &ChannelMonitorQuotaFetcher{cache: make(map[int64]monitorQuotaCacheEntry)} + if usage != nil { + f.usage = usage } + if cnQuota != nil { + f.cnQuota = cnQuota + } + if cnBalance != nil { + f.cnBalance = cnBalance + } + if accounts != nil { + f.accounts = accounts + } + return f } // LoadAccount 加载账号(不走缓存)。供 Create/Update 时校验 diff --git a/backend/internal/service/wire.go b/backend/internal/service/wire.go index d6bb56fac4..88c3a1f185 100644 --- a/backend/internal/service/wire.go +++ b/backend/internal/service/wire.go @@ -907,6 +907,7 @@ var ProviderSet = wire.NewSet( ProvideBalanceNotifyService, ProvideChannelMonitorService, ProvideChannelMonitorRunner, + NewChannelMonitorQuotaFetcher, ProvideChannelMonitorV2Service, ProvideChannelMonitorV2Aggregator, NewChannelMonitorRequestTemplateService, @@ -965,13 +966,20 @@ func ProvideChannelMonitorService( // 通过 SetScheduler 注入回 service 后再 Start,确保启动时加载所有 enabled monitor, // 后续 CRUD 也能即时同步任务表。Runner.Stop 由 cleanup function 调用。 // settingService 用于 runner 每次 fire 读取功能开关。 -func ProvideChannelMonitorRunner(svc *ChannelMonitorService, settingService *SettingService) *ChannelMonitorRunner { +// quotaFetcher(账号侧用量聚合)也在此注入:accountUsage/CN 服务在 wire 图中 +// 晚于 channelMonitorService 构造,走 setter 注入避免调整既有构造顺序。 +func ProvideChannelMonitorRunner( + svc *ChannelMonitorService, + settingService *SettingService, + quotaFetcher *ChannelMonitorQuotaFetcher, +) *ChannelMonitorRunner { r := NewChannelMonitorRunner(svc, settingService) if svc != nil { // Ensure runtime reader is set even if ProvideChannelMonitorService // was constructed without settings (tests / alternate providers). svc.SetRuntimeReader(settingService) svc.SetScheduler(r) + svc.SetQuotaFetcher(quotaFetcher) } r.Start() return r