feat(channel-monitor-v2): 接入模式开关、路由门控与依赖注入

注册 admin/user 路由与 feature/mode 守卫,串联 Wire DI 与设置读写,
公开 channel_monitor_mode 与 hide_throughput 等运行时标志。
This commit is contained in:
IanShaw027
2026-08-07 11:05:32 +08:00
parent ead73264c7
commit a5beecb92a
20 changed files with 480 additions and 23 deletions
+8 -1
View File
@@ -110,6 +110,7 @@ func provideCleanup(
backupSvc *service.BackupService,
paymentOrderExpiry *service.PaymentOrderExpiryService,
channelMonitorRunner *service.ChannelMonitorRunner,
channelMonitorV2Aggregator *service.ChannelMonitorV2Aggregator,
quotaFlusher *service.UserPlatformQuotaUsageFlusher,
upstreamBillingProbe *service.UpstreamBillingProbeService,
ollamaCloudUsage *service.OllamaCloudUsageService,
@@ -319,7 +320,13 @@ func provideCleanup(
}
return nil
}},
{"ChannelMonitorRunner", func() error {
{"ChannelMonitorV2Aggregator", func() error {
if channelMonitorV2Aggregator != nil {
channelMonitorV2Aggregator.Stop()
}
return nil
}},
{"ChannelMonitorRunner", func() error {
if channelMonitorRunner != nil {
channelMonitorRunner.Stop()
}
+14 -3
View File
@@ -174,8 +174,11 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) {
announcementService := service.NewAnnouncementService(announcementRepository, announcementReadRepository, userRepository, userSubscriptionRepository)
announcementHandler := handler.NewAnnouncementHandler(announcementService)
channelMonitorRepository := repository.NewChannelMonitorRepository(client, db)
channelMonitorService := service.ProvideChannelMonitorService(channelMonitorRepository, secretEncryptor)
channelMonitorService := service.ProvideChannelMonitorService(channelMonitorRepository, secretEncryptor, settingService)
channelMonitorUserHandler := handler.NewChannelMonitorUserHandler(channelMonitorService, settingService)
channelMonitorV2Repository := repository.NewChannelMonitorV2Repository(db)
channelMonitorV2Service := service.ProvideChannelMonitorV2Service(channelMonitorV2Repository, settingService)
channelMonitorV2Handler := handler.NewChannelMonitorV2Handler(channelMonitorV2Service)
dashboardAggregationRepository := repository.NewDashboardAggregationRepository(db)
dashboardStatsCache := repository.NewDashboardCache(redisClient, configConfig)
dashboardService := service.NewDashboardService(usageLogRepository, dashboardAggregationRepository, dashboardStatsCache, configConfig)
@@ -308,7 +311,7 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) {
batchImageHandler := handler.ProvideBatchImageHandler(batchImagePublicService, batchImageDownloadService, batchImageCleanupService, openAIGatewayHandler)
idempotencyCoordinator := service.ProvideIdempotencyCoordinator(idempotencyRepository, configConfig)
idempotencyCleanupService := service.ProvideIdempotencyCleanupService(idempotencyRepository, configConfig)
handlers := handler.ProvideHandlers(authHandler, userHandler, apiKeyHandler, usageHandler, redeemHandler, subscriptionHandler, announcementHandler, channelMonitorUserHandler, adminHandlers, gatewayHandler, openAIGatewayHandler, handlerSettingHandler, totpHandler, passkeyHandler, handlerPaymentHandler, paymentWebhookHandler, availableChannelHandler, modelPlazaHandler, asyncImageHandler, batchImageHandler, idempotencyCoordinator, idempotencyCleanupService)
handlers := handler.ProvideHandlers(authHandler, userHandler, apiKeyHandler, usageHandler, redeemHandler, subscriptionHandler, announcementHandler, channelMonitorUserHandler, channelMonitorV2Handler, adminHandlers, gatewayHandler, openAIGatewayHandler, handlerSettingHandler, totpHandler, passkeyHandler, handlerPaymentHandler, paymentWebhookHandler, availableChannelHandler, modelPlazaHandler, asyncImageHandler, batchImageHandler, idempotencyCoordinator, idempotencyCleanupService)
jwtAuthMiddleware := middleware.NewJWTAuthMiddleware(authService, userService, settingService, auditLogService)
optionalJWTAuthMiddleware := middleware.NewOptionalJWTAuthMiddleware(authService, userService, settingService, auditLogService)
adminAuthMiddleware := middleware.NewAdminAuthMiddleware(authService, userService, settingService, auditLogService)
@@ -331,8 +334,9 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) {
scheduledTestRunnerService := service.ProvideScheduledTestRunnerService(scheduledTestPlanRepository, scheduledTestService, accountTestService, rateLimitService, configConfig)
paymentOrderExpiryService := service.ProvidePaymentOrderExpiryService(paymentService, leaderLockCache, db)
channelMonitorRunner := service.ProvideChannelMonitorRunner(channelMonitorService, settingService)
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, openAICodexVersionSyncService, proxyExpiryService, subscriptionExpiryService, usageCleanupService, idempotencyCleanupService, batchImageCleanupService, batchImageWorkerRuntime, pricingService, emailQueueService, billingCacheService, usageRecordWorkerPool, subscriptionService, oAuthService, openAIOAuthService, geminiOAuthService, antigravityOAuthService, grokOAuthService, openAIGatewayService, scheduledTestRunnerService, backupService, paymentOrderExpiryService, channelMonitorRunner, userPlatformQuotaUsageFlusher, upstreamBillingProbeService, ollamaCloudUsageService, auditLogService, promptService)
v := provideCleanup(client, redisClient, opsMetricsCollector, opsAggregationService, opsAlertEvaluatorService, opsCleanupService, opsScheduledReportService, opsSystemLogSink, opsService, opsIngressRejectAggregator, apiKeyService, authCacheInvalidationWorker, schedulerSnapshotService, tokenRefreshService, accountExpiryService, 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)
application := &Application{
Server: httpServer,
PromptAudit: promptService,
@@ -398,6 +402,7 @@ func provideCleanup(
backupSvc *service.BackupService,
paymentOrderExpiry *service.PaymentOrderExpiryService,
channelMonitorRunner *service.ChannelMonitorRunner,
channelMonitorV2Aggregator *service.ChannelMonitorV2Aggregator,
quotaFlusher *service.UserPlatformQuotaUsageFlusher,
upstreamBillingProbe *service.UpstreamBillingProbeService,
ollamaCloudUsage *service.OllamaCloudUsageService,
@@ -606,6 +611,12 @@ func provideCleanup(
}
return nil
}},
{"ChannelMonitorV2Aggregator", func() error {
if channelMonitorV2Aggregator != nil {
channelMonitorV2Aggregator.Stop()
}
return nil
}},
{"ChannelMonitorRunner", func() error {
if channelMonitorRunner != nil {
channelMonitorRunner.Stop()
+1
View File
@@ -88,6 +88,7 @@ func TestProvideCleanup_WithMinimalDependencies_NoPanic(t *testing.T) {
nil, // backupSvc
nil, // paymentOrderExpiry
nil, // channelMonitorRunner
nil, // channelMonitorV2Aggregator
nil, // quotaFlusher
nil, // upstreamBillingProbe
nil, // ollamaCloudUsage
@@ -370,7 +370,9 @@ func (h *SettingHandler) GetSettings(c *gin.Context) {
PaymentAlipayMobilePrecreateDeepLink: paymentCfg.AlipayMobilePrecreateDeepLink,
ChannelMonitorEnabled: settings.ChannelMonitorEnabled,
ChannelMonitorMode: settings.ChannelMonitorMode,
ChannelMonitorDefaultIntervalSeconds: settings.ChannelMonitorDefaultIntervalSeconds,
ChannelMonitorHideThroughput: settings.ChannelMonitorHideThroughput,
AvailableChannelsEnabled: settings.AvailableChannelsEnabled,
@@ -327,8 +327,10 @@ type UpdateSettingsRequest struct {
PaymentAlipayMobilePrecreateDeepLink *bool `json:"payment_alipay_mobile_precreate_deep_link"`
// Channel Monitor feature switch
ChannelMonitorEnabled *bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds *int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorEnabled *bool `json:"channel_monitor_enabled"`
ChannelMonitorMode *string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds *int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput *bool `json:"channel_monitor_hide_throughput"`
// Available Channels feature switch (user-facing)
AvailableChannelsEnabled *bool `json:"available_channels_enabled"`
@@ -1854,12 +1856,24 @@ func (h *SettingHandler) UpdateSettings(c *gin.Context) {
}
return previousSettings.ChannelMonitorEnabled
}(),
ChannelMonitorMode: func() string {
if req.ChannelMonitorMode != nil {
return *req.ChannelMonitorMode
}
return previousSettings.ChannelMonitorMode
}(),
ChannelMonitorDefaultIntervalSeconds: func() int {
if req.ChannelMonitorDefaultIntervalSeconds != nil {
return *req.ChannelMonitorDefaultIntervalSeconds
}
return previousSettings.ChannelMonitorDefaultIntervalSeconds
}(),
ChannelMonitorHideThroughput: func() bool {
if req.ChannelMonitorHideThroughput != nil {
return *req.ChannelMonitorHideThroughput
}
return previousSettings.ChannelMonitorHideThroughput
}(),
AvailableChannelsEnabled: func() bool {
if req.AvailableChannelsEnabled != nil {
return *req.AvailableChannelsEnabled
@@ -2291,7 +2305,9 @@ func (h *SettingHandler) UpdateSettings(c *gin.Context) {
PaymentAlipayMobilePrecreateDeepLink: updatedPaymentCfg.AlipayMobilePrecreateDeepLink,
ChannelMonitorEnabled: updatedSettings.ChannelMonitorEnabled,
ChannelMonitorMode: updatedSettings.ChannelMonitorMode,
ChannelMonitorDefaultIntervalSeconds: updatedSettings.ChannelMonitorDefaultIntervalSeconds,
ChannelMonitorHideThroughput: updatedSettings.ChannelMonitorHideThroughput,
AvailableChannelsEnabled: updatedSettings.AvailableChannelsEnabled,
+8 -4
View File
@@ -300,8 +300,10 @@ type SystemSettings struct {
AccountQuotaNotifyEmails []NotifyEmailEntry `json:"account_quota_notify_emails"`
// Channel Monitor feature switch
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
// Available Channels feature switch (user-facing aggregate view)
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
@@ -398,8 +400,10 @@ type PublicSettings struct {
BalanceLowNotifyThreshold float64 `json:"balance_low_notify_threshold"`
BalanceLowNotifyRechargeURL string `json:"balance_low_notify_recharge_url"`
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
+1
View File
@@ -53,6 +53,7 @@ type Handlers struct {
Subscription *SubscriptionHandler
Announcement *AnnouncementHandler
ChannelMonitor *ChannelMonitorUserHandler
ChannelMonitorV2 *ChannelMonitorV2Handler
Admin *AdminHandlers
Gateway *GatewayHandler
OpenAIGateway *OpenAIGatewayHandler
@@ -103,7 +103,9 @@ func (h *SettingHandler) GetPublicSettings(c *gin.Context) {
BalanceLowNotifyRechargeURL: settings.BalanceLowNotifyRechargeURL,
ChannelMonitorEnabled: settings.ChannelMonitorEnabled,
ChannelMonitorMode: settings.ChannelMonitorMode,
ChannelMonitorDefaultIntervalSeconds: settings.ChannelMonitorDefaultIntervalSeconds,
ChannelMonitorHideThroughput: settings.ChannelMonitorHideThroughput,
AvailableChannelsEnabled: settings.AvailableChannelsEnabled,
+3
View File
@@ -175,6 +175,7 @@ func ProvideHandlers(
subscriptionHandler *SubscriptionHandler,
announcementHandler *AnnouncementHandler,
channelMonitorUserHandler *ChannelMonitorUserHandler,
channelMonitorV2Handler *ChannelMonitorV2Handler,
adminHandlers *AdminHandlers,
gatewayHandler *GatewayHandler,
openaiGatewayHandler *OpenAIGatewayHandler,
@@ -199,6 +200,7 @@ func ProvideHandlers(
Subscription: subscriptionHandler,
Announcement: announcementHandler,
ChannelMonitor: channelMonitorUserHandler,
ChannelMonitorV2: channelMonitorV2Handler,
Admin: adminHandlers,
Gateway: gatewayHandler,
OpenAIGateway: openaiGatewayHandler,
@@ -225,6 +227,7 @@ var ProviderSet = wire.NewSet(
NewSubscriptionHandler,
NewAnnouncementHandler,
NewChannelMonitorUserHandler,
NewChannelMonitorV2Handler,
ProvideGatewayHandler,
ProvideOpenAIGatewayHandler,
NewTotpHandler,
+1
View File
@@ -98,6 +98,7 @@ var ProviderSet = wire.NewSet(
NewTLSFingerprintProfileRepository,
NewChannelRepository,
NewChannelMonitorRepository,
NewChannelMonitorV2Repository,
NewChannelMonitorRequestTemplateRepository,
NewContentModerationRepository,
NewAffiliateRepository,
+68 -3
View File
@@ -1,10 +1,10 @@
// Package routes provides HTTP route registration and handlers.
package routes
import (
"github.com/Wei-Shaw/sub2api/internal/handler"
"github.com/Wei-Shaw/sub2api/internal/server/middleware"
"github.com/Wei-Shaw/sub2api/internal/service"
"github.com/Wei-Shaw/sub2api/internal/pkg/response"
"github.com/gin-gonic/gin"
)
@@ -106,7 +106,8 @@ func RegisterAdminRoutes(
registerChannelRoutes(admin, h)
// 渠道监控
registerChannelMonitorRoutes(admin, h)
registerChannelMonitorRoutes(admin, h, settingService)
registerChannelMonitorV2Routes(admin, h, settingService)
// 风控中心
registerContentModerationRoutes(admin, h)
@@ -729,8 +730,10 @@ func registerChannelRoutes(admin *gin.RouterGroup, h *handler.Handlers) {
}
}
func registerChannelMonitorRoutes(admin *gin.RouterGroup, h *handler.Handlers) {
func registerChannelMonitorRoutes(admin *gin.RouterGroup, h *handler.Handlers, settingService *service.SettingService) {
guard := channelMonitorAdminFeatureGuard(settingService)
monitors := admin.Group("/channel-monitors")
monitors.Use(guard)
{
monitors.GET("", h.Admin.ChannelMonitor.List)
monitors.POST("", h.Admin.ChannelMonitor.Create)
@@ -743,6 +746,7 @@ func registerChannelMonitorRoutes(admin *gin.RouterGroup, h *handler.Handlers) {
}
templates := admin.Group("/channel-monitor-templates")
templates.Use(guard)
{
templates.GET("", h.Admin.ChannelMonitorTemplate.List)
templates.POST("", h.Admin.ChannelMonitorTemplate.Create)
@@ -773,3 +777,64 @@ func registerAffiliateRoutes(admin *gin.RouterGroup, h *handler.Handlers) {
}
}
}
func registerChannelMonitorV2Routes(admin *gin.RouterGroup, h *handler.Handlers, settingService *service.SettingService) {
// Config GET/PUT: feature enabled only (operators can prepare V2 before flipping mode).
// Read/matrix endpoints: require mode=v2 so V1 deployments do not serve passive data.
featureGuard := channelMonitorAdminFeatureGuard(settingService)
modeV2Guard := channelMonitorModeV2Guard(settingService)
monitor := admin.Group("/channel-monitor-v2")
{
config := monitor.Group("")
config.Use(featureGuard)
{
config.GET("/config", h.ChannelMonitorV2.GetConfig)
config.PUT("/config", h.ChannelMonitorV2.UpdateConfig)
}
reads := monitor.Group("")
reads.Use(modeV2Guard)
{
reads.GET("/dimensions", h.ChannelMonitorV2.Dimensions)
reads.GET("/snapshot", h.ChannelMonitorV2.AdminSnapshot)
reads.GET("/models", h.ChannelMonitorV2.AdminModels)
reads.GET("/matrix", h.ChannelMonitorV2.AdminMatrix)
reads.GET("/errors", h.ChannelMonitorV2.Errors)
reads.GET("/users", h.ChannelMonitorV2.AdminUsers)
}
}
}
func channelMonitorAdminFeatureGuard(settingService *service.SettingService) gin.HandlerFunc {
return func(c *gin.Context) {
if settingService != nil && settingService.GetChannelMonitorRuntime(c.Request.Context()).Enabled {
c.Next()
return
}
response.ErrorFrom(c, service.ErrChannelMonitorDisabled)
c.Abort()
}
}
// channelMonitorModeV2Guard requires feature enabled and channel_monitor_mode=v2.
func channelMonitorModeV2Guard(settingService *service.SettingService) gin.HandlerFunc {
return func(c *gin.Context) {
if settingService == nil {
response.ErrorFrom(c, service.ErrChannelMonitorDisabled)
c.Abort()
return
}
rt := settingService.GetChannelMonitorRuntime(c.Request.Context())
if !rt.Enabled {
response.ErrorFrom(c, service.ErrChannelMonitorDisabled)
c.Abort()
return
}
if !rt.PassiveAggregationAllowed() {
response.ErrorFrom(c, service.ErrChannelMonitorModeMismatch)
c.Abort()
return
}
c.Next()
}
}
@@ -0,0 +1,173 @@
package routes
import (
"context"
"net/http"
"net/http/httptest"
"testing"
"github.com/Wei-Shaw/sub2api/internal/config"
"github.com/Wei-Shaw/sub2api/internal/service"
"github.com/gin-gonic/gin"
"github.com/stretchr/testify/require"
)
// channelMonitorRouteSettingRepoStub is a minimal SettingRepository for route guards.
type channelMonitorRouteSettingRepoStub struct {
values map[string]string
}
func (s *channelMonitorRouteSettingRepoStub) Get(context.Context, string) (*service.Setting, error) {
panic("unexpected Get call")
}
func (s *channelMonitorRouteSettingRepoStub) GetValue(_ context.Context, key string) (string, error) {
return s.values[key], nil
}
func (s *channelMonitorRouteSettingRepoStub) Set(context.Context, string, string) error {
panic("unexpected Set call")
}
func (s *channelMonitorRouteSettingRepoStub) GetMultiple(_ context.Context, keys []string) (map[string]string, error) {
out := make(map[string]string, len(keys))
for _, key := range keys {
if value, ok := s.values[key]; ok {
out[key] = value
}
}
return out, nil
}
func (s *channelMonitorRouteSettingRepoStub) SetMultiple(context.Context, map[string]string) error {
panic("unexpected SetMultiple call")
}
func (s *channelMonitorRouteSettingRepoStub) GetAll(context.Context) (map[string]string, error) {
panic("unexpected GetAll call")
}
func (s *channelMonitorRouteSettingRepoStub) Delete(context.Context, string) error {
panic("unexpected Delete call")
}
func newChannelMonitorRouteSettings(enabled bool) *service.SettingService {
value := "false"
if enabled {
value = "true"
}
return service.NewSettingService(&channelMonitorRouteSettingRepoStub{
values: map[string]string{
service.SettingKeyChannelMonitorEnabled: value,
},
}, &config.Config{})
}
func newChannelMonitorModeSettings(enabled bool, mode string) *service.SettingService {
enabledVal := "false"
if enabled {
enabledVal = "true"
}
return service.NewSettingService(&channelMonitorRouteSettingRepoStub{
values: map[string]string{
service.SettingKeyChannelMonitorEnabled: enabledVal,
service.SettingKeyChannelMonitorMode: mode,
},
}, &config.Config{})
}
func TestChannelMonitorAdminFeatureGuard(t *testing.T) {
tests := []struct {
name string
svc *service.SettingService
wantStatus int
}{
{
name: "nil setting service blocks",
svc: nil,
wantStatus: http.StatusForbidden,
},
{
name: "disabled blocks",
svc: newChannelMonitorRouteSettings(false),
wantStatus: http.StatusForbidden,
},
{
name: "enabled allows",
svc: newChannelMonitorRouteSettings(true),
wantStatus: http.StatusOK,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gin.SetMode(gin.TestMode)
router := gin.New()
router.Use(channelMonitorAdminFeatureGuard(tt.svc))
router.GET("/test", func(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"ok": true})
})
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/test", nil)
router.ServeHTTP(rec, req)
require.Equal(t, tt.wantStatus, rec.Code)
if tt.wantStatus == http.StatusForbidden {
require.Contains(t, rec.Body.String(), "CHANNEL_MONITOR_DISABLED")
}
})
}
}
func TestChannelMonitorModeV2Guard(t *testing.T) {
tests := []struct {
name string
svc *service.SettingService
wantStatus int
wantCode string
}{
{
name: "nil blocks as disabled",
svc: nil,
wantStatus: http.StatusForbidden,
wantCode: "CHANNEL_MONITOR_DISABLED",
},
{
name: "feature off blocks",
svc: newChannelMonitorModeSettings(false, service.ChannelMonitorModeV2),
wantStatus: http.StatusForbidden,
wantCode: "CHANNEL_MONITOR_DISABLED",
},
{
name: "mode v1 blocks with mode mismatch",
svc: newChannelMonitorModeSettings(true, service.ChannelMonitorModeV1),
wantStatus: http.StatusForbidden,
wantCode: "CHANNEL_MONITOR_MODE_MISMATCH",
},
{
name: "mode v2 allows",
svc: newChannelMonitorModeSettings(true, service.ChannelMonitorModeV2),
wantStatus: http.StatusOK,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gin.SetMode(gin.TestMode)
router := gin.New()
router.Use(channelMonitorModeV2Guard(tt.svc))
router.GET("/test", func(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"ok": true})
})
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(context.Background(), http.MethodGet, "/test", nil)
router.ServeHTTP(rec, req)
require.Equal(t, tt.wantStatus, rec.Code)
if tt.wantCode != "" {
require.Contains(t, rec.Body.String(), tt.wantCode)
}
})
}
}
+13
View File
@@ -141,5 +141,18 @@ func RegisterUserRoutes(
monitors.GET("", h.ChannelMonitor.List)
monitors.GET("/:id/status", h.ChannelMonitor.GetStatus)
}
// V2 passive views require feature on + mode=v2.
monitorV2 := authenticated.Group("/channel-monitor-v2")
monitorV2.Use(panelRateLimiter.Heavy())
monitorV2.Use(channelMonitorModeV2Guard(settingService))
{
monitorV2.GET("/dimensions", h.ChannelMonitorV2.Dimensions)
monitorV2.GET("/snapshot", h.ChannelMonitorV2.Snapshot)
monitorV2.GET("/models", h.ChannelMonitorV2.Models)
monitorV2.GET("/matrix", h.ChannelMonitorV2.Matrix)
monitorV2.GET("/errors", h.ChannelMonitorV2.Errors)
monitorV2.GET("/users", h.ChannelMonitorV2.Users)
}
}
}
@@ -393,10 +393,24 @@ const (
// When false: runner skips scheduling and user-facing endpoints return an empty list.
SettingKeyChannelMonitorEnabled = "channel_monitor_enabled"
// SettingKeyChannelMonitorMode selects exclusive implementation:
// "v1" active probes, "v2" passive aggregation. Default "v2".
SettingKeyChannelMonitorMode = "channel_monitor_mode"
// ChannelMonitorModeV1/V2 are the only accepted mode values.
ChannelMonitorModeV1 = "v1"
ChannelMonitorModeV2 = "v2"
// SettingKeyChannelMonitorDefaultIntervalSeconds controls the default interval (seconds)
// pre-filled when creating a new channel monitor from the admin UI. Range: [15, 3600].
SettingKeyChannelMonitorDefaultIntervalSeconds = "channel_monitor_default_interval_seconds"
// SettingKeyChannelMonitorHideThroughput hides RPM/TPM (and similar absolute
// throughput rates) from non-admin user-facing monitor APIs and UI, so users
// cannot reverse-estimate fleet volume from rates × window length.
// Default false (show rates). Admin endpoints always keep full metrics.
SettingKeyChannelMonitorHideThroughput = "channel_monitor_hide_throughput"
// SettingKeyAvailableChannelsEnabled is a DB-backed soft switch for the "Available Channels"
// user-facing aggregate view. When false: user endpoint returns an empty list and the
// sidebar entry is hidden. Defaults to false (opt-in feature).
@@ -185,7 +185,9 @@ func (s *SettingService) InitializeDefaultSettings(ctx context.Context) error {
// Channel monitor defaults (enabled, 60s)
SettingKeyChannelMonitorEnabled: "true",
SettingKeyChannelMonitorMode: ChannelMonitorModeV2,
SettingKeyChannelMonitorDefaultIntervalSeconds: "60",
SettingKeyChannelMonitorHideThroughput: "false",
// Available channels feature (default disabled; opt-in)
SettingKeyAvailableChannelsEnabled: "false",
@@ -781,9 +783,11 @@ func (s *SettingService) parseSettings(settings map[string]string) *SystemSettin
// Channel monitor feature (default: enabled, 60s)
result.ChannelMonitorEnabled = !isFalseSettingValue(settings[SettingKeyChannelMonitorEnabled])
result.ChannelMonitorMode = normalizeChannelMonitorMode(settings[SettingKeyChannelMonitorMode])
result.ChannelMonitorDefaultIntervalSeconds = parseChannelMonitorInterval(
settings[SettingKeyChannelMonitorDefaultIntervalSeconds],
)
result.ChannelMonitorHideThroughput = settings[SettingKeyChannelMonitorHideThroughput] == "true"
// Available channels feature (default: disabled; strict true)
result.AvailableChannelsEnabled = settings[SettingKeyAvailableChannelsEnabled] == "true"
+53 -4
View File
@@ -227,7 +227,9 @@ func (s *SettingService) GetPublicSettings(ctx context.Context) (*PublicSettings
SettingKeyBalanceLowNotifyRechargeURL,
SettingKeyAccountQuotaNotifyEnabled,
SettingKeyChannelMonitorEnabled,
SettingKeyChannelMonitorMode,
SettingKeyChannelMonitorDefaultIntervalSeconds,
SettingKeyChannelMonitorHideThroughput,
SettingKeyAvailableChannelsEnabled,
SettingKeyModelPlazaEnabled,
SettingKeyModelPlazaRequireAuth,
@@ -348,7 +350,9 @@ func (s *SettingService) GetPublicSettings(ctx context.Context) (*PublicSettings
BalanceLowNotifyRechargeURL: settings[SettingKeyBalanceLowNotifyRechargeURL],
ChannelMonitorEnabled: !isFalseSettingValue(settings[SettingKeyChannelMonitorEnabled]),
ChannelMonitorMode: normalizeChannelMonitorMode(settings[SettingKeyChannelMonitorMode]),
ChannelMonitorDefaultIntervalSeconds: parseChannelMonitorInterval(settings[SettingKeyChannelMonitorDefaultIntervalSeconds]),
ChannelMonitorHideThroughput: settings[SettingKeyChannelMonitorHideThroughput] == "true",
AvailableChannelsEnabled: settings[SettingKeyAvailableChannelsEnabled] == "true",
@@ -369,8 +373,21 @@ const (
channelMonitorIntervalMin = 15
channelMonitorIntervalMax = 3600
channelMonitorIntervalFallback = 60
defaultChannelMonitorMode = ChannelMonitorModeV2
)
// normalizeChannelMonitorMode accepts only v1/v2; empty/invalid → v2.
func normalizeChannelMonitorMode(raw string) string {
switch strings.ToLower(strings.TrimSpace(raw)) {
case ChannelMonitorModeV1:
return ChannelMonitorModeV1
case ChannelMonitorModeV2, "":
return ChannelMonitorModeV2
default:
return defaultChannelMonitorMode
}
}
// parseChannelMonitorInterval parses the stored string and clamps to [15, 3600].
// Empty / invalid input falls back to channelMonitorIntervalFallback.
func parseChannelMonitorInterval(raw string) int {
@@ -396,25 +413,53 @@ func clampChannelMonitorInterval(v int) int {
}
// ChannelMonitorRuntime is the lightweight view of the channel monitor feature
// consumed by the runner and user-facing handlers.
// consumed by the runner, V2 aggregator, and user-facing handlers.
type ChannelMonitorRuntime struct {
Enabled bool
Mode string // ChannelMonitorModeV1 or ChannelMonitorModeV2
DefaultIntervalSeconds int
// HideThroughput: when true, user-facing V2 APIs omit RPM/TPM scale signals.
HideThroughput bool
}
// ActiveProbesAllowed reports whether V1 active provider probes may run.
func (r ChannelMonitorRuntime) ActiveProbesAllowed() bool {
return r.Enabled && r.Mode == ChannelMonitorModeV1
}
// PassiveAggregationAllowed reports whether V2 passive aggregation may run.
func (r ChannelMonitorRuntime) PassiveAggregationAllowed() bool {
return r.Enabled && r.Mode == ChannelMonitorModeV2
}
// GetChannelMonitorRuntime reads the channel monitor feature flags directly from
// the settings store. Fail-open: on error returns Enabled=true with the default interval.
// the settings store. Fail-open: on error returns Enabled=true, Mode=v2, default interval.
func (s *SettingService) GetChannelMonitorRuntime(ctx context.Context) ChannelMonitorRuntime {
if s == nil || s.settingRepo == nil {
return ChannelMonitorRuntime{
Enabled: true,
Mode: defaultChannelMonitorMode,
DefaultIntervalSeconds: channelMonitorIntervalFallback,
}
}
vals, err := s.settingRepo.GetMultiple(ctx, []string{
SettingKeyChannelMonitorEnabled,
SettingKeyChannelMonitorMode,
SettingKeyChannelMonitorDefaultIntervalSeconds,
SettingKeyChannelMonitorHideThroughput,
})
if err != nil {
return ChannelMonitorRuntime{Enabled: true, DefaultIntervalSeconds: channelMonitorIntervalFallback}
return ChannelMonitorRuntime{
Enabled: true,
Mode: defaultChannelMonitorMode,
DefaultIntervalSeconds: channelMonitorIntervalFallback,
}
}
return ChannelMonitorRuntime{
Enabled: !isFalseSettingValue(vals[SettingKeyChannelMonitorEnabled]),
Mode: normalizeChannelMonitorMode(vals[SettingKeyChannelMonitorMode]),
DefaultIntervalSeconds: parseChannelMonitorInterval(vals[SettingKeyChannelMonitorDefaultIntervalSeconds]),
HideThroughput: vals[SettingKeyChannelMonitorHideThroughput] == "true",
}
}
@@ -552,7 +597,10 @@ type PublicSettingsInjectionPayload struct {
// that hid the "可用渠道" menu on page refresh.
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
// ChannelMonitorHideThroughput is public so the user UI can hide RPM/TPM
// without waiting for API redaction alone (defense in depth).
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
ModelPlazaEnabled bool `json:"model_plaza_enabled"`
ModelPlazaRequireAuth bool `json:"model_plaza_require_auth"`
AffiliateEnabled bool `json:"affiliate_enabled"`
@@ -628,6 +676,7 @@ func (s *SettingService) GetPublicSettingsForInjection(ctx context.Context) (any
ChannelMonitorEnabled: settings.ChannelMonitorEnabled,
ChannelMonitorDefaultIntervalSeconds: settings.ChannelMonitorDefaultIntervalSeconds,
ChannelMonitorHideThroughput: settings.ChannelMonitorHideThroughput,
AvailableChannelsEnabled: settings.AvailableChannelsEnabled,
ModelPlazaEnabled: settings.ModelPlazaEnabled,
ModelPlazaRequireAuth: settings.ModelPlazaRequireAuth,
@@ -10,6 +10,7 @@ import (
"github.com/Wei-Shaw/sub2api/internal/config"
infraerrors "github.com/Wei-Shaw/sub2api/internal/pkg/errors"
"golang.org/x/sync/singleflight"
"sync"
)
var (
@@ -78,6 +79,9 @@ type SettingService struct {
// instance owns its own cache, no shared package-level state.
openAIQuotaAutoPauseSettingsCache atomic.Value // *cachedOpenAIQuotaAutoPauseSettings
openAIQuotaAutoPauseSettingsSF singleflight.Group
channelMonitorRuntimeListenersMu sync.Mutex
channelMonitorRuntimeListeners []func()
}
// DefaultPlatformQuotaSetting 单 platform 三档限额(nil = 沿用上层;0 = 显式禁用;>0 = 上限)
@@ -302,6 +306,53 @@ func (s *SettingService) SetOnUpdateCallback(callback func()) {
s.onUpdate = callback
}
// SubscribeChannelMonitorRuntime registers a listener that is invoked after
// settings are successfully persisted (and process caches refreshed).
// Used by ChannelMonitorRunner / ChannelMonitorV2Aggregator for immediate
// mode flips without waiting for poll intervals.
func (s *SettingService) SubscribeChannelMonitorRuntime(listener func()) (unsubscribe func()) {
if s == nil || listener == nil {
return func() {}
}
s.channelMonitorRuntimeListenersMu.Lock()
s.channelMonitorRuntimeListeners = append(s.channelMonitorRuntimeListeners, listener)
idx := len(s.channelMonitorRuntimeListeners) - 1
s.channelMonitorRuntimeListenersMu.Unlock()
return func() {
s.channelMonitorRuntimeListenersMu.Lock()
defer s.channelMonitorRuntimeListenersMu.Unlock()
if idx < 0 || idx >= len(s.channelMonitorRuntimeListeners) {
return
}
s.channelMonitorRuntimeListeners[idx] = nil
}
}
func (s *SettingService) notifyChannelMonitorRuntimeListeners() {
if s == nil {
return
}
s.channelMonitorRuntimeListenersMu.Lock()
listeners := make([]func(), 0, len(s.channelMonitorRuntimeListeners))
for _, l := range s.channelMonitorRuntimeListeners {
if l != nil {
listeners = append(listeners, l)
}
}
s.channelMonitorRuntimeListenersMu.Unlock()
for _, l := range listeners {
func(fn func()) {
defer func() {
if recovered := recover(); recovered != nil {
// keep settings path healthy
}
}()
fn()
}(l)
}
}
// SetVersion sets the application version for injection into public settings
func (s *SettingService) SetVersion(version string) {
s.version = version
@@ -411,9 +411,11 @@ func (s *SettingService) buildSystemSettingsUpdates(ctx context.Context, setting
// Channel monitor feature switch
updates[SettingKeyChannelMonitorEnabled] = strconv.FormatBool(settings.ChannelMonitorEnabled)
updates[SettingKeyChannelMonitorMode] = normalizeChannelMonitorMode(settings.ChannelMonitorMode)
if v := clampChannelMonitorInterval(settings.ChannelMonitorDefaultIntervalSeconds); v > 0 {
updates[SettingKeyChannelMonitorDefaultIntervalSeconds] = strconv.Itoa(v)
}
updates[SettingKeyChannelMonitorHideThroughput] = strconv.FormatBool(settings.ChannelMonitorHideThroughput)
// Available channels feature switch
updates[SettingKeyAvailableChannelsEnabled] = strconv.FormatBool(settings.AvailableChannelsEnabled)
@@ -681,6 +683,7 @@ func (s *SettingService) refreshCachedSettings(settings *SystemSettings) {
if s.onUpdate != nil {
s.onUpdate() // Invalidate cache after settings update
}
s.notifyChannelMonitorRuntimeListeners()
}
func (s *SettingService) defaultRewriteMessageCacheControl() bool {
+8 -4
View File
@@ -196,8 +196,10 @@ type SystemSettings struct {
OpsMetricsIntervalSeconds int
// Channel Monitor feature
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
// Available Channels feature (user-facing aggregate view)
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
@@ -362,8 +364,10 @@ type PublicSettings struct {
BalanceLowNotifyRechargeURL string
// Channel Monitor feature
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorEnabled bool `json:"channel_monitor_enabled"`
ChannelMonitorMode string `json:"channel_monitor_mode"`
ChannelMonitorDefaultIntervalSeconds int `json:"channel_monitor_default_interval_seconds"`
ChannelMonitorHideThroughput bool `json:"channel_monitor_hide_throughput"`
// Available Channels feature (user-facing aggregate view)
AvailableChannelsEnabled bool `json:"available_channels_enabled"`
+35 -2
View File
@@ -1,6 +1,7 @@
package service
import (
"os"
"context"
"database/sql"
"time"
@@ -845,6 +846,8 @@ var ProviderSet = wire.NewSet(
ProvideBalanceNotifyService,
ProvideChannelMonitorService,
ProvideChannelMonitorRunner,
ProvideChannelMonitorV2Service,
ProvideChannelMonitorV2Aggregator,
NewChannelMonitorRequestTemplateService,
ProvideUserPlatformQuotaUsageFlusher,
)
@@ -886,11 +889,15 @@ func ProvidePaymentOrderExpiryService(paymentSvc *PaymentService, lockCache Lead
// ProvideChannelMonitorService 创建渠道监控服务(CRUD + RunCheck + 用户视图聚合)。
// 加密器复用 wire 中已注入的 SecretEncryptor(AES-256-GCM)。
// settingService gates RunCheck via channel_monitor_enabled + channel_monitor_mode.
func ProvideChannelMonitorService(
repo ChannelMonitorRepository,
encryptor SecretEncryptor,
settingService *SettingService,
) *ChannelMonitorService {
return NewChannelMonitorService(repo, encryptor)
svc := NewChannelMonitorService(repo, encryptor)
svc.SetRuntimeReader(settingService)
return svc
}
// ProvideChannelMonitorRunner 创建并启动渠道监控调度器。
@@ -899,7 +906,33 @@ func ProvideChannelMonitorService(
// settingService 用于 runner 每次 fire 读取功能开关。
func ProvideChannelMonitorRunner(svc *ChannelMonitorService, settingService *SettingService) *ChannelMonitorRunner {
r := NewChannelMonitorRunner(svc, settingService)
svc.SetScheduler(r)
if svc != nil {
// Ensure runtime reader is set even if ProvideChannelMonitorService
// was constructed without settings (tests / alternate providers).
svc.SetRuntimeReader(settingService)
svc.SetScheduler(r)
}
r.Start()
return r
}
// ProvideChannelMonitorV2Service wires settings for user-facing privacy flags
// (e.g. hide RPM/TPM throughput).
func ProvideChannelMonitorV2Service(repo ChannelMonitorV2Repository, settingService *SettingService) *ChannelMonitorV2Service {
svc := NewChannelMonitorV2Service(repo)
svc.SetRuntimeReader(settingService)
return svc
}
// ProvideChannelMonitorV2Aggregator starts the passive minute-rollup worker.
// Aggregation only runs when channel_monitor_enabled=true and mode=v2 (and V2 config enabled).
// Set CHANNEL_MONITOR_V2_DISABLE_AGGREGATOR=1 to skip Start (local demo with seeded facts).
func ProvideChannelMonitorV2Aggregator(repo ChannelMonitorV2Repository, db *sql.DB, settingService *SettingService) *ChannelMonitorV2Aggregator {
aggregator := NewChannelMonitorV2Aggregator(repo, db, settingService)
if os.Getenv("CHANNEL_MONITOR_V2_DISABLE_AGGREGATOR") == "1" {
return aggregator
}
aggregator.Start()
return aggregator
}