diff --git a/internal/apps/risk_control/logics.go b/internal/apps/risk_control/logics.go index 7e409d2d..f2bcf684 100644 --- a/internal/apps/risk_control/logics.go +++ b/internal/apps/risk_control/logics.go @@ -21,13 +21,13 @@ const ( ) // InitLogWriter 初始化日志写入通道和后台写入协程 -func InitLogWriter() { +func InitLogWriter(ctx context.Context) { if !config.Config.ClickHouse.Enabled { return } logChan = make(chan *UserAccessLog, defaultQueueSize) - go startBatchWorker() + go startBatchWorker(context.WithoutCancel(ctx)) } // IsBufferFull 检查当前本地缓冲队列是否已满 @@ -52,7 +52,7 @@ func QueueAccessLog(logItem *UserAccessLog) { } } -func startBatchWorker() { +func startBatchWorker(ctx context.Context) { ticker := time.NewTicker(flushInterval) defer ticker.Stop() @@ -67,7 +67,6 @@ func startBatchWorker() { return } - ctx := context.Background() b, err := db.ChConn.PrepareBatch(ctx, "INSERT INTO w_user_access_logs (id, user_id, path, method, ip, user_agent, headers, status, latency, created_at)") if err != nil { logger.ErrorF(ctx, "[RiskControl] Prepare ClickHouse batch failed: %v", err) diff --git a/internal/bootstrap/bootstrap.go b/internal/bootstrap/bootstrap.go new file mode 100644 index 00000000..22523a91 --- /dev/null +++ b/internal/bootstrap/bootstrap.go @@ -0,0 +1,88 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package bootstrap wires cross-module integrations and process-level subsystem initialization. +// All registrations use sync.Once so entry points can call them safely without import-order side effects. +package bootstrap + +import ( + "context" + "sync" + + admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push" + "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" + "github.com/Rain-kl/Wavelet/internal/apps/risk_control" + taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers" + "github.com/Rain-kl/Wavelet/pkg/logger" +) + +// Options selects role-specific runtime bootstrap steps for the current process. +type Options struct { + // API enables HTTP-only subsystems such as the ClickHouse access-log writer. + API bool +} + +var ( + registerTasksOnce sync.Once + registerPushDomainEventsOnce sync.Once + registerTaskListenersOnce sync.Once + initRuntimeOnce sync.Once +) + +// RegisterTasks registers all built-in task handlers and metadata. +func RegisterTasks() { + registerTasksOnce.Do(func() { + taskhandlers.Register() + }) +} + +// RegisterPushDomainEvents wires push notification handlers for domain events. +func RegisterPushDomainEvents() { + registerPushDomainEventsOnce.Do(func() { + custom_events.Register() + }) +} + +// RegisterTaskListeners wires operational listeners to task framework hooks. +func RegisterTaskListeners() { + registerTaskListenersOnce.Do(func() { + admin_push.RegisterTaskListeners() + }) +} + +// RegisterAPI wires integrations required by the HTTP API process. +func RegisterAPI() { + RegisterTasks() + RegisterPushDomainEvents() +} + +// RegisterWorker wires integrations required by the task worker process. +func RegisterWorker() { + RegisterTasks() + RegisterTaskListeners() +} + +// RegisterScheduler wires integrations required by the task scheduler process. +func RegisterScheduler() { + RegisterTasks() +} + +// RegisterAll wires integrations for fused mode (API + Worker + Scheduler). +func RegisterAll() { + RegisterTasks() + RegisterPushDomainEvents() + RegisterTaskListeners() +} + +// Init runs shared runtime bootstrap exactly once per process. +// Call from cmd entry points after wiring registration and database migration, not from router. +func Init(ctx context.Context, opts Options) { + initRuntimeOnce.Do(func() { + if err := admin_push.SyncEvents(ctx); err != nil { + logger.ErrorF(ctx, "[Bootstrap] sync push events failed: %v", err) + } + if opts.API { + risk_control.InitLogWriter(ctx) + } + }) +} \ No newline at end of file diff --git a/internal/bootstrap/bootstrap_test.go b/internal/bootstrap/bootstrap_test.go new file mode 100644 index 00000000..ef678c03 --- /dev/null +++ b/internal/bootstrap/bootstrap_test.go @@ -0,0 +1,34 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package bootstrap + +import ( + "context" + "testing" + + admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/testhelper" +) + +func TestInitSyncsPushEventsOnce(t *testing.T) { + dbConn, _, cleanup := testhelper.SetupTestEnvironment(t) + defer cleanup() + + if err := dbConn.AutoMigrate(&model.PushEvent{}); err != nil { + t.Fatalf("auto migrate push events failed: %v", err) + } + + ctx := context.Background() + Init(ctx, Options{}) + Init(ctx, Options{API: true}) + + var count int64 + if err := dbConn.Model(&model.PushEvent{}).Count(&count).Error; err != nil { + t.Fatalf("count push events failed: %v", err) + } + if count != int64(len(admin_push.BuiltInEvents)) { + t.Fatalf("push event count = %d, want %d", count, len(admin_push.BuiltInEvents)) + } +} \ No newline at end of file diff --git a/internal/cmd/all.go b/internal/cmd/all.go index 10dde12b..b29f7a93 100644 --- a/internal/cmd/all.go +++ b/internal/cmd/all.go @@ -9,6 +9,7 @@ import ( "log" "sync" + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/router" "github.com/Rain-kl/Wavelet/internal/task/scheduler" "github.com/Rain-kl/Wavelet/internal/task/worker" @@ -20,6 +21,8 @@ var allCmd = &cobra.Command{ Short: "以融合模式同时启动 API、Worker 和 Scheduler", Run: func(_ *cobra.Command, _ []string) { log.Println("[All] 融合模式启动") + bootstrap.RegisterAll() + runBootstrap(bootstrap.Options{API: true}) var wg sync.WaitGroup diff --git a/internal/cmd/api.go b/internal/cmd/api.go index e4b745f0..373721af 100644 --- a/internal/cmd/api.go +++ b/internal/cmd/api.go @@ -5,6 +5,7 @@ package cmd import ( + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/router" "github.com/spf13/cobra" ) @@ -13,6 +14,8 @@ var apiCmd = &cobra.Command{ Use: "api", Short: "wavelet API", Run: func(_ *cobra.Command, _ []string) { + bootstrap.RegisterAPI() + runBootstrap(bootstrap.Options{API: true}) router.Serve() }, } diff --git a/internal/cmd/bootstrap.go b/internal/cmd/bootstrap.go new file mode 100644 index 00000000..635cec94 --- /dev/null +++ b/internal/cmd/bootstrap.go @@ -0,0 +1,17 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package cmd + +import ( + "context" + + "github.com/Rain-kl/Wavelet/internal/bootstrap" + "github.com/Rain-kl/Wavelet/pkg/trace" +) + +func runBootstrap(opts bootstrap.Options) { + ctx, span := trace.Start(context.Background(), "bootstrap.Init") + defer span.End() + bootstrap.Init(ctx, opts) +} \ No newline at end of file diff --git a/internal/cmd/scheduler.go b/internal/cmd/scheduler.go index 48795533..f2715a4a 100644 --- a/internal/cmd/scheduler.go +++ b/internal/cmd/scheduler.go @@ -7,6 +7,7 @@ package cmd import ( "log" + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/task/scheduler" "github.com/spf13/cobra" @@ -16,6 +17,7 @@ var schedulerCmd = &cobra.Command{ Use: "scheduler", Short: "wavelet Scheduler", Run: func(_ *cobra.Command, _ []string) { + runBootstrap(bootstrap.Options{}) log.Println("[Scheduler] 启动定时任务调度服务") if err := scheduler.StartScheduler(); err != nil { log.Fatalf("[调度器] 启动失败: %v", err) diff --git a/internal/cmd/worker.go b/internal/cmd/worker.go index 61f9d608..196d17f7 100644 --- a/internal/cmd/worker.go +++ b/internal/cmd/worker.go @@ -7,6 +7,7 @@ package cmd import ( "log" + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/task/worker" "github.com/spf13/cobra" @@ -16,6 +17,7 @@ var workerCmd = &cobra.Command{ Use: "worker", Short: "wavelet Worker", Run: func(_ *cobra.Command, _ []string) { + runBootstrap(bootstrap.Options{}) log.Println("[Worker] 启动任务处理服务") if err := worker.StartWorker(); err != nil { log.Fatalf("[工作器] 启动失败: %v", err) diff --git a/internal/router/router.go b/internal/router/router.go index b059c4a7..c21fecac 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -15,13 +15,10 @@ import ( "syscall" "time" - admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push" "github.com/Rain-kl/Wavelet/internal/apps/risk_control" router_root "github.com/Rain-kl/Wavelet/internal/router/root" v1 "github.com/Rain-kl/Wavelet/internal/router/v1" - // Swagger 文档生成 - _ "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" "github.com/Rain-kl/Wavelet/internal/apps/oauth" "github.com/Rain-kl/Wavelet/internal/config" otel_trace "github.com/Rain-kl/Wavelet/pkg/trace" @@ -38,14 +35,6 @@ func Serve() { gin.SetMode(gin.ReleaseMode) } - // 初始化 ClickHouse 异步日志写入器 - risk_control.InitLogWriter() - - // 运行内置事件同步 - if err := admin_push.SyncEvents(context.Background()); err != nil { - log.Printf("[API] sync push events failed: %v\n", err) - } - // 初始化路由 r := gin.New() r.Use(gin.Recovery())