mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-07 16:16:37 +08:00
refactor(bootstrap): move runtime init from router to cmd layer
Extract SyncEvents and InitLogWriter from router.Serve into bootstrap.Init called from cmd entry points with trace-aware context. Preserve existing Register* wiring for task and push domain integrations.
This commit is contained in:
@@ -21,13 +21,13 @@ const (
|
|||||||
)
|
)
|
||||||
|
|
||||||
// InitLogWriter 初始化日志写入通道和后台写入协程
|
// InitLogWriter 初始化日志写入通道和后台写入协程
|
||||||
func InitLogWriter() {
|
func InitLogWriter(ctx context.Context) {
|
||||||
if !config.Config.ClickHouse.Enabled {
|
if !config.Config.ClickHouse.Enabled {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
logChan = make(chan *UserAccessLog, defaultQueueSize)
|
logChan = make(chan *UserAccessLog, defaultQueueSize)
|
||||||
go startBatchWorker()
|
go startBatchWorker(context.WithoutCancel(ctx))
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsBufferFull 检查当前本地缓冲队列是否已满
|
// IsBufferFull 检查当前本地缓冲队列是否已满
|
||||||
@@ -52,7 +52,7 @@ func QueueAccessLog(logItem *UserAccessLog) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func startBatchWorker() {
|
func startBatchWorker(ctx context.Context) {
|
||||||
ticker := time.NewTicker(flushInterval)
|
ticker := time.NewTicker(flushInterval)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
@@ -67,7 +67,6 @@ func startBatchWorker() {
|
|||||||
return
|
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)")
|
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 {
|
if err != nil {
|
||||||
logger.ErrorF(ctx, "[RiskControl] Prepare ClickHouse batch failed: %v", err)
|
logger.ErrorF(ctx, "[RiskControl] Prepare ClickHouse batch failed: %v", err)
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -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))
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/bootstrap"
|
||||||
"github.com/Rain-kl/Wavelet/internal/router"
|
"github.com/Rain-kl/Wavelet/internal/router"
|
||||||
"github.com/Rain-kl/Wavelet/internal/task/scheduler"
|
"github.com/Rain-kl/Wavelet/internal/task/scheduler"
|
||||||
"github.com/Rain-kl/Wavelet/internal/task/worker"
|
"github.com/Rain-kl/Wavelet/internal/task/worker"
|
||||||
@@ -20,6 +21,8 @@ var allCmd = &cobra.Command{
|
|||||||
Short: "以融合模式同时启动 API、Worker 和 Scheduler",
|
Short: "以融合模式同时启动 API、Worker 和 Scheduler",
|
||||||
Run: func(_ *cobra.Command, _ []string) {
|
Run: func(_ *cobra.Command, _ []string) {
|
||||||
log.Println("[All] 融合模式启动")
|
log.Println("[All] 融合模式启动")
|
||||||
|
bootstrap.RegisterAll()
|
||||||
|
runBootstrap(bootstrap.Options{API: true})
|
||||||
|
|
||||||
var wg sync.WaitGroup
|
var wg sync.WaitGroup
|
||||||
|
|
||||||
|
|||||||
@@ -5,6 +5,7 @@
|
|||||||
package cmd
|
package cmd
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/bootstrap"
|
||||||
"github.com/Rain-kl/Wavelet/internal/router"
|
"github.com/Rain-kl/Wavelet/internal/router"
|
||||||
"github.com/spf13/cobra"
|
"github.com/spf13/cobra"
|
||||||
)
|
)
|
||||||
@@ -13,6 +14,8 @@ var apiCmd = &cobra.Command{
|
|||||||
Use: "api",
|
Use: "api",
|
||||||
Short: "wavelet API",
|
Short: "wavelet API",
|
||||||
Run: func(_ *cobra.Command, _ []string) {
|
Run: func(_ *cobra.Command, _ []string) {
|
||||||
|
bootstrap.RegisterAPI()
|
||||||
|
runBootstrap(bootstrap.Options{API: true})
|
||||||
router.Serve()
|
router.Serve()
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -7,6 +7,7 @@ package cmd
|
|||||||
import (
|
import (
|
||||||
"log"
|
"log"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/bootstrap"
|
||||||
"github.com/Rain-kl/Wavelet/internal/task/scheduler"
|
"github.com/Rain-kl/Wavelet/internal/task/scheduler"
|
||||||
|
|
||||||
"github.com/spf13/cobra"
|
"github.com/spf13/cobra"
|
||||||
@@ -16,6 +17,7 @@ var schedulerCmd = &cobra.Command{
|
|||||||
Use: "scheduler",
|
Use: "scheduler",
|
||||||
Short: "wavelet Scheduler",
|
Short: "wavelet Scheduler",
|
||||||
Run: func(_ *cobra.Command, _ []string) {
|
Run: func(_ *cobra.Command, _ []string) {
|
||||||
|
runBootstrap(bootstrap.Options{})
|
||||||
log.Println("[Scheduler] 启动定时任务调度服务")
|
log.Println("[Scheduler] 启动定时任务调度服务")
|
||||||
if err := scheduler.StartScheduler(); err != nil {
|
if err := scheduler.StartScheduler(); err != nil {
|
||||||
log.Fatalf("[调度器] 启动失败: %v", err)
|
log.Fatalf("[调度器] 启动失败: %v", err)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ package cmd
|
|||||||
import (
|
import (
|
||||||
"log"
|
"log"
|
||||||
|
|
||||||
|
"github.com/Rain-kl/Wavelet/internal/bootstrap"
|
||||||
"github.com/Rain-kl/Wavelet/internal/task/worker"
|
"github.com/Rain-kl/Wavelet/internal/task/worker"
|
||||||
|
|
||||||
"github.com/spf13/cobra"
|
"github.com/spf13/cobra"
|
||||||
@@ -16,6 +17,7 @@ var workerCmd = &cobra.Command{
|
|||||||
Use: "worker",
|
Use: "worker",
|
||||||
Short: "wavelet Worker",
|
Short: "wavelet Worker",
|
||||||
Run: func(_ *cobra.Command, _ []string) {
|
Run: func(_ *cobra.Command, _ []string) {
|
||||||
|
runBootstrap(bootstrap.Options{})
|
||||||
log.Println("[Worker] 启动任务处理服务")
|
log.Println("[Worker] 启动任务处理服务")
|
||||||
if err := worker.StartWorker(); err != nil {
|
if err := worker.StartWorker(); err != nil {
|
||||||
log.Fatalf("[工作器] 启动失败: %v", err)
|
log.Fatalf("[工作器] 启动失败: %v", err)
|
||||||
|
|||||||
@@ -15,13 +15,10 @@ import (
|
|||||||
"syscall"
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push"
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
||||||
router_root "github.com/Rain-kl/Wavelet/internal/router/root"
|
router_root "github.com/Rain-kl/Wavelet/internal/router/root"
|
||||||
v1 "github.com/Rain-kl/Wavelet/internal/router/v1"
|
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/apps/oauth"
|
||||||
"github.com/Rain-kl/Wavelet/internal/config"
|
"github.com/Rain-kl/Wavelet/internal/config"
|
||||||
otel_trace "github.com/Rain-kl/Wavelet/pkg/trace"
|
otel_trace "github.com/Rain-kl/Wavelet/pkg/trace"
|
||||||
@@ -38,14 +35,6 @@ func Serve() {
|
|||||||
gin.SetMode(gin.ReleaseMode)
|
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 := gin.New()
|
||||||
r.Use(gin.Recovery())
|
r.Use(gin.Recovery())
|
||||||
|
|||||||
Reference in New Issue
Block a user