diff --git a/.agent/skills/clickhouse-batchwriter/SKILL.md b/.agent/skills/clickhouse-batchwriter/SKILL.md index ce67d278..a83c2120 100644 --- a/.agent/skills/clickhouse-batchwriter/SKILL.md +++ b/.agent/skills/clickhouse-batchwriter/SKILL.md @@ -18,7 +18,8 @@ DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运 | Model | `internal/model/analytics/` | 列定义、`TableName()`、`BatchInsertSQL()`(及可选 `InsertColumns()`) | | Repository | `internal/repository/analytics/` | `BatchInsert*` / `BatchInsertNodeAccessLogs` 等;`PrepareBatch` + 多行 `Append` + 一次 `Send` | | Apps | `internal/apps//` | 采集、入队、背压;`FlushFunc` 只调 repository,不写 SQL、不 `PrepareBatch` | -| 装配 | `internal/bootstrap/bootstrap.go` | 进程启动时 `Writer.Start`、停机时 `Writer.Stop`(与 `risk_control.InitLogWriter` 同级) | +| 装配 | `internal/bootstrap/bootstrap.go` | 进程启动时调用 `Writer.Start`;初始化时需调用 `lifecycle.OnShutdown` 挂载停机钩子 | +| 生命周期 | `internal/lifecycle/lifecycle.go` | 统一协调全局并发优雅停机,业务包无需在 `bootstrap.go` 中硬编码 `Stop` 逻辑 | **禁止**在 Handler / middleware 内直接 `db.ChConn.PrepareBatch`;**禁止**在 repository 内启动 goroutine 或维护全局 channel(队列生命周期由 apps + bootstrap 或专用 writer 包负责)。 @@ -73,8 +74,8 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush - `db.ChConn == nil` 返回明确错误 - 一次 `PrepareBatch` → 循环 `Append` → 一次 `Send` 4. **Writer 胶水**(`internal/apps//` 或 `internal/repository/analytics/_writer.go`): - - `New` + `Start` 在 bootstrap 注册 - - 业务路径 `TryEnqueue`;HTTP 背压用 `IsFull()` + - `New` + `Start`,并在初始化逻辑内通过 `lifecycle.OnShutdown("your_writer_name", Stop)` 注册停机回调 + - 业务路径 `TryEnqueue`;HTTP 背压用 `IsFull()` 5. **测试**: - repository:mock `ChConn` 验证 `BatchInsertSQL` 与 append 列数 - batchwriter:`go test ./internal/db/batchwriter` @@ -134,7 +135,7 @@ func RegisterAPI(ctx context.Context) { ``` - `RegisterAPI` / `RegisterAll`:`Start` -- 进程优雅停机:带超时的 `Stop(ctx)` +- 进程优雅停机:业务模块在初始化时调用 `lifecycle.OnShutdown` 注册,由 `bootstrap.Stop()` 代理 `lifecycle.Stop()` 并发停机。 - 使用 `sync.Once` 保证幂等 ## 验证清单 @@ -158,4 +159,5 @@ make code-check - OpenFlare 写入胶水:`internal/apps/openflare/chwriter/writer.go` - 节点访问日志 repository:`internal/repository/analytics/node_access_log_writer.go` - 可观测 repository:`internal/repository/analytics/node_observability_writer.go` +- 生命周期管理器:`internal/lifecycle/lifecycle.go` - Bootstrap:`internal/bootstrap/bootstrap.go` \ No newline at end of file diff --git a/internal/apps/openflare/chwriter/writer.go b/internal/apps/openflare/chwriter/writer.go index eee83561..8e8f0a03 100644 --- a/internal/apps/openflare/chwriter/writer.go +++ b/internal/apps/openflare/chwriter/writer.go @@ -13,6 +13,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db/batchwriter" + "github.com/Rain-kl/Wavelet/internal/lifecycle" analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/Rain-kl/Wavelet/pkg/logger" @@ -65,6 +66,8 @@ func Init(ctx context.Context) { frpsWriter.Start(ctx) frpcWriter.Start(ctx) nodeAccessLogWriter.Start(ctx) + + lifecycle.OnShutdown("openflare_chwriter", Stop) }) } diff --git a/internal/apps/risk_control/logics.go b/internal/apps/risk_control/logics.go index 23379421..8cba0ade 100644 --- a/internal/apps/risk_control/logics.go +++ b/internal/apps/risk_control/logics.go @@ -9,6 +9,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db/batchwriter" + "github.com/Rain-kl/Wavelet/internal/lifecycle" "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/Rain-kl/Wavelet/pkg/logger" @@ -60,6 +61,7 @@ func InitLogWriter(ctx context.Context) { writer.Start(ctx) logWriter = writer + lifecycle.OnShutdown("risk_control_log_writer", StopLogWriter) } // StopLogWriter stops the ClickHouse access-log batch writer and drains pending logs. diff --git a/internal/bootstrap/bootstrap.go b/internal/bootstrap/bootstrap.go index d3e3cbde..a18a9c55 100644 --- a/internal/bootstrap/bootstrap.go +++ b/internal/bootstrap/bootstrap.go @@ -13,6 +13,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" "github.com/Rain-kl/Wavelet/internal/apps/risk_control" + "github.com/Rain-kl/Wavelet/internal/lifecycle" taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers" "github.com/Rain-kl/Wavelet/pkg/logger" ) @@ -91,12 +92,7 @@ func Init(ctx context.Context, opts Options) { // Stop stops all batch writers and background resources. func Stop(ctx context.Context) { - if err := risk_control.StopLogWriter(ctx); err != nil { - logger.ErrorF(ctx, "[Bootstrap] stop risk_control log writer failed: %v", err) - } - if err := chwriter.Stop(ctx); err != nil { - logger.ErrorF(ctx, "[Bootstrap] stop chwriter failed: %v", err) - } + lifecycle.Stop(ctx) } // ResetInitRuntimeOnceForTest clears initRuntimeOnce so Init can run again in unit tests. diff --git a/internal/lifecycle/lifecycle.go b/internal/lifecycle/lifecycle.go new file mode 100644 index 00000000..edfe2d60 --- /dev/null +++ b/internal/lifecycle/lifecycle.go @@ -0,0 +1,66 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package lifecycle manages global application and business component shutdown hooks. +package lifecycle + +import ( + "context" + "log" + "sync" +) + +// ShutdownFunc defines the signature for a graceful shutdown callback. +type ShutdownFunc func(ctx context.Context) error + +type hook struct { + name string + fn ShutdownFunc +} + +var ( + hooks []hook + mu sync.Mutex +) + +// OnShutdown registers a callback to be run during graceful shutdown. +func OnShutdown(name string, fn ShutdownFunc) { + mu.Lock() + defer mu.Unlock() + hooks = append(hooks, hook{name: name, fn: fn}) +} + +// Stop executes all registered shutdown hooks concurrently and waits for completion or context timeout. +func Stop(ctx context.Context) { + mu.Lock() + localHooks := make([]hook, len(hooks)) + copy(localHooks, hooks) + mu.Unlock() + + var wg sync.WaitGroup + for _, h := range localHooks { + wg.Add(1) + go func(name string, fn ShutdownFunc) { + defer wg.Done() + log.Printf("[Lifecycle] stopping %s...\n", name) + if err := fn(ctx); err != nil { + log.Printf("[Lifecycle] stop %s failed: %v\n", name, err) + } else { + log.Printf("[Lifecycle] %s stopped successfully\n", name) + } + }(h.name, h.fn) + } + + done := make(chan struct{}) + go func() { + wg.Wait() + close(done) + }() + + select { + case <-done: + log.Println("[Lifecycle] all services stopped gracefully") + case <-ctx.Done(): + log.Printf("[Lifecycle] shutdown timed out: %v\n", ctx.Err()) + } +}