From a5beecb92ace6f584e6ef9c9c2d800a5710a4e8d Mon Sep 17 00:00:00 2001 From: IanShaw027 Date: Fri, 7 Aug 2026 11:01:22 +0800 Subject: [PATCH] =?UTF-8?q?feat(channel-monitor-v2):=20=E6=8E=A5=E5=85=A5?= =?UTF-8?q?=E6=A8=A1=E5=BC=8F=E5=BC=80=E5=85=B3=E3=80=81=E8=B7=AF=E7=94=B1?= =?UTF-8?q?=E9=97=A8=E6=8E=A7=E4=B8=8E=E4=BE=9D=E8=B5=96=E6=B3=A8=E5=85=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 注册 admin/user 路由与 feature/mode 守卫,串联 Wire DI 与设置读写, 公开 channel_monitor_mode 与 hide_throughput 等运行时标志。 --- backend/cmd/server/wire.go | 9 +- backend/cmd/server/wire_gen.go | 17 +- backend/cmd/server/wire_gen_test.go | 1 + .../internal/handler/admin/setting_handler.go | 2 + .../handler/admin/setting_handler_update.go | 20 +- backend/internal/handler/dto/settings.go | 12 +- backend/internal/handler/handler.go | 1 + backend/internal/handler/setting_handler.go | 2 + backend/internal/handler/wire.go | 3 + backend/internal/repository/wire.go | 1 + backend/internal/server/routes/admin.go | 71 ++++++- .../channel_monitor_feature_gate_test.go | 173 ++++++++++++++++++ backend/internal/server/routes/user.go | 13 ++ backend/internal/service/domain_constants.go | 14 ++ backend/internal/service/setting_parse.go | 4 + backend/internal/service/setting_public.go | 57 +++++- backend/internal/service/setting_service.go | 51 ++++++ backend/internal/service/setting_update.go | 3 + backend/internal/service/settings_view.go | 12 +- backend/internal/service/wire.go | 37 +++- 20 files changed, 480 insertions(+), 23 deletions(-) create mode 100644 backend/internal/server/routes/channel_monitor_feature_gate_test.go diff --git a/backend/cmd/server/wire.go b/backend/cmd/server/wire.go index 2e4b0d902c..8bd9ac5b16 100644 --- a/backend/cmd/server/wire.go +++ b/backend/cmd/server/wire.go @@ -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() } diff --git a/backend/cmd/server/wire_gen.go b/backend/cmd/server/wire_gen.go index 2369e09049..eb4fdfecf8 100644 --- a/backend/cmd/server/wire_gen.go +++ b/backend/cmd/server/wire_gen.go @@ -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() diff --git a/backend/cmd/server/wire_gen_test.go b/backend/cmd/server/wire_gen_test.go index 805ba5353a..dff3a9129e 100644 --- a/backend/cmd/server/wire_gen_test.go +++ b/backend/cmd/server/wire_gen_test.go @@ -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 diff --git a/backend/internal/handler/admin/setting_handler.go b/backend/internal/handler/admin/setting_handler.go index a9d8b0d88e..da21a0f734 100644 --- a/backend/internal/handler/admin/setting_handler.go +++ b/backend/internal/handler/admin/setting_handler.go @@ -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, diff --git a/backend/internal/handler/admin/setting_handler_update.go b/backend/internal/handler/admin/setting_handler_update.go index 5619641c53..99eb1f21b6 100644 --- a/backend/internal/handler/admin/setting_handler_update.go +++ b/backend/internal/handler/admin/setting_handler_update.go @@ -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, diff --git a/backend/internal/handler/dto/settings.go b/backend/internal/handler/dto/settings.go index 6eccb92c8e..e849f74940 100644 --- a/backend/internal/handler/dto/settings.go +++ b/backend/internal/handler/dto/settings.go @@ -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"` diff --git a/backend/internal/handler/handler.go b/backend/internal/handler/handler.go index 9334ffe26a..359aff0071 100644 --- a/backend/internal/handler/handler.go +++ b/backend/internal/handler/handler.go @@ -53,6 +53,7 @@ type Handlers struct { Subscription *SubscriptionHandler Announcement *AnnouncementHandler ChannelMonitor *ChannelMonitorUserHandler + ChannelMonitorV2 *ChannelMonitorV2Handler Admin *AdminHandlers Gateway *GatewayHandler OpenAIGateway *OpenAIGatewayHandler diff --git a/backend/internal/handler/setting_handler.go b/backend/internal/handler/setting_handler.go index 991d9beed1..3886f82c2a 100644 --- a/backend/internal/handler/setting_handler.go +++ b/backend/internal/handler/setting_handler.go @@ -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, diff --git a/backend/internal/handler/wire.go b/backend/internal/handler/wire.go index 73b36dc0cc..36fdb85280 100644 --- a/backend/internal/handler/wire.go +++ b/backend/internal/handler/wire.go @@ -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, diff --git a/backend/internal/repository/wire.go b/backend/internal/repository/wire.go index 84cecf59e4..34c1a1b9d2 100644 --- a/backend/internal/repository/wire.go +++ b/backend/internal/repository/wire.go @@ -98,6 +98,7 @@ var ProviderSet = wire.NewSet( NewTLSFingerprintProfileRepository, NewChannelRepository, NewChannelMonitorRepository, + NewChannelMonitorV2Repository, NewChannelMonitorRequestTemplateRepository, NewContentModerationRepository, NewAffiliateRepository, diff --git a/backend/internal/server/routes/admin.go b/backend/internal/server/routes/admin.go index 161571cd53..b174abbf1f 100644 --- a/backend/internal/server/routes/admin.go +++ b/backend/internal/server/routes/admin.go @@ -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() + } +} diff --git a/backend/internal/server/routes/channel_monitor_feature_gate_test.go b/backend/internal/server/routes/channel_monitor_feature_gate_test.go new file mode 100644 index 0000000000..1881daa4e5 --- /dev/null +++ b/backend/internal/server/routes/channel_monitor_feature_gate_test.go @@ -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) + } + }) + } +} diff --git a/backend/internal/server/routes/user.go b/backend/internal/server/routes/user.go index bcb2a8c917..65ecaa037d 100644 --- a/backend/internal/server/routes/user.go +++ b/backend/internal/server/routes/user.go @@ -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) + } } } diff --git a/backend/internal/service/domain_constants.go b/backend/internal/service/domain_constants.go index c0905e11b9..3db920bf47 100644 --- a/backend/internal/service/domain_constants.go +++ b/backend/internal/service/domain_constants.go @@ -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). diff --git a/backend/internal/service/setting_parse.go b/backend/internal/service/setting_parse.go index 05be77086f..03361da0cf 100644 --- a/backend/internal/service/setting_parse.go +++ b/backend/internal/service/setting_parse.go @@ -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" diff --git a/backend/internal/service/setting_public.go b/backend/internal/service/setting_public.go index ff11fb6897..d3c583ef12 100644 --- a/backend/internal/service/setting_public.go +++ b/backend/internal/service/setting_public.go @@ -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, diff --git a/backend/internal/service/setting_service.go b/backend/internal/service/setting_service.go index 42489ed53b..9319972717 100644 --- a/backend/internal/service/setting_service.go +++ b/backend/internal/service/setting_service.go @@ -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 diff --git a/backend/internal/service/setting_update.go b/backend/internal/service/setting_update.go index 825222fa94..6f38a4d366 100644 --- a/backend/internal/service/setting_update.go +++ b/backend/internal/service/setting_update.go @@ -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 { diff --git a/backend/internal/service/settings_view.go b/backend/internal/service/settings_view.go index 30024ebbe2..28fa75fce6 100644 --- a/backend/internal/service/settings_view.go +++ b/backend/internal/service/settings_view.go @@ -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"` diff --git a/backend/internal/service/wire.go b/backend/internal/service/wire.go index 0e9f61c95a..9b384941ef 100644 --- a/backend/internal/service/wire.go +++ b/backend/internal/service/wire.go @@ -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 +}