diff --git a/internal/apps/admin/push/task_listener.go b/internal/apps/admin/push/task_listener.go index 39fbcd69..5dafd187 100644 --- a/internal/apps/admin/push/task_listener.go +++ b/internal/apps/admin/push/task_listener.go @@ -15,8 +15,9 @@ import ( "github.com/Rain-kl/Wavelet/pkg/logger" ) -func init() { - task.OnTaskCompleted = handleTaskCompleted +// RegisterTaskListeners subscribes push notification handlers to task completion events. +func RegisterTaskListeners() { + task.OnTaskCompleted(handleTaskCompleted) } // handleTaskCompleted handles task completions and triggers appropriate push events. diff --git a/internal/apps/admin/task/routers.go b/internal/apps/admin/task/routers.go index 580a2261..974fd762 100644 --- a/internal/apps/admin/task/routers.go +++ b/internal/apps/admin/task/routers.go @@ -13,7 +13,6 @@ import ("fmt" "github.com/Rain-kl/Wavelet/internal/apps/admin" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/task" - taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers" "github.com/Rain-kl/Wavelet/internal/task/scheduler" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" @@ -21,10 +20,6 @@ import ("fmt" "github.com/Rain-kl/Wavelet/internal/common/response") -func init() { - taskhandlers.Register() -} - // ListTaskTypes 获取支持的任务类型列表 // @Summary 获取支持的任务类型 // @Description 返回系统支持的所有可调度任务类型列表,包括任务名称、描述、是否支持时间范围等元数据,需要管理员权限 diff --git a/internal/bootstrap/bootstrap.go b/internal/bootstrap/bootstrap.go new file mode 100644 index 00000000..3471ea19 --- /dev/null +++ b/internal/bootstrap/bootstrap.go @@ -0,0 +1,65 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package bootstrap wires cross-module integrations at the application composition root. +// All registrations use sync.Once so entry points can call them safely without import-order side effects. +package bootstrap + +import ( + "sync" + + admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push" + "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" + taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers" +) + +var ( + registerTasksOnce sync.Once + registerPushDomainEventsOnce sync.Once + registerTaskListenersOnce 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() +} \ No newline at end of file diff --git a/internal/cmd/all.go b/internal/cmd/all.go index 10dde12b..dabebbd6 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,7 @@ var allCmd = &cobra.Command{ Short: "以融合模式同时启动 API、Worker 和 Scheduler", Run: func(_ *cobra.Command, _ []string) { log.Println("[All] 融合模式启动") + bootstrap.RegisterAll() var wg sync.WaitGroup diff --git a/internal/router/router.go b/internal/router/router.go index 415de15f..ec7849a8 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -16,8 +16,8 @@ import ( "time" 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" + "github.com/Rain-kl/Wavelet/internal/bootstrap" router_root "github.com/Rain-kl/Wavelet/internal/router/root" v1 "github.com/Rain-kl/Wavelet/internal/router/v1" @@ -40,8 +40,8 @@ func Serve() { // 初始化 ClickHouse 异步日志写入器 risk_control.InitLogWriter() - // 装配推送模块与领域事件的集成(组合根显式注册,避免 init 副作用) - custom_events.Register() + // 组合根显式装配跨模块集成(避免 init 副作用与 import 顺序依赖) + bootstrap.RegisterAPI() // 运行内置事件同步 if err := admin_push.SyncEvents(context.Background()); err != nil { diff --git a/internal/task/executor.go b/internal/task/executor.go index 74e23117..e41ae65e 100644 --- a/internal/task/executor.go +++ b/internal/task/executor.go @@ -26,8 +26,16 @@ import ( // handlerRegistry 已注册的任务处理器 var handlerRegistry = make(map[string]TaskHandler) -// OnTaskCompleted is a hook called when a task execution completes. -var OnTaskCompleted func(ctx context.Context, execution *model.TaskExecution, result *TaskResult, execErr error) +// CompletedHandler is called when a task execution completes. +type CompletedHandler func(ctx context.Context, execution *model.TaskExecution, result *TaskResult, execErr error) + +var taskCompletedHandlers []CompletedHandler + +// OnTaskCompleted registers a handler for task completion events. +// Handlers must be registered during application bootstrap before processing tasks. +func OnTaskCompleted(handler CompletedHandler) { + taskCompletedHandlers = append(taskCompletedHandlers, handler) +} // RegisterHandler 注册任务处理器 // 传入任务类型标识(对应 constants.go 中的 AsynqTask 常量)和 TaskHandler 实现 @@ -405,9 +413,17 @@ func completeTaskExecution(ctx context.Context, execution *model.TaskExecution, } } - if OnTaskCompleted != nil { - asyncCtx := context.WithoutCancel(ctx) - go OnTaskCompleted(asyncCtx, execution, result, execErr) + notifyTaskCompleted(ctx, execution, result, execErr) +} + +func notifyTaskCompleted(ctx context.Context, execution *model.TaskExecution, result *TaskResult, execErr error) { + if len(taskCompletedHandlers) == 0 { + return + } + + asyncCtx := context.WithoutCancel(ctx) + for _, handler := range taskCompletedHandlers { + go handler(asyncCtx, execution, result, execErr) } } diff --git a/internal/task/scheduler/scheduler.go b/internal/task/scheduler/scheduler.go index 2d1a1718..478a0511 100644 --- a/internal/task/scheduler/scheduler.go +++ b/internal/task/scheduler/scheduler.go @@ -11,6 +11,7 @@ import ( "syscall" "time" + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/task" "github.com/Rain-kl/Wavelet/pkg/logger" @@ -32,6 +33,8 @@ func GetAsynqClient() *asynq.Client { // StartScheduler 启动调度器 (该函数阻塞,直到调度器退出) func StartScheduler() error { + bootstrap.RegisterScheduler() + var err error schedulerOnce.Do(func() { quitChan = make(chan struct{}) diff --git a/internal/task/worker/worker.go b/internal/task/worker/worker.go index 72e2668d..b66d4f65 100644 --- a/internal/task/worker/worker.go +++ b/internal/task/worker/worker.go @@ -7,22 +7,18 @@ package worker import ( "time" + "github.com/Rain-kl/Wavelet/internal/bootstrap" "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/task" - taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers" "github.com/hibiken/asynq" ) // workerShutdownTimeout Worker 优雅关闭超时时间 const workerShutdownTimeout = 3 * time.Minute -func init() { - // 注册所有任务处理器 - taskhandlers.Register() -} - // StartWorker 启动任务处理服务器 func StartWorker() error { + bootstrap.RegisterWorker() asynqServer := asynq.NewServer( task.RedisOpt, asynq.Config{