mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
修复退出异常问题
This commit is contained in:
@@ -6,7 +6,9 @@ package scheduler
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os/signal"
|
||||||
"sync"
|
"sync"
|
||||||
|
"syscall"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/logger"
|
"github.com/Rain-kl/Wavelet/internal/logger"
|
||||||
@@ -33,6 +35,7 @@ func StartScheduler() error {
|
|||||||
var err error
|
var err error
|
||||||
schedulerOnce.Do(func() {
|
schedulerOnce.Do(func() {
|
||||||
quitChan = make(chan struct{})
|
quitChan = make(chan struct{})
|
||||||
|
done := quitChan
|
||||||
|
|
||||||
// 初始化并运行首次调度
|
// 初始化并运行首次调度
|
||||||
if err = ReloadScheduler(); err != nil {
|
if err = ReloadScheduler(); err != nil {
|
||||||
@@ -40,8 +43,12 @@ func StartScheduler() error {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 阻塞等待
|
signalCtx, stopSignals := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||||
<-quitChan
|
defer stopSignals()
|
||||||
|
|
||||||
|
if waitForStop(done, signalCtx.Done()) {
|
||||||
|
StopScheduler()
|
||||||
|
}
|
||||||
})
|
})
|
||||||
return err
|
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
|
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))
|
logger.InfoF(context.Background(), "[Scheduler] 成功重新加载定时任务,共注册 %d 个活动任务", len(schedules))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func waitForStop(done <-chan struct{}, signals <-chan struct{}) bool {
|
||||||
|
select {
|
||||||
|
case <-done:
|
||||||
|
return false
|
||||||
|
case <-signals:
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user