feat(monitor): wire quota fetcher & expose check_mode in handlers
- 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
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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})
|
||||
}
|
||||
|
||||
@@ -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 时校验
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user