mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-27 12:53:35 +00:00
fix(plugin): StopAll 并行停插件 —— 修关停必然超时被 SIGKILL
## 现象 systemd 每次都报 `State 'stop-sigterm' timed out. Killing.` 进程组里 23 个子进程插件**全退完了**,最后那条 `[homed] stopped` 仍打不出来,然后被 SIGKILL。 ## 根因(算出来的,不是猜的) 单个插件的 Stop 最坏预算: 5s(CallContext plugin.stop)+ 5s(等 exited)+ 2s(Kill 后收割)= 12s 串行停 23 个 ⇒ 23 × 12s = 276s,而 systemd 只给 90s。 ⇒ 关停必然超时。线上每一条 stop 记录都是 timed out,无一例外。 ## 改法 StopAll 改为并行:取插件快照 + 各自的 SDK 句柄后**立即释放 registry 锁**, 每个插件一个 goroutine,等全部完成再释放共享段。 三处必须小心的点(都是并行化引入的新风险): 1. **先释放 registry 锁再并行停**。p.Stop() 会触发 markExited → onExit → ReclaimOwner,那条链要读共享内存段。 持着锁并行跑,若某插件的 onExit 需要拿 registry 锁(摘通道等) 就是自死锁。 2. **stop handler 的快照要在清空 sdkRefs 之前取**。handler 挂在 PluginSDK 上(r.sdkRefs),先清空就再也拿不到了。 3. **单个插件 panic 不带崩关停**(那会让剩下的插件全停不掉), 也不静默吞(留日志)。 ## 判据:5 条 + 3 组变异 + race - TestStopAllStopsInParallel:8 个插件各 120ms,串行 960ms / 并行 120ms。 判据直接量耗时,串行实现必然变红(实测串行时 962ms)。 - TestStopAllSurvivesPanickingPlugin:panic 不外冒、其它插件照停、不死锁 (带 10s 超时,死锁会超时而不是挂住测试) - TestStopAllFreezesAutoRestart / StopsEachPluginExactlyOnce: 冻结自动重启、每个插件恰好 Stop 一次(重复会二次释放共享段) - TestStopAllRunsStopHandlerBeforeStop:handler 必须先于 Stop, 走真实的 PluginSDK.RegisterStopHandler 路径 变异:退回串行 → 判红并打出实测耗时;去掉 handler 调用 → 顺序判红; 假装并行只清空 → 4 条判红。 `-race` 通过(并行化必须验锁,这是本次改动的头号风险)。 全量 41 包绿。
This commit is contained in:
@ -707,17 +707,40 @@ func (r *Registry) runOnRemoveHandlers(name string) {
|
||||
}
|
||||
}
|
||||
|
||||
// stopTarget 是 StopAll 并行化时的一个停止单元:插件本体 + 它的 SDK 句柄。
|
||||
type stopTarget struct {
|
||||
plugin sdk.Plugin
|
||||
sdkRef *sdk.PluginSDK
|
||||
}
|
||||
|
||||
func (r *Registry) StopAll() {
|
||||
// 关停开始即冻结自动重启:否则「Stop 触发退出 → 崩溃判定 → 重新 spawn」
|
||||
// 会在内核正在关停时把子进程又拉起来,段已拆而进程还在,直接 SIGBUS。
|
||||
r.shuttingDown.Store(true)
|
||||
|
||||
// 取插件快照后**立即释放 registry 锁**,再并行停。
|
||||
//
|
||||
// 为何必须并行:串行时最坏耗时 = Σ(每个插件) = 5s(plugin.stop 调用)
|
||||
// + 5s(等退出) + 2s(收割) = 12s;线上有 23 个子进程插件,
|
||||
// 即 276s,而 systemd 只给 90s ⇒ 关停必然 timed out 然后 SIGKILL。
|
||||
// 实测确实每次都超时(线上日志里 23 个插件全退完了,
|
||||
// 最后那条 "[homed] stopped" 仍打不出来)。
|
||||
//
|
||||
// 为何先释放锁:p.Stop() 会触发 markExited → onExit → ReclaimOwner,
|
||||
// 那条链要读共享内存段。持着 registry 锁并行跑,若某插件的 onExit
|
||||
// 回调需要拿 registry 锁(如摘通道),就是自死锁。
|
||||
r.mu.Lock()
|
||||
for _, p := range r.instances {
|
||||
r.runStopHandlers(p.Name())
|
||||
if err := p.Stop(); err != nil {
|
||||
log.Printf("[plugin] stop %s: %v", p.Name(), err)
|
||||
snapshot := make([]sdk.Plugin, len(r.instances))
|
||||
copy(snapshot, r.instances)
|
||||
// stop handler 挂在 PluginSDK 上(r.sdkRefs),必须**在清空 sdkRefs 之前**
|
||||
// 把 handler 跑掉 —— 否则下面并行 goroutine 里就找不到它了。
|
||||
// runStopHandlers 自己不加锁(调用方持锁),这里正是持锁状态。
|
||||
stoppers := make([]stopTarget, 0, len(snapshot))
|
||||
for _, p := range snapshot {
|
||||
if p == nil {
|
||||
continue
|
||||
}
|
||||
stoppers = append(stoppers, stopTarget{plugin: p, sdkRef: r.sdkRefs[p.Name()]})
|
||||
}
|
||||
r.plugins = make(map[string]sdk.Plugin)
|
||||
r.instances = nil
|
||||
@ -725,6 +748,30 @@ func (r *Registry) StopAll() {
|
||||
r.sdkRefs = make(map[string]*sdk.PluginSDK)
|
||||
r.mu.Unlock()
|
||||
|
||||
// 并行停:每个插件一个 goroutine,等全部完成。
|
||||
// 单个插件 panic 不带崩整个关停(那会让剩下的插件全停不掉),
|
||||
// 也不静默吞掉(留下日志)。
|
||||
var wg sync.WaitGroup
|
||||
for _, t := range stoppers {
|
||||
wg.Add(1)
|
||||
go func(tg stopTarget) {
|
||||
defer wg.Done()
|
||||
defer func() {
|
||||
if rec := recover(); rec != nil {
|
||||
log.Printf("[plugin] stop %s panic: %v", tg.plugin.Name(), rec)
|
||||
}
|
||||
}()
|
||||
// stop handler(解绑通道等)必须先于 Stop:见 runStopHandlers 注释。
|
||||
if tg.sdkRef != nil {
|
||||
tg.sdkRef.RunStopHandlers()
|
||||
}
|
||||
if err := tg.plugin.Stop(); err != nil {
|
||||
log.Printf("[plugin] stop %s: %v", tg.plugin.Name(), err)
|
||||
}
|
||||
}(t)
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// 共享段在全部子进程退出后再释放:插件还持有映射时拆段,
|
||||
// 它们下一次访问就是 SIGBUS。在锁外调用:Close 不需 registry 锁,
|
||||
// 而持锁调它会与 onProcCrash 路径(子进程退出回调)产生锁序风险。
|
||||
|
||||
190
internal/plugin/stopall_parallel_test.go
Normal file
190
internal/plugin/stopall_parallel_test.go
Normal file
@ -0,0 +1,190 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
)
|
||||
|
||||
// ===== StopAll 并行化 =====
|
||||
//
|
||||
// 串行停 23 个插件、每个最坏 5s(plugin.stop 调用)+5s(等退出)+2s(收割)
|
||||
// = 最坏 276s,而 systemd 只给 90s ⇒ 关停几乎必然被 SIGKILL。
|
||||
// 线上实测:每次 stop 都 "State 'stop-sigterm' timed out",
|
||||
// 进程组里 23 个插件全退完了,最后那条 [homed] stopped 仍打不出来。
|
||||
//
|
||||
// 并行后最坏约等于单个插件的预算(约 12s),而不是 N 倍。
|
||||
//
|
||||
// ★ 并行化最大的风险是死锁与重复释放:Stop 会触发 markExited →
|
||||
// onExit → ReclaimOwner,后者要读共享内存段。所以判据同时盯
|
||||
// "真的并行"与"不死锁、不错杀"。
|
||||
|
||||
// stubPlugin 是最小可用插件。
|
||||
type stubPlugin struct {
|
||||
name string
|
||||
delay time.Duration
|
||||
stopped atomic.Bool
|
||||
panicOnStop bool
|
||||
stopCount atomic.Int32
|
||||
}
|
||||
|
||||
func newStubPlugin(name string, delay time.Duration) *stubPlugin {
|
||||
return &stubPlugin{name: name, delay: delay}
|
||||
}
|
||||
|
||||
func (s *stubPlugin) Name() string { return s.name }
|
||||
func (s *stubPlugin) Start(*sdk.PluginSDK) error { return nil }
|
||||
func (s *stubPlugin) Stop() error {
|
||||
s.stopCount.Add(1)
|
||||
if s.panicOnStop {
|
||||
panic("stub: 故意 panic")
|
||||
}
|
||||
if s.delay > 0 {
|
||||
time.Sleep(s.delay)
|
||||
}
|
||||
s.stopped.Store(true)
|
||||
return nil
|
||||
}
|
||||
|
||||
// StopAll 必须真的并行,否则退回串行 = 关停超时。
|
||||
func TestStopAllStopsInParallel(t *testing.T) {
|
||||
const n = 8
|
||||
const each = 120 * time.Millisecond
|
||||
|
||||
var plugins []sdk.Plugin
|
||||
var stubs []*stubPlugin
|
||||
for i := 0; i < n; i++ {
|
||||
s := newStubPlugin("p", each)
|
||||
_ = i
|
||||
stubs = append(stubs, s)
|
||||
plugins = append(plugins, s)
|
||||
}
|
||||
r := &Registry{instances: plugins}
|
||||
|
||||
start := time.Now()
|
||||
r.StopAll()
|
||||
elapsed := time.Since(start)
|
||||
|
||||
serial := n * each
|
||||
// 串行实现会耗时 serial;留一半余量仍能可靠区分。
|
||||
if elapsed > serial/2 {
|
||||
t.Errorf("StopAll 是串行的:%d 个插件各 %v 用了 %v(串行≈%v,并行应≈%v)",
|
||||
n, each, elapsed, serial, each)
|
||||
}
|
||||
for _, s := range stubs {
|
||||
if !s.stopped.Load() {
|
||||
t.Errorf("插件 %s 没被停掉", s.name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 一个插件 panic 不能带崩整个关停,也不能让其它插件停不掉。
|
||||
func TestStopAllSurvivesPanickingPlugin(t *testing.T) {
|
||||
bad := newStubPlugin("bad", 0)
|
||||
bad.panicOnStop = true
|
||||
good := newStubPlugin("good", 0)
|
||||
r := &Registry{instances: []sdk.Plugin{bad, good}}
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
defer func() {
|
||||
if rec := recover(); rec != nil {
|
||||
t.Errorf("插件的 panic 不该冒到关停流程上:%v", rec)
|
||||
}
|
||||
}()
|
||||
r.StopAll()
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("StopAll 卡死(可能死锁)")
|
||||
}
|
||||
if !good.stopped.Load() {
|
||||
t.Error("一个插件 panic 不该让其它插件停不掉")
|
||||
}
|
||||
}
|
||||
|
||||
// 关停期间必须冻结自动重启,否则崩溃判定会在关停中把插件重新拉起
|
||||
// (段已拆而进程还在 → SIGBUS)。
|
||||
func TestStopAllFreezesAutoRestart(t *testing.T) {
|
||||
s := newStubPlugin("p", 0)
|
||||
r := &Registry{instances: []sdk.Plugin{s}}
|
||||
r.StopAll()
|
||||
if !r.shuttingDown.Load() {
|
||||
t.Error("StopAll 之后 shuttingDown 应为 true")
|
||||
}
|
||||
if r.instances != nil {
|
||||
t.Error("StopAll 之后 instances 应被清空")
|
||||
}
|
||||
}
|
||||
|
||||
// 每个插件恰好 Stop 一次:重复调用会二次释放共享段/重复跑 stop handler。
|
||||
func TestStopAllStopsEachPluginExactlyOnce(t *testing.T) {
|
||||
var plugins []sdk.Plugin
|
||||
var stubs []*stubPlugin
|
||||
for i := 0; i < 5; i++ {
|
||||
s := newStubPlugin("p", 0)
|
||||
stubs = append(stubs, s)
|
||||
plugins = append(plugins, s)
|
||||
}
|
||||
r := &Registry{instances: plugins}
|
||||
r.StopAll()
|
||||
for _, s := range stubs {
|
||||
if c := s.stopCount.Load(); c != 1 {
|
||||
t.Errorf("Stop 被调 %d 次,应恰好 1 次", c)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// stop handler 必须在 Stop 之前跑完(handler 负责解绑通道等),
|
||||
// 且并行化后这个顺序不能被破坏。
|
||||
//
|
||||
// handler 挂在 PluginSDK 上(RegisterStopHandler),由 runStopHandlers
|
||||
// 通过 r.sdkRefs 取出执行 —— 所以判据必须真的构造一个 PluginSDK,
|
||||
// 否则测的是一条不存在的注册路径。
|
||||
func TestStopAllRunsStopHandlerBeforeStop(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var order []string
|
||||
|
||||
s := newStubPlugin("p", 0)
|
||||
r := &Registry{
|
||||
instances: []sdk.Plugin{s},
|
||||
sdkRefs: map[string]*sdk.PluginSDK{},
|
||||
}
|
||||
sdkn := sdk.New("p", sdk.SDKConfig{})
|
||||
sdkn.RegisterStopHandler(func() {
|
||||
mu.Lock()
|
||||
order = append(order, "handler")
|
||||
mu.Unlock()
|
||||
})
|
||||
r.sdkRefs["p"] = sdkn
|
||||
|
||||
probe := &orderProbePlugin{stubPlugin: s, onStop: func() {
|
||||
mu.Lock()
|
||||
order = append(order, "stop")
|
||||
mu.Unlock()
|
||||
}}
|
||||
r.instances = []sdk.Plugin{probe}
|
||||
|
||||
r.StopAll()
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if len(order) != 2 || order[0] != "handler" || order[1] != "stop" {
|
||||
t.Errorf("stop handler 必须先于 Stop 执行,实际顺序 %v", order)
|
||||
}
|
||||
}
|
||||
|
||||
type orderProbePlugin struct {
|
||||
*stubPlugin
|
||||
onStop func()
|
||||
}
|
||||
|
||||
func (o *orderProbePlugin) Stop() error {
|
||||
o.onStop()
|
||||
return o.stubPlugin.Stop()
|
||||
}
|
||||
Reference in New Issue
Block a user