mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-03 15:06:36 +08:00
refactor(task): replace init registration with bootstrap wiring
Introduce internal/bootstrap as the composition root with sync.Once guards for task handler registration and push listener wiring. Replace the single OnTaskCompleted global hook with multi-subscriber handlers and remove init()-driven side effects from worker, admin task, and push.
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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{})
|
||||
|
||||
@@ -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{
|
||||
|
||||
Reference in New Issue
Block a user