diff --git a/internal/plugin/registry.go b/internal/plugin/registry.go index 8f1ea5c..ed9fb89 100644 --- a/internal/plugin/registry.go +++ b/internal/plugin/registry.go @@ -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 路径(子进程退出回调)产生锁序风险。 diff --git a/internal/plugin/stopall_parallel_test.go b/internal/plugin/stopall_parallel_test.go new file mode 100644 index 0000000..2d3b4c2 --- /dev/null +++ b/internal/plugin/stopall_parallel_test.go @@ -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() +}