From 73fe71166b387cb8c861ccd08c2ba843392af313 Mon Sep 17 00:00:00 2001 From: ndou <447662456@qq.com> Date: Wed, 5 Aug 2026 20:58:37 +0800 Subject: [PATCH] fix: stop collectservice deadline ticker --- .../domain/collectservice/common/common.go | 40 ++++++++++++++++++- .../collectservice/common/common_test.go | 27 +++++++++++++ 2 files changed, 65 insertions(+), 2 deletions(-) create mode 100644 manager/backend/services/msg-etl/domain/collectservice/common/common_test.go diff --git a/manager/backend/services/msg-etl/domain/collectservice/common/common.go b/manager/backend/services/msg-etl/domain/collectservice/common/common.go index 1e7ec8b..4790aba 100644 --- a/manager/backend/services/msg-etl/domain/collectservice/common/common.go +++ b/manager/backend/services/msg-etl/domain/collectservice/common/common.go @@ -6,25 +6,61 @@ import ( pkgmq "gitee.com/OpenCloudOS/ocmanager/manager/backend/pkg/library/mq" compb "gitee.com/OpenCloudOS/ocmanager/manager/backend/protocol/common" "google.golang.org/protobuf/proto" + "sync" "sync/atomic" "time" ) var ( nextDeadlineSec int64 // 缓存当前时间1小时后的时间戳 + tickerMu sync.Mutex + tickerStop chan struct{} + tickerDone chan struct{} ) // Init 初始化 func Init() { atomic.StoreInt64(&nextDeadlineSec, time.Now().Add(time.Hour).Unix()) + tickerMu.Lock() + defer tickerMu.Unlock() + if tickerStop != nil { + return + } + stop := make(chan struct{}) + done := make(chan struct{}) + tickerStop = stop + tickerDone = done go func() { + defer close(done) ticker := time.NewTicker(time.Second) - for range ticker.C { - atomic.StoreInt64(&nextDeadlineSec, time.Now().Add(time.Hour).Unix()) + defer ticker.Stop() + for { + select { + case <-ticker.C: + atomic.StoreInt64(&nextDeadlineSec, time.Now().Add(time.Hour).Unix()) + case <-stop: + return + } } }() } +// Stop 停止 Init 启动的 deadline 刷新 goroutine。 +func Stop() { + tickerMu.Lock() + stop := tickerStop + done := tickerDone + if stop == nil { + tickerMu.Unlock() + return + } + tickerStop = nil + tickerDone = nil + close(stop) + tickerMu.Unlock() + <-done +} + // MarshalMetricsMessage 序列化指标转发消息;返回独立副本,可安全在异步 goroutine 中发送。 func MarshalMetricsMessage(msg *compb.MetricsMessage) ([]byte, error) { b, err := proto.Marshal(msg) diff --git a/manager/backend/services/msg-etl/domain/collectservice/common/common_test.go b/manager/backend/services/msg-etl/domain/collectservice/common/common_test.go new file mode 100644 index 0000000..7e7c730 --- /dev/null +++ b/manager/backend/services/msg-etl/domain/collectservice/common/common_test.go @@ -0,0 +1,27 @@ +package common + +import ( + "testing" + "time" +) + +func TestInitStopTickerLifecycle(t *testing.T) { + Init() + Init() + + done := make(chan struct{}) + go func() { + Stop() + Stop() + close(done) + }() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("Stop should return after closing the ticker goroutine") + } + + Init() + Stop() +} -- Gitee