diff --git a/.auto/ideas.md b/.auto/ideas.md index c0da48f5..ab3102cf 100644 --- a/.auto/ideas.md +++ b/.auto/ideas.md @@ -73,7 +73,7 @@ - 已修:analytics/node_access_log_filter.go、analytics/access_log_filter.go、 logstore/postgres_store.go×2(PG/SQLite 加 ESCAPE '\',CH 用默认反斜杠转义)。 新助手 pkg/util/like.go EscapeLike + 单测。 -- 遗留同类(低风险,用户名/关键词搜索):repository/upload.go:53 keyword contains、 - repository/user.go:73/76/188 username/email 前缀+contains、task_execution.go:155 - task_type 前缀。user.go:228 `base+"-%"` 与 config_version.go:65 为系统生成模式, - 刻意通配勿动。GORM 站点加 ESCAPE 子句即可复用 EscapeLike。 +- Run #47 已收尾全部 GORM 站点:upload.go keyword、user.go:73/76/188/229 + (含 OAuth uniqueUsername base 转义——外部输入含 _ 曾误报用户名冲突)、 + task_execution.go task_type 前缀。均加显式 ESCAPE '\'。 +- 刻意保留:upload.go:199 `image/%`(系统常量)、config_version.go:65(系统生成)。 diff --git a/.auto/log.jsonl b/.auto/log.jsonl index 224fb429..0ece6a69 100644 --- a/.auto/log.jsonl +++ b/.auto/log.jsonl @@ -45,3 +45,4 @@ {"run":44,"commit":"63007fc","metric":8,"metrics":{"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":92,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"全仓 race 扫描发现 upload/cache 监听器 DATA RACE:捕获 redis 客户端消除全局读竞争 + Stop 等待 done + 同型监听器(oauth×2/repository×2)加固","timestamp":1787711092906,"segment":0,"confidence":null,"asi":{"hypothesis":"全仓 go test -race 可能暴露并发 bug(此前仅局部验证)","next_action_hint":"继续扫其他模块;可考虑把 -race 纳入周期性检查","result":"发现并修复 1 个真实 DATA RACE;修复后全仓 -race 0 竞争,metric 持平 8","root_cause":"upload/cache 监听器 goroutine 读可变全局 db.Redis,与 testhelper 清理置 nil 竞争;testhelper 导入 upload/cache 有循环依赖,故用启动时捕获客户端的根因修复(oauth/repository 同型监听器一并加固),并补 StopUploadMetaCacheListener 同步等待 done"}} {"run":45,"commit":"63007fc","metric":8,"metrics":{"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":70,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"discard","description":"探索轮:索引对齐/前端请求瀑布/BasicAuth 注入面三假设均证伪,无代码变更","timestamp":1787711474404,"segment":0,"confidence":null,"asi":{"hypothesis":"SQLite 迁移缺 PG 同款索引;前端存在串行请求瀑布;nginx BasicAuth 密码有注入面","next_action_hint":"代码库已高度收敛;下轮可考虑 observability 查询构造器审计或周期性重跑 -race","rollback_reason":"纯探索无代码变更,无需回滚","result":"三个假设均无产出:①索引对比(修正提取正则后)PG/SQLite 完全对齐,SQLite 仅多 legacy w_* 冗余索引;②前端 await Service 均在事件处理器非渲染期;③BasicAuth 密码经 base64 编码(字母表无元字符)无注入面","lessons":"grep 提取 SQL 时注意 IF NOT EXISTS 变体,否则产生假缺口"}} {"run":46,"commit":"2cb3392","metric":8,"metrics":{"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":106,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"LIKE 过滤器转义修复:日志搜索含 %/_ 的输入不再被当通配符;pkg/util 新增 EscapeLike 共享助手 + 单测","timestamp":1787712116152,"segment":0,"confidence":null,"asi":{"hypothesis":"日志搜索 LIKE 过滤器不转义 %/_/\\,含下划线的路径/主机名搜索结果错误","next_action_hint":"同类遗留站点(upload/user/task_execution GORM 搜索)已记 ideas.md,可作后续轮次","result":"修复 4 个站点:analytics 两处 CH 过滤器 + logstore postgres_store 两处(PG/SQLite 加 ESCAPE '\\')。新增 pkg/util/like.go EscapeLike + 单测。metric 持平 8,全部测试通过","scope_decision":"GORM 实体搜索站(upload keyword、user username/email)同 bug 类但低风险且可能依赖现有通配语义,本轮不动"}} +{"run":47,"commit":"3528323","metric":8,"metrics":{"eslint_errors":0,"eslint_problems":0,"eslint_warnings":0,"golint_canonicalheader":0,"golint_errname":0,"golint_errorlint":1,"golint_exhaustive":0,"golint_forcetypeassert":0,"golint_gosec":0,"golint_intrange":0,"golint_modernize":3,"golint_nilnil":3,"golint_perfsprint":0,"golint_prealloc":0,"golint_recvcheck":1,"golint_test_testifylint":0,"golint_test_thelper":0,"golint_test_total":0,"golint_test_usetesting":0,"golint_total":8,"golint_usestdlibvars":0,"golint_vetx_total":0,"golint_wastedassign":0,"measure_s":102,"tsc_errors":0,"vitest_failed":0,"vitest_total":126},"status":"keep","description":"GORM 实体搜索 LIKE 转义收尾:6 站点复用 EscapeLike + 显式 ESCAPE 子句,含 OAuth 用户名冲突误报修复","timestamp":1787712555794,"segment":0,"confidence":null,"asi":{"hypothesis":"GORM 实体搜索站与 #46 日志搜索同 bug 类:LIKE 模式不转义通配符","next_action_hint":"LIKE 类已全部收尾;下轮可考虑 ideas.md 的测试可运行性方向或周期性全仓 -race 重跑","result":"6 站点修复(upload keyword、user username/email 前缀+contains、OAuth uniqueUsername base、task_type 前缀),PG/SQLite 加显式 ESCAPE。系统常量模式刻意保留(upload.go:199 image/%)。metric 持平 8,测试全绿","scope_decision":"uniqueUsername 的 base 来自 OAuth 用户信息属外部输入,含 _ 会误报用户名冲突——虽是系统生成后缀模式也需转义 base 本身"}} diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go index df5209b1..6a6b4731 100644 --- a/internal/apps/admin/logs/routers.go +++ b/internal/apps/admin/logs/routers.go @@ -18,6 +18,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/logger" + "github.com/Rain-kl/Wavelet/pkg/util" "github.com/gin-gonic/gin" "github.com/Rain-kl/Wavelet/internal/shared/response" @@ -106,7 +107,7 @@ func HandleLogWebSocket(c *gin.Context) { // 在独立 goroutine 中读取客户端消息(保持连接活跃 + 检测断开) done := make(chan struct{}) - go func() { + util.Go(func() { defer close(done) for { _, _, err := conn.ReadMessage() @@ -114,7 +115,7 @@ func HandleLogWebSocket(c *gin.Context) { return } } - }() + }) // 主循环:推送日志 for { diff --git a/internal/apps/admin/push/events.go b/internal/apps/admin/push/events.go index ca5687c2..66f47f19 100644 --- a/internal/apps/admin/push/events.go +++ b/internal/apps/admin/push/events.go @@ -17,6 +17,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/pkg/logger" pkgpush "github.com/Rain-kl/Wavelet/pkg/push" + "github.com/Rain-kl/Wavelet/pkg/util" "gorm.io/gorm" ) @@ -75,7 +76,7 @@ var DefaultTrigger = &EventTrigger{} //nolint:contextcheck func (t *EventTrigger) Trigger(ctx context.Context, meta EventMetadata, body map[string]any) { asyncCtx := context.WithoutCancel(ctx) - go func() { + util.Go(func() { if body == nil { body = make(map[string]any) } @@ -99,7 +100,7 @@ func (t *EventTrigger) Trigger(ctx context.Context, meta EventMetadata, body map flatBody := getFlatBody(body) msg, _ := t.buildMessage(&event, meta, flatBody, body) t.enqueuePushTasks(asyncCtx, meta, &event, msg, flatBody) - }() + }) } func (t *EventTrigger) buildMessage(event *model.PushEvent, meta EventMetadata, flatBody map[string]any, body map[string]any) (NotificationMessage, string) { diff --git a/internal/apps/admin/updater/routers.go b/internal/apps/admin/updater/routers.go index 2439a10f..91cd999e 100644 --- a/internal/apps/admin/updater/routers.go +++ b/internal/apps/admin/updater/routers.go @@ -9,6 +9,7 @@ import ( "time" "github.com/Rain-kl/Wavelet/pkg/logger" + "github.com/Rain-kl/Wavelet/pkg/util" "github.com/gin-gonic/gin" "github.com/Rain-kl/Wavelet/internal/shared/response" @@ -58,11 +59,11 @@ func ApplyUpdate(c *gin.Context) { logger.InfoF(c.Request.Context(), "[Updater] upgrade prepared; restarting with %s", stagedBinary) c.JSON(http.StatusOK, response.OKNil()) - go func() { + util.Go(func() { time.Sleep(time.Second) if err := replaceAndRestart(executable, stagedBinary); err != nil { defaultManager.finishUpgrade() logger.ErrorF(context.Background(), "[Updater] replace and restart failed: %v", err) } - }() + }) } diff --git a/internal/apps/agent/agent/runner.go b/internal/apps/agent/agent/runner.go index 8ef4f3ce..b2141483 100644 --- a/internal/apps/agent/agent/runner.go +++ b/internal/apps/agent/agent/runner.go @@ -18,6 +18,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/agent/state" "github.com/Rain-kl/Wavelet/internal/apps/agent/wsclient" edgeheartbeat "github.com/Rain-kl/Wavelet/internal/apps/edge/heartbeat" + "github.com/Rain-kl/Wavelet/pkg/util" ) // HeartbeatService handles node registration and periodic heartbeat reporting. @@ -200,12 +201,12 @@ func (r *Runner) startWebSocket(ctx context.Context, nodeID string) (<-chan erro return nil, err } done := make(chan error, 1) - go func() { + util.Go(func() { defer func() { _ = conn.Close() }() done <- r.runWebSocket(ctx, nodeID, conn) - }() + }) return done, nil } @@ -251,7 +252,7 @@ func (r *Runner) runWebSocket(ctx context.Context, nodeID string, conn protocol. childCtx, cancel := context.WithCancel(ctx) defer cancel() - go func() { + util.Go(func() { for { select { case <-childCtx.Done(): @@ -264,7 +265,7 @@ func (r *Runner) runWebSocket(ctx context.Context, nodeID string, conn protocol. } } } - }() + }) wsConn, ok := conn.(*wsclient.Connection) if !ok { diff --git a/internal/apps/cap/runtime_settings.go b/internal/apps/cap/runtime_settings.go index 76399342..3cd7b2ea 100644 --- a/internal/apps/cap/runtime_settings.go +++ b/internal/apps/cap/runtime_settings.go @@ -17,6 +17,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -190,7 +191,7 @@ func startRuntimeSettingsInvalidationListener() { return } - go func() { + util.Go(func() { pubsub := db.Redis.Subscribe(context.Background(), repository.SystemConfigInvalidationChannel) defer func() { _ = pubsub.Close() @@ -208,5 +209,5 @@ func startRuntimeSettingsInvalidationListener() { InvalidateRuntimeSettings() } } - }() + }) } diff --git a/internal/apps/flared/frpc/manager.go b/internal/apps/flared/frpc/manager.go index 43af7544..1dfa0302 100644 --- a/internal/apps/flared/frpc/manager.go +++ b/internal/apps/flared/frpc/manager.go @@ -22,6 +22,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/flared/config" service "github.com/Rain-kl/Wavelet/pkg/protocol" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -180,7 +181,7 @@ func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath } m.processes[relayID] = proc - go func() { + util.Go(func() { backoff := 1 * time.Second const maxBackoff = 60 * time.Second @@ -262,7 +263,7 @@ func (m *Manager) restartProcess(ctx context.Context, relayID string, configPath } t.Stop() } - }() + }) } // Stop cancels and removes all managed frpc processes. diff --git a/internal/apps/oauth/cache.go b/internal/apps/oauth/cache.go index 696f43ba..b8d7d348 100644 --- a/internal/apps/oauth/cache.go +++ b/internal/apps/oauth/cache.go @@ -14,6 +14,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/pkg/cache/ram" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -60,7 +61,7 @@ func startTokenCacheInvalidationListener() { tokenListenerDone = make(chan struct{}) redisClient := db.Redis // 捕获当前客户端:goroutine 不读可变全局,避免与测试置空 db.Redis 竞争 - go func() { + util.Go(func() { listenerCtx := tokenListenerCtx defer close(tokenListenerDone) @@ -69,10 +70,10 @@ func startTokenCacheInvalidationListener() { _ = pubsub.Close() }() - go func() { + util.Go(func() { <-listenerCtx.Done() _ = pubsub.Close() - }() + }) for msg := range pubsub.Channel() { tokenHash := msg.Payload @@ -82,7 +83,7 @@ func startTokenCacheInvalidationListener() { tokenRAM.Invalidate(tokenHash) } } - }() + }) } func publishTokenRAMInvalidation(ctx context.Context, tokenHash string) { @@ -104,7 +105,7 @@ func startUserCacheInvalidationListener() { userListenerDone = make(chan struct{}) redisClient := db.Redis // 捕获当前客户端:goroutine 不读可变全局,避免与测试置空 db.Redis 竞争 - go func() { + util.Go(func() { listenerCtx := userListenerCtx defer close(userListenerDone) @@ -113,10 +114,10 @@ func startUserCacheInvalidationListener() { _ = pubsub.Close() }() - go func() { + util.Go(func() { <-listenerCtx.Done() _ = pubsub.Close() - }() + }) for msg := range pubsub.Channel() { userIDStr := msg.Payload @@ -126,7 +127,7 @@ func startUserCacheInvalidationListener() { userRAM.Invalidate(userID) } } - }() + }) } func publishUserRAMInvalidation(ctx context.Context, userID uint64) { diff --git a/internal/apps/openflare/pages/source_sync.go b/internal/apps/openflare/pages/source_sync.go index 117ebaf2..a7f93122 100644 --- a/internal/apps/openflare/pages/source_sync.go +++ b/internal/apps/openflare/pages/source_sync.go @@ -20,6 +20,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/pkg/logger" + "github.com/Rain-kl/Wavelet/pkg/util" "gorm.io/gorm" "gorm.io/gorm/clause" ) @@ -91,7 +92,7 @@ func startSourceLeaseHeartbeat( workCtx, cancel := context.WithCancel(ctx) done := make(chan error, 1) heartbeat := &sourceLeaseHeartbeat{cancel: cancel, done: done} - go func() { + util.Go(func() { ticker := time.NewTicker(interval) defer ticker.Stop() for { @@ -117,7 +118,7 @@ func startSourceLeaseHeartbeat( } } } - }() + }) return workCtx, heartbeat, nil } diff --git a/internal/apps/upload/cache/access_cache.go b/internal/apps/upload/cache/access_cache.go index d513b1a5..4b39c1fd 100644 --- a/internal/apps/upload/cache/access_cache.go +++ b/internal/apps/upload/cache/access_cache.go @@ -17,6 +17,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/pkg/util" ) const fileAccessInvalidationChannel = "upload:file_access_invalidation" @@ -57,7 +58,7 @@ func startAccessCacheInvalidationListener() { return } - go func() { + util.Go(func() { pubsub := redis.Subscribe( context.Background(), objectstore.ConfigInvalidationChannel, @@ -70,7 +71,7 @@ func startAccessCacheInvalidationListener() { for range pubsub.Channel() { ResetAccessCaches() } - }() + }) } // IsFilePublic reports whether uploadType is in the public access whitelist. diff --git a/internal/apps/upload/cache/meta_cache.go b/internal/apps/upload/cache/meta_cache.go index a04177b0..809036ee 100644 --- a/internal/apps/upload/cache/meta_cache.go +++ b/internal/apps/upload/cache/meta_cache.go @@ -13,6 +13,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" "github.com/Rain-kl/Wavelet/pkg/cache/ram" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -54,17 +55,17 @@ func startUploadMetaCacheInvalidationListener() { // 捕获当前客户端:goroutine 不再读可变全局 db.Redis,测试置空/替换全局时不会数据竞争 redisClient := db.Redis - go func() { + util.Go(func() { defer close(uploadMetaListenerDone) pubsub := redisClient.Subscribe(uploadMetaListenerCtx, uploadMetaInvalidationChan) defer func() { _ = pubsub.Close() }() - go func() { + util.Go(func() { <-uploadMetaListenerCtx.Done() _ = pubsub.Close() - }() + }) for msg := range pubsub.Channel() { var payload uploadMetaInvalidationMessage @@ -74,7 +75,7 @@ func startUploadMetaCacheInvalidationListener() { } uploadMetaRAM.Invalidate(payload.ID) } - }() + }) } func publishUploadMetaRAMInvalidation(ctx context.Context, id uint64) { diff --git a/internal/apps/upload/task/storage_migration.go b/internal/apps/upload/task/storage_migration.go index 9e2ea119..e4153511 100644 --- a/internal/apps/upload/task/storage_migration.go +++ b/internal/apps/upload/task/storage_migration.go @@ -24,6 +24,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/infra/task" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -101,7 +102,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T }() //nolint:contextcheck,gosec - go func() { + util.Go(func() { ticker := time.NewTicker(renewalInterval) defer ticker.Stop() for { @@ -116,7 +117,7 @@ func (h *MigrationHandler) Execute(ctx context.Context, payload []byte) (*task.T return } } - }() + }) } active, err := objectstore.LoadConfig(ctx) diff --git a/internal/infra/objectstore/storage.go b/internal/infra/objectstore/storage.go index 984a69bd..8822ec7e 100644 --- a/internal/infra/objectstore/storage.go +++ b/internal/infra/objectstore/storage.go @@ -16,6 +16,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/pkg/util" "gorm.io/gorm" ) @@ -86,7 +87,7 @@ func startPubSubListener() { if db.Redis == nil { return } - go func() { + util.Go(func() { pubsub := db.Redis.Subscribe(context.Background(), ConfigInvalidationChannel) defer func() { _ = pubsub.Close() @@ -96,7 +97,7 @@ func startPubSubListener() { for range ch { ResetCache() } - }() + }) } // Active returns the configured active driver and backend, using an in-memory cache with 5s TTL. diff --git a/internal/platform/lifecycle/lifecycle.go b/internal/platform/lifecycle/lifecycle.go index edfe2d60..333e30d4 100644 --- a/internal/platform/lifecycle/lifecycle.go +++ b/internal/platform/lifecycle/lifecycle.go @@ -8,6 +8,8 @@ import ( "context" "log" "sync" + + "github.com/Rain-kl/Wavelet/pkg/util" ) // ShutdownFunc defines the signature for a graceful shutdown callback. @@ -52,10 +54,10 @@ func Stop(ctx context.Context) { } done := make(chan struct{}) - go func() { + util.Go(func() { wg.Wait() close(done) - }() + }) select { case <-done: diff --git a/internal/repository/auth_source_cache.go b/internal/repository/auth_source_cache.go index d988d230..4ea564dd 100644 --- a/internal/repository/auth_source_cache.go +++ b/internal/repository/auth_source_cache.go @@ -13,6 +13,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/pkg/cache/ram" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -120,7 +121,7 @@ func startAuthSourceCacheInvalidationListener() { authSourceListenerDone = make(chan struct{}) redisClient := db.Redis // 捕获当前客户端:goroutine 不读可变全局,避免与测试置空 db.Redis 竞争 - go func() { + util.Go(func() { listenerCtx := authSourceListenerCtx defer close(authSourceListenerDone) @@ -129,16 +130,16 @@ func startAuthSourceCacheInvalidationListener() { _ = pubsub.Close() }() - go func() { + util.Go(func() { <-listenerCtx.Done() _ = pubsub.Close() - }() + }) for range pubsub.Channel() { authSourceActiveRAM.InvalidateAll() authSourceByNameRAM.InvalidateAll() } - }() + }) } func publishAuthSourceRAMInvalidation(ctx context.Context) { diff --git a/internal/repository/system_config_cache.go b/internal/repository/system_config_cache.go index 30ec1d3e..07f5db36 100644 --- a/internal/repository/system_config_cache.go +++ b/internal/repository/system_config_cache.go @@ -14,6 +14,7 @@ import ( db "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/pkg/cache/ram" + "github.com/Rain-kl/Wavelet/pkg/util" ) const ( @@ -106,7 +107,7 @@ func startSystemConfigCacheInvalidationListener() { systemConfigListenerDone = make(chan struct{}) redisClient := db.Redis // 捕获当前客户端:goroutine 不读可变全局,避免与测试置空 db.Redis 竞争 - go func() { + util.Go(func() { listenerCtx := systemConfigListenerCtx defer close(systemConfigListenerDone) @@ -115,10 +116,10 @@ func startSystemConfigCacheInvalidationListener() { _ = pubsub.Close() }() - go func() { + util.Go(func() { <-listenerCtx.Done() _ = pubsub.Close() - }() + }) for msg := range pubsub.Channel() { var payload systemConfigBroadcastMessage @@ -134,7 +135,7 @@ func startSystemConfigCacheInvalidationListener() { ram.Delete(payload.Type, key) } } - }() + }) } // StopSystemConfigCacheListener stops the Redis Pub/Sub subscription listener and resets the sync.Once guard. diff --git a/internal/router/router.go b/internal/router/router.go index ef821e72..456f49d2 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -23,6 +23,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/oauth" "github.com/Rain-kl/Wavelet/internal/infra/config" otel_trace "github.com/Rain-kl/Wavelet/pkg/trace" + "github.com/Rain-kl/Wavelet/pkg/util" "github.com/gin-contrib/sessions" "github.com/gin-contrib/sessions/redis" "github.com/gin-gonic/gin" @@ -93,12 +94,12 @@ func Serve(onStarted func()) { onStarted() } - go func() { + util.Go(func() { log.Printf("[API] server listening on %s\n", config.Config.App.Addr) if err := srv.Serve(listener); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Fatalf("[API] server failed: %v\n", err) } - }() + }) quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) diff --git a/pkg/push/email.go b/pkg/push/email.go index 98ff5f64..42d86c08 100644 --- a/pkg/push/email.go +++ b/pkg/push/email.go @@ -10,6 +10,8 @@ import ( "net" "net/smtp" "strings" + + "github.com/Rain-kl/Wavelet/pkg/util" ) func init() { @@ -86,9 +88,9 @@ func (p *EmailPusher) Send(ctx context.Context, cfg Config, target string, body // 异步超时处理 errChan := make(chan error, 1) - go func() { + util.Go(func() { errChan <- smtp.SendMail(host+":"+port, auth, from, []string{to}, msg) - }() + }) select { case <-ctx.Done(): diff --git a/pkg/util/goroutine.go b/pkg/util/goroutine.go new file mode 100644 index 00000000..96dc36ee --- /dev/null +++ b/pkg/util/goroutine.go @@ -0,0 +1,31 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package util + +import ( + "log/slog" + "runtime" + "runtime/debug" +) + +// Go runs fn in a new goroutine and recovers panics, so a background task +// cannot crash the whole process. The panic is logged together with the +// util.Go call site. Use it for every fire-and-forget / long-lived +// background goroutine; HTTP handlers are already covered by gin.Recovery. +func Go(fn func()) { + pc, file, line, _ := runtime.Caller(1) + go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("panic recovered in background goroutine", + "caller", runtime.FuncForPC(pc).Name(), + "file", file, + "line", line, + "panic", r, + "stack", string(debug.Stack())) + } + }() + fn() + }() +} diff --git a/pkg/util/goroutine_test.go b/pkg/util/goroutine_test.go new file mode 100644 index 00000000..ae49aa68 --- /dev/null +++ b/pkg/util/goroutine_test.go @@ -0,0 +1,31 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package util + +import ( + "testing" + "time" +) + +func TestGoRunsFn(t *testing.T) { + done := make(chan struct{}) + Go(func() { close(done) }) + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("fn was not run") + } +} + +func TestGoRecoversPanic(t *testing.T) { + done := make(chan struct{}) + Go(func() { + defer close(done) + panic("boom") + }) + <-done + // Give the recovering goroutine a moment to finish logging; the test + // only fails if the panic had propagated and crashed the process. + time.Sleep(10 * time.Millisecond) +} diff --git a/pkg/wsclient/client.go b/pkg/wsclient/client.go index 32ec5a70..9af9dca3 100644 --- a/pkg/wsclient/client.go +++ b/pkg/wsclient/client.go @@ -16,6 +16,7 @@ import ( "sync" "time" + "github.com/Rain-kl/Wavelet/pkg/util" "golang.org/x/net/websocket" ) @@ -198,13 +199,13 @@ func (conn *Connection) RunReceiveLoop(ctx context.Context, handler MessageHandl doneChan := make(chan struct{}) defer close(doneChan) - go func() { + util.Go(func() { select { case <-ctx.Done(): _ = conn.Close() case <-doneChan: } - }() + }) if err := handler.OnConnect(ctx); err != nil { handler.OnClose(err)