From c9642fb5e27c065b86d7659b435d681d0d5446f1 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 12 Jun 2026 14:32:27 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E9=80=80=E5=87=BA=E5=BC=82?= =?UTF-8?q?=E5=B8=B8=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/task/scheduler/scheduler.go | 30 +++++++++++----- internal/task/scheduler/scheduler_test.go | 42 +++++++++++++++++++++++ 2 files changed, 64 insertions(+), 8 deletions(-) create mode 100644 internal/task/scheduler/scheduler_test.go diff --git a/internal/task/scheduler/scheduler.go b/internal/task/scheduler/scheduler.go index 90b69219..fd5a2589 100644 --- a/internal/task/scheduler/scheduler.go +++ b/internal/task/scheduler/scheduler.go @@ -6,7 +6,9 @@ package scheduler import ( "context" "fmt" + "os/signal" "sync" + "syscall" "time" "github.com/Rain-kl/Wavelet/internal/logger" @@ -33,6 +35,7 @@ func StartScheduler() error { var err error schedulerOnce.Do(func() { quitChan = make(chan struct{}) + done := quitChan // 初始化并运行首次调度 if err = ReloadScheduler(); err != nil { @@ -40,8 +43,12 @@ func StartScheduler() error { return } - // 阻塞等待 - <-quitChan + signalCtx, stopSignals := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stopSignals() + + if waitForStop(done, signalCtx.Done()) { + StopScheduler() + } }) return err } @@ -114,14 +121,21 @@ func ReloadScheduler() error { } } - // 5. 替换全局调度器并异步启动 + // 5. 启动并替换全局调度器。进程信号由 StartScheduler 统一处理。 + if err := newScheduler.Start(); err != nil { + return fmt.Errorf("start scheduler failed: %w", err) + } activeScheduler = newScheduler - go func() { - if err := activeScheduler.Run(); err != nil { - logger.ErrorF(context.Background(), "[Scheduler] 调度器运行错误: %v", err) - } - }() logger.InfoF(context.Background(), "[Scheduler] 成功重新加载定时任务,共注册 %d 个活动任务", len(schedules)) return nil } + +func waitForStop(done <-chan struct{}, signals <-chan struct{}) bool { + select { + case <-done: + return false + case <-signals: + return true + } +} diff --git a/internal/task/scheduler/scheduler_test.go b/internal/task/scheduler/scheduler_test.go new file mode 100644 index 00000000..c9c2307e --- /dev/null +++ b/internal/task/scheduler/scheduler_test.go @@ -0,0 +1,42 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package scheduler + +import "testing" + +func TestWaitForStop(t *testing.T) { + tests := []struct { + name string + closeDone bool + closeSignal bool + wantSignal bool + }{ + { + name: "explicit stop", + closeDone: true, + }, + { + name: "process signal", + closeSignal: true, + wantSignal: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + done := make(chan struct{}) + signals := make(chan struct{}) + if tt.closeDone { + close(done) + } + if tt.closeSignal { + close(signals) + } + + if got := waitForStop(done, signals); got != tt.wantSignal { + t.Errorf("waitForStop() = %t, want %t", got, tt.wantSignal) + } + }) + } +}