fix(proc): 子进程崩溃自愈 + 集中台账 + 注册面摘除

根因:子进程插件被 kill 后,内核只发了一个无人订阅的事件,
工具/stage handler/IO 通道全留在注册表里指向死进程,
模型继续调用只吃 ErrProcessExited,没有任何路径把插件拉回来。

## 四层修复

### 1. 专职 waitLoop(进程收割)
- 每个子进程配一根 waitLoop goroutine,是 cmd.Wait() 的唯一调用点
- 不再依赖 stdout EOF 判定死亡(孙子进程继承 stdout 时 EOF 永不到来)
- 手工 os.Pipe 替代 cmd.StdinPipe/StdoutPipe,避免 waitLoop 与
  os/exec 的内部关闭竞争
- host.go: Host.Supervisor(),Host.Close() 先 StopAll 再拆段

### 2. 集中台账 Supervisor
- proc/supervisor.go: 插件 Spawn 握手成功即 track,进程退出即 untrack
- StopAll: 并发发 plugin.stop 走优雅路径,到期仍在的一律 Kill
- 关停后才完成握手的进程被立即结束,不会活过内核
- 消除「孤儿进程持共享段映射 → SIGBUS」的隐患

### 3. 注册面摘除(detachPlugin)
- 新增 StageHost.UnregisterPluginStages:摘除指定插件的全部 stage handler
- 新增 Registry.pluginChannels 台账:记录每个插件注册的 IO 通道
- 三条路径统一走 detachPlugin:Disable / ReloadOne / RemovePlugin
- StopAndUnload 漏了 IO 通道也一并补上

### 4. 自动重启
- onProcCrash 从「只发事件」改为「摘注册面 → 从注册表移除 → 异步排重启」
- scheduleProcRestart: 窗口 5 分钟内最多 3 次,线性退避 1s/2s/3s
- 超限停手留日志;重启前复核是否已被 Disable 或被其他路径加载
- 崩溃计数窗口过期自动归零

### 5. 主动停止 vs 崩溃的区分
- proc.Plugin 新增 stopping 标志:Stop()/Close() 里 Set(true)
- handleExit 读 stopping 标志,主动停止不上报 onCrash
- 防止重载/禁用/卸载被误判为崩溃触发多余重启

### 6. Linux Pdeathsig 兜底
- procattr_linux.go: SysProcAttr.Pdeathsig = SIGKILL
- 兜 homed 自身被 SIGKILL/OOM 时子进程变孤儿的场景
- macOS/Windows 无等价物,空实现

### 7. pluginmgr 升级
- PluginManager 接口新增 PluginRuntime / ListPluginRuntimes
- plugin_list 输出运行态:loaded / alive / pid / crash_count / channel
- 新增 plugin_status: 全量运行期快照 + dead/unhealthy 汇总
- 新增 plugin_restart: 无条件重启单个插件(plgreload 不动未改二进制的插件)

### 测试
- process_test.go: 3 例(grandchild stdout 感知 / Supervisor track-untrack /
  StopAll 无孤儿)
- crash_recovery_test.go: 8 例(detach 三项齐全 / 通道重注册 / 崩溃不阻塞 /
  退避阈值 / 窗口过期 / 关停中跳过 / PluginRuntime 通道识别)
- stages_plugin_test.go: 4 例(stage 按插件摘除 / 空 stage 清理 / 空名 no-op /
  工具+stage 双摘后可重新注册同名)
This commit is contained in:
JianFeeeee
2026-09-03 12:37:58 +08:00
parent ad1b05e6ef
commit 56498d3592
18 changed files with 1649 additions and 74 deletions

View File

@ -45,6 +45,12 @@ type Host struct {
stageMu sync.Mutex
coordMu sync.Mutex
coord *stageCoordinator
// sup 是内核侧唯一的子进程台账,与共享段同生命周期。
//
// 放在 Host 而不是 registry 的理由:能拿到 Host 的地方就能拿到台账,
// 而 Host 本就是「全部子进程插件共享的那一份内核侧状态」。
sup *Supervisor
}
// NewHost 创建共享段(平台层 allocShm + 布局初始化)。
@ -81,6 +87,7 @@ func NewHost() (*Host, error) {
evtRing.Init()
return &Host{
sup: NewSupervisor(),
memfd: memfd,
data: data,
seg: seg,
@ -102,7 +109,15 @@ func NewHost() (*Host, error) {
const shmDefaultSize = 256 * 1024
// Close 释放共享段(StageContext + 事件环)。
// Supervisor 返回子进程台账(供 registry 查询/关停)。
func (h *Host) Supervisor() *Supervisor { return h.sup }
func (h *Host) Close() error {
// 先停全部子进程再拆段:插件还持有映射时 unmap,
// 它们下一次访问共享段就是 SIGBUS。
if h.sup != nil {
h.sup.StopAll(0)
}
var firstErr error
if h.data != nil {
if err := freeShm(h.memfd, h.data); err != nil && firstErr == nil {

View File

@ -6,6 +6,7 @@ import (
"fmt"
"log"
"sync"
"sync/atomic"
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
)
@ -43,6 +44,14 @@ type Plugin struct {
caps *capabilitySet
stopOnce sync.Once
// stopping 标记「本次退出是内核主动发起的」,用于压掉 onCrash。
//
// 必要性:Stop() 宽限期超时与 Close() 都走 Process.Kill(),
// 而 Kill 产生的 `signal: killed` 是非 nil 的 waitErr——若不区分,
// 重载/禁用/卸载这些**内核自己发起**的停止会被 handleExit 当成崩溃上报,
// 触发一轮多余的自动重启(重载路径下等于把刚装好的插件又推倒一次)。
stopping atomic.Bool
}
// New 创建子进程插件(不启动进程)。
@ -67,6 +76,31 @@ func New(name, bin, dir string, config map[string]interface{}, host *Host, onCra
// Name 实现 sdk.Plugin。
func (p *Plugin) Name() string { return p.name }
// PID 返回子进程号;未启动或已退出返回 0。
// 供 pluginmgr 呈现「插件实际在跑哪个进程」。
func (p *Plugin) PID() int {
if p.proc == nil {
return 0
}
if !p.Alive() {
return 0
}
return p.proc.PID()
}
// Alive 报告子进程是否仍存活。
func (p *Plugin) Alive() bool {
if p.proc == nil {
return false
}
select {
case <-p.proc.Exited():
return false
default:
return true
}
}
// Start 启动子进程并完成注册。
//
// core 是内核为该插件构建的能力面(internal/sdk.PluginSDK 天然满足 CoreSDK)。
@ -99,6 +133,7 @@ func (p *Plugin) Start(core CoreSDK) error {
EvtRingSize: evtTotalSize,
Handler: p.handler.Handle,
OnExit: p.handleExit,
Supervisor: p.host.Supervisor(),
})
if err != nil {
return err
@ -124,6 +159,7 @@ func (p *Plugin) Start(core CoreSDK) error {
// Stop 优雅停止(实现 sdk.Plugin)。
func (p *Plugin) Stop() error {
p.stopping.Store(true)
var err error
p.stopOnce.Do(func() {
if p.proc != nil {
@ -138,6 +174,7 @@ func (p *Plugin) Stop() error {
// **这里是真 kill + wait**——对比 cabi 路径的 Close 只做 dlclose,
// 而 dlclose 对 Go c-shared 是 no-op(§1.1,热重载静默失效的根因)。
func (p *Plugin) Close() error {
p.stopping.Store(true)
var err error
p.stopOnce.Do(func() {
if p.proc != nil {
@ -156,6 +193,11 @@ func (p *Plugin) handleExit(name string, err error) {
if p.host != nil && p.host.ForceReleaseLock(name) {
log.Printf("[proc] %s 退出,内核已释放其持有的 stage 锁", name)
}
// 内核主动停止(Stop/Close,含宽限期超时后的 Kill)不算崩溃:
// 否则重载/禁用/卸载都会误触发自动重启。
if p.stopping.Load() {
return
}
if err != nil && p.onCrash != nil {
p.onCrash(name, err)
}

View File

@ -0,0 +1,30 @@
//go:build linux
package proc
import (
"os/exec"
"syscall"
)
// applyProcAttr 让子进程在父进程(homed)死亡时收到 SIGKILL。
//
// 这是**最后一道兜底**,不是主路径:正常关停走 Supervisor.StopAll。
// 它兜的是内核自身异常终止的场景——homed 被 SIGKILL、段错误、OOM——
// 此时没有任何 Go 代码有机会运行,Supervisor 也来不及 StopAll,
// 子进程会被 init 收养成孤儿:
// - 继续持有已被 unmap 的共享段映射,下次访问即 SIGBUS;
// - 与新启动的 homed 抢同一份外部资源(qq 的 WS 会话、browser 的
// chromium profile 锁),表现为"重启后插件时好时坏"。
//
// Pdeathsig 由内核在父进程退出时投递,不依赖任何用户态代码,
// 因此在 homed 被 SIGKILL 的情况下依然生效。
//
// 仅 Linux 有此机制。macOS/Windows 无等价物,回退为空实现(procattr_other.go):
// 那两个平台上孤儿风险依旧存在,靠 StopAll 覆盖正常关停路径。
func applyProcAttr(cmd *exec.Cmd) {
if cmd.SysProcAttr == nil {
cmd.SysProcAttr = &syscall.SysProcAttr{}
}
cmd.SysProcAttr.Pdeathsig = syscall.SIGKILL
}

View File

@ -0,0 +1,15 @@
//go:build !linux
package proc
import "os/exec"
// applyProcAttr 在非 Linux 平台是空实现。
//
// macOS 没有 Pdeathsig(kqueue 的 NOTE_EXIT 要求父进程存活才能监听,
// 恰好在父进程被 SIGKILL 时失效);Windows 的 Job Object 可做到类似效果,
// 但需要额外的句柄管理,且 Windows 侧尚未真机验证(§12.5),不在此引入。
//
// 后果:这两个平台上 homed 被强杀时子进程会成为孤儿。
// 正常关停路径(Supervisor.StopAll)不受影响。
func applyProcAttr(cmd *exec.Cmd) {}

View File

@ -36,6 +36,11 @@ type Process struct {
stdin *bufio.Writer
stdout io.ReadCloser
// stdinFile / stdoutFile 是父进程侧的管道端(手工 os.Pipe,非 cmd.StdinPipe)。
// 持有它们才能在退出时主动 Close,逼 readLoop 从 Scan 里出来。
stdinFile *os.File
stdoutFile *os.File
// writeMu 串行化 stdin 写入:NDJSON 帧不能交错,否则对端解析错乱。
writeMu sync.Mutex
@ -48,18 +53,27 @@ type Process struct {
// handler 处理插件反向发起的调用(51 个 core.* method)。
handler RequestHandler
// exited 在 readLoop 检测到 EOF/进程退出后关闭,用于唤醒所有等待者。
// exited 在进程被收割后关闭,用于唤醒所有等待者。
exited chan struct{}
exitOnce sync.Once
exitErr atomic.Pointer[error]
readerWG sync.WaitGroup
waiterWG sync.WaitGroup
readyOnce sync.Once
ready chan struct{}
// waitErr 由**唯一的** waitLoop 写入:cmd.Wait() 的返回值。
// waitDone 关闭后 waitErr 才可读。
waitErr error
waitDone chan struct{}
// onExit 在进程退出时回调(内核用它喂 plugin_health.recordCrash,
// 以及 ForceRelease 释放该插件持有的 stage 锁)。
onExit func(name string, err error)
// sup 是内核的集中进程表(可为 nil,单测直接 Spawn 时)。
sup *Supervisor
// shmSize 是握手时告知插件的共享段大小(0 表示本插件不用共享段)。
shmSize int
// evtRingSize 是事件环段大小(0 表示不支持事件环)。
@ -87,6 +101,8 @@ type Options struct {
Handler RequestHandler
// OnExit 进程退出回调。
OnExit func(name string, err error)
// Supervisor 是内核的集中进程表;为 nil 时不纳管(单测路径)。
Supervisor *Supervisor
// HandshakeTimeout 建链超时,默认 10s。
HandshakeTimeout time.Duration
}
@ -97,6 +113,9 @@ const (
// stopGracePeriod 是发出 plugin.stop 后等待进程自行退出的时间。
// 超时则 Kill——**这是"真正的取消"**,对比 cgo 路径超时后线程永久泄漏。
stopGracePeriod = 5 * time.Second
// killReapTimeout 是 SIGKILL 后等待 waitLoop 收割的上限。
// 正常情况 wait4 微秒级返回;超过说明卡在不可中断的内核态。
killReapTimeout = 2 * time.Second
)
// ErrProcessExited 表示子进程已退出,调用无法完成。
@ -120,35 +139,68 @@ func Spawn(name, bin string, opts Options) (*Process, error) {
cmd.Env = append(os.Environ(), opts.Env...)
}
cmd.ExtraFiles = opts.ExtraFiles
applyProcAttr(cmd)
stdinPipe, err := cmd.StdinPipe()
// 管道手工创建而非用 cmd.StdinPipe/StdoutPipe。
//
// 原因:cmd.Wait() 会等待并**关闭** StdinPipe/StdoutPipe 创建的管道,
// 且文档明确要求“读完再 Wait”。既然现在有一根专职的 waitLoop 立即
// Wait(不等 readLoop),就必须自己控制管道生命期,否则会与
// os/exec 的内部关闭竞争,在 readLoop 里读到 "file already closed"。
stdinR, stdinW, err := os.Pipe()
if err != nil {
return nil, fmt.Errorf("proc: %s stdin 管道: %w", name, err)
}
stdoutPipe, err := cmd.StdoutPipe()
stdoutR, stdoutW, err := os.Pipe()
if err != nil {
stdinR.Close()
stdinW.Close()
return nil, fmt.Errorf("proc: %s stdout 管道: %w", name, err)
}
cmd.Stdin = stdinR
cmd.Stdout = stdoutW
p := &Process{
name: name,
bin: bin,
dir: opts.Dir,
cmd: cmd,
stdin: bufio.NewWriter(stdinPipe),
stdout: stdoutPipe,
stdin: bufio.NewWriter(stdinW),
stdout: stdoutR,
stdinFile: stdinW,
stdoutFile: stdoutR,
pending: make(map[uint64]chan *Response),
handler: opts.Handler,
exited: make(chan struct{}),
ready: make(chan struct{}),
waitDone: make(chan struct{}),
onExit: opts.OnExit,
sup: opts.Supervisor,
shmSize: opts.ShmSize,
evtRingSize: opts.EvtRingSize,
}
if err := cmd.Start(); err != nil {
stdinR.Close()
stdinW.Close()
stdoutR.Close()
stdoutW.Close()
return nil, fmt.Errorf("proc: 启动 %s (%s): %w", name, bin, err)
}
// 子进程已继承它们,父进程侧关掉对端。
// stdoutW 必须关:否则子进程死后写端仍被父进程持有,readLoop 永不到 EOF。
stdinR.Close()
stdoutW.Close()
// 专职收割协程:这是 cmd.Wait() 的**唯一**调用点。
//
// 为何不能靠 readLoop 的 EOF:EOF 只说明 stdout 写端全部关闭,而插件
// fork 出去的孙子进程(browser 拉 chromium、editdoc 拉 python)继承着
// 同一个 stdout:插件本体死了但孙子还持有写端,EOF 就不来,
// 内核完全感知不到插件已死(进程表里是僵尸,注册表里一切正常)。
// wait 直接盯进程本身,不受 fd 继承影响。
p.waiterWG.Add(1)
go p.waitLoop()
p.readerWG.Add(1)
go p.readLoop()
@ -160,6 +212,9 @@ func Spawn(name, bin string, opts Options) (*Process, error) {
p.Kill()
return nil, err
}
if p.sup != nil {
p.sup.track(p)
}
return p, nil
}
@ -270,22 +325,54 @@ func (p *Process) readLoop() {
log.Printf("[proc] %s 读取 stdout 出错: %v", p.name, err)
}
// stdout 关闭(EOF)意味着进程结束——2.5ms 内即可感知(实验 6)。
// stdout 关闭(EOF)通常意味着进程结束——2.5ms 内即可感知(实验 6)。
//
// 但 EOF **不是**权威信号:插件 fork 的孙子进程继承同一 stdout 写端时,
// 插件本体死了 EOF 也不会到。真正的死亡判定在 waitLoop。
// 这里只等 waitLoop 的结果(若进程确实已退,它立即就给)。
<-p.waitDone
p.markExited()
}
// markExited 回收进程、唤醒所有等待者、触发 onExit 回调。
// waitLoop 是内核侧**唯一**的 cmd.Wait() 调用点,每个子进程一根。
//
// 这是「把 panic 捕获换成进程退出检测」的落点(§2.3):
// plugin_health 的 recordCrash / 冷却 / 自愈 / pendingReloads 全部逻辑复用,
// 只是信号源从 recover() 变成进程退出。
// 为何需要专职协程而不是靠 readLoop 的 EOF:
// 1. **EOF 不等于进程死**。插件用 exec.Command 拉起的孙子进程(browser 拉
// chromium、editdoc 拉 python)默认继承插件的 stdout。插件被 kill 后
// 孙子还活着持有写端,readLoop 就永远阻在 Scan 上——内核根本不知道
// 插件已经死了,工具调用一直超时,自愈也永不触发。
// 2. **不收割就是僵尸进程**。不调 Wait 的已退出子进程以 Z 状态占着 PID 槽位。
// 3. **反应速度**。Wait 底层是 wait4(2),内核侧退出即返回(微秒级),
// 比任何轮询健康检查都快,也不消耗 CPU。
func (p *Process) waitLoop() {
defer p.waiterWG.Done()
p.waitErr = p.cmd.Wait()
close(p.waitDone)
// 主动拆管道:若孙子进程仍持有 stdout 写端,readLoop 不会自己退,
// 关掉读端逼它从 Scan 里出来(报 file already closed,已预期)。
if p.stdoutFile != nil {
_ = p.stdoutFile.Close()
}
if p.stdinFile != nil {
_ = p.stdinFile.Close()
}
p.markExited()
}
// markExited 唤醒所有等待者、触发 onExit 回调(幂等,两条路径可并发调用)。
//
// 这是「把 panic 捕获换成进程退出检测」的落点(§2.3)。
// 注意:不在此处调 cmd.Wait()——它属于 waitLoop,Wait 并非并发安全,
// 两处调会报 "wait: no child processes" 或丢失真实退出码。
func (p *Process) markExited() {
p.exitOnce.Do(func() {
waitErr := p.cmd.Wait()
if waitErr != nil {
e := fmt.Errorf("插件进程 %s 异常退出: %w", p.name, waitErr)
<-p.waitDone // 保证 waitErr 可读
if p.waitErr != nil {
e := fmt.Errorf("插件进程 %s 异常退出: %w", p.name, p.waitErr)
p.exitErr.Store(&e)
log.Printf("[proc] %s 退出: %v", p.name, waitErr)
log.Printf("[proc] %s 退出: %v", p.name, p.waitErr)
} else {
log.Printf("[proc] %s 正常退出", p.name)
}
@ -305,6 +392,9 @@ func (p *Process) markExited() {
}
close(p.exited)
if p.sup != nil {
p.sup.untrack(p.name)
}
if p.onExit != nil {
p.onExit(p.name, p.ExitError())
}
@ -491,11 +581,14 @@ func (p *Process) Kill() error {
return nil
}
err := p.cmd.Process.Kill()
// 等 readLoop 观察到 EOF 并完成 Wait/清理
// 等 waitLoop 收割完成。不再在此兜底调 markExited:
// cmd.Wait 只能由 waitLoop 调一次,两处调会报 "wait: no child processes"。
select {
case <-p.exited:
case <-time.After(2 * time.Second):
p.markExited() // 兜底:极端情况下强制走清理
case <-time.After(killReapTimeout):
// SIGKILL 后仍未收割:进程卡在不可中断的内核态(D 状态,如 NFS I/O)。
// 不能无限等,否则重载路径整体挂死;留日志供定位。
log.Printf("[proc] %s SIGKILL 后 %v 仍未被收割(进程可能卡在内核态)", p.name, killReapTimeout)
}
p.readerWG.Wait()
if err != nil && !errors.Is(err, os.ErrProcessDone) {

View File

@ -332,3 +332,150 @@ func TestProcess_SpawnRequiresHandler(t *testing.T) {
t.Fatal("缺少 Handler 应报错(插件无法回调内核)")
}
}
// 插件死亡但孙子进程仍持有 stdout 写端时,内核必须仍能感知退出。
//
// 这是「EOF 不等于进程死亡」的回归测试。旧实现只在 readLoop 读到 EOF 后
// 才 markExited,而 exec.Command 起的孙子进程默认继承插件的 stdout:
// 插件本体退出后写端仍被孙子持有,EOF 永不到来,于是
// - 在途调用挂到自己的超时;
// - OnExit 不触发 → 崩溃计数、工具摘除、自动重启全都不发生;
// - 进程表里插件已是僵尸,注册表里却一切正常。
// 生产上 browser 拉 chromium、editdoc 拉 python 正是这个形状。
// 现在由专职 waitLoop 直接 wait4(2) 判定,不再依赖 fd 生命周期。
func TestProcess_ExitDetectedDespiteInheritedStdout(t *testing.T) {
if _, err := exec.LookPath("sleep"); err != nil {
t.Skip("环境无 sleep,跳过")
}
bin := buildTestPlugin(t, "forkplugin.go")
exitCh := make(chan error, 1)
p, err := Spawn("fork", bin, Options{
Handler: noopHandler,
OnExit: func(name string, err error) { exitCh <- err },
})
if err != nil {
t.Fatalf("Spawn: %v", err)
}
defer p.Kill()
// 让插件本体退出(孙子 sleep 300 仍活着,继续持有 stdout 写端)
if _, callErr := p.Call(MethodToolInvoke, ToolInvokeParams{Name: "die"}); callErr == nil {
t.Error("插件退出时在途调用应返回错误")
}
select {
case exitErr := <-exitCh:
if exitErr == nil {
t.Error("非零退出码应报告为错误(供崩溃计数使用)")
}
case <-time.After(5 * time.Second):
t.Fatal("孙子进程持有 stdout 时未能感知插件退出——退化回只靠 EOF 判定")
}
if _, err := p.Call(MethodToolInvoke, ToolInvokeParams{Name: "x"}); !errors.Is(err, ErrProcessExited) {
t.Errorf("退出后调用应返回 ErrProcessExited,实际 %v", err)
}
}
// Supervisor 台账:握手成功即在册,进程退出即注销。
func TestSupervisor_TrackAndUntrack(t *testing.T) {
bin := buildTestPlugin(t, "echoplugin.go")
sup := NewSupervisor()
p, err := Spawn("echo", bin, Options{Handler: noopHandler, Supervisor: sup})
if err != nil {
t.Fatalf("Spawn: %v", err)
}
if sup.Count() != 1 {
t.Fatalf("握手成功后应在册,实际 %d", sup.Count())
}
got, ok := sup.Get("echo")
if !ok || got.PID() != p.PID() {
t.Errorf("台账里的进程应是刚 spawn 的那个")
}
list := sup.List()
if len(list) != 1 || !list[0].Alive || list[0].PID != p.PID() {
t.Errorf("List 应报告存活与 PID,实际 %+v", list)
}
if err := p.Stop(); err != nil {
t.Fatalf("Stop: %v", err)
}
// 退出回调在 markExited 里注销,等它落地
deadline := time.Now().Add(3 * time.Second)
for sup.Count() != 0 && time.Now().Before(deadline) {
time.Sleep(10 * time.Millisecond)
}
if sup.Count() != 0 {
t.Errorf("进程退出后应注销,实际仍有 %d 个在册", sup.Count())
}
}
// StopAll 必须停掉全部在册子进程——内核关停时不留孤儿。
func TestSupervisor_StopAllLeavesNoSurvivor(t *testing.T) {
bin := buildTestPlugin(t, "echoplugin.go")
sup := NewSupervisor()
var procs []*Process
for i := 0; i < 3; i++ {
p, err := Spawn(fmt.Sprintf("echo%d", i), bin, Options{Handler: noopHandler, Supervisor: sup})
if err != nil {
t.Fatalf("Spawn %d: %v", i, err)
}
procs = append(procs, p)
}
if sup.Count() != 3 {
t.Fatalf("应有 3 个在册,实际 %d", sup.Count())
}
sup.StopAll(5 * time.Second)
for _, p := range procs {
select {
case <-p.Exited():
case <-time.After(2 * time.Second):
t.Errorf("%s 未被 StopAll 停掉(会成为孤儿进程)", p.Name())
}
}
}
// 卡死插件(不响应 plugin.stop)必须在 StopAll 的预算内被强杀。
func TestSupervisor_StopAllKillsUnresponsive(t *testing.T) {
bin := buildTestPlugin(t, "hangplugin.go")
sup := NewSupervisor()
p, err := Spawn("hang", bin, Options{Handler: noopHandler, Supervisor: sup})
if err != nil {
t.Fatalf("Spawn: %v", err)
}
// 预算给足以覆盖 stopGracePeriod,之后剩下的一律 Kill
sup.StopAll(500 * time.Millisecond)
select {
case <-p.Exited():
case <-time.After(10 * time.Second):
t.Error("不响应 plugin.stop 的插件应被强制结束,否则 homed 关停会被它拖住")
}
}
// 关停后完成握手的进程不得留存:立即被结束,不能活过内核。
func TestSupervisor_TrackAfterCloseKillsProcess(t *testing.T) {
bin := buildTestPlugin(t, "echoplugin.go")
sup := NewSupervisor()
sup.StopAll(time.Second) // 置 closed
p, err := Spawn("late", bin, Options{Handler: noopHandler, Supervisor: sup})
if err != nil {
t.Fatalf("Spawn: %v", err)
}
if sup.Count() != 0 {
t.Errorf("关停后不应再纳管新进程,实际在册 %d", sup.Count())
}
select {
case <-p.Exited():
case <-time.After(3 * time.Second):
t.Error("关停后冒出的进程应被立即结束")
}
}

View File

@ -0,0 +1,163 @@
package proc
import (
"fmt"
"log"
"sort"
"sync"
"time"
)
// Supervisor 是内核侧**唯一**的子进程台账。
//
// 为什么必须有它,而不是让每个 Plugin 各自管好自己的 Process:
//
// 1. **没有台账就没有"全部子进程"这个概念**。内核关停时只能遍历 registry 的
// 插件表逐个 Stop,而 registry 表是按插件名索引的——握手失败、Start 中途
// 出错、或刚 spawn 还没进表就崩了的进程,registry 根本不知道它们存在,
// 那些进程会变成孤儿(ppid=1)继续跑,还持有共享段映射。
// 2. **诊断面缺失**。此前 `/api/manager/status` 之类的接口拿不到"实跑几个子进程、
// 各自 PID 多少、活了多久、崩过几次",运维只能 ps | grep。
// 3. **收割保证**。每个 Process 自带一根 waitLoop 立即 wait4(2),Supervisor
// 只负责登记/注销与聚合视图;两者配合才能做到"进程一死内核立刻知道"。
//
// 生命周期:Spawn 成功握手后 track,Process.markExited 里 untrack。
type Supervisor struct {
mu sync.RWMutex
procs map[string]*Process
// closed 后拒绝新的 track,防止关停竞态里又冒出新进程。
closed bool
}
// NewSupervisor 创建空台账。
func NewSupervisor() *Supervisor {
return &Supervisor{procs: make(map[string]*Process)}
}
// track 登记一个已握手成功的子进程。
//
// 同名覆盖是正常情况(重载:旧进程 untrack 早于或晚于新进程 track 都可能,
// 取决于 Kill 与 Spawn 的交错),故不报错,只在真覆盖时留日志。
func (s *Supervisor) track(p *Process) {
if p == nil {
return
}
s.mu.Lock()
defer s.mu.Unlock()
if s.closed {
// 关停途中还有进程完成握手:立即结束它,不让它活过内核。
go p.Kill()
return
}
if old, ok := s.procs[p.name]; ok && old != p {
log.Printf("[proc] 台账中 %s 已有 pid=%d,被 pid=%d 覆盖", p.name, old.PID(), p.PID())
}
s.procs[p.name] = p
}
// untrack 注销(进程已退出)。只有当表里那一项确实是它时才删,
// 避免重载时新进程被旧进程的退出回调误删。
func (s *Supervisor) untrack(name string) {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.procs, name)
}
// Get 按插件名取子进程句柄。
func (s *Supervisor) Get(name string) (*Process, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
p, ok := s.procs[name]
return p, ok
}
// Count 返回在册子进程数。
func (s *Supervisor) Count() int {
s.mu.RLock()
defer s.mu.RUnlock()
return len(s.procs)
}
// ProcInfo 是单个子进程的运行期快照。
type ProcInfo struct {
Name string `json:"name"`
PID int `json:"pid"`
Alive bool `json:"alive"`
Bin string `json:"bin"`
}
// List 返回全部在册子进程的快照(按插件名排序,便于稳定展示)。
func (s *Supervisor) List() []ProcInfo {
s.mu.RLock()
out := make([]ProcInfo, 0, len(s.procs))
for name, p := range s.procs {
alive := true
select {
case <-p.Exited():
alive = false
default:
}
out = append(out, ProcInfo{Name: name, PID: p.PID(), Alive: alive, Bin: p.bin})
}
s.mu.RUnlock()
sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name })
return out
}
// StopAll 停止全部在册子进程:先并发发 plugin.stop 走优雅路径,
// 到期仍在的一律 Kill。
//
// 这是内核关停时**必须**调的:不调则子进程被 init 收养成孤儿,
// 继续持有共享段映射(段已被内核 unmap,它们下次访问就是 SIGBUS),
// 并且下次 homed 启动时同名插件会与残留进程抢同一份外部资源
// (qq 的 WS 连接、browser 的 chromium profile 锁)。
func (s *Supervisor) StopAll(timeout time.Duration) {
s.mu.Lock()
s.closed = true
procs := make([]*Process, 0, len(s.procs))
for _, p := range s.procs {
procs = append(procs, p)
}
s.mu.Unlock()
if len(procs) == 0 {
return
}
log.Printf("[proc] 关停 %d 个子进程插件", len(procs))
var wg sync.WaitGroup
for _, p := range procs {
wg.Add(1)
go func(pr *Process) {
defer wg.Done()
if err := pr.Stop(); err != nil {
log.Printf("[proc] 停止 %s: %v", pr.Name(), err)
}
}(p)
}
done := make(chan struct{})
go func() { wg.Wait(); close(done) }()
if timeout <= 0 {
timeout = stopGracePeriod * 2
}
select {
case <-done:
case <-time.After(timeout):
// 优雅停止没在预算内完成:剩下的直接 Kill。
// 不能无限等——homed 关停被单个卡住的插件拖住比杀掉它更糟。
var stuck []string
for _, p := range procs {
select {
case <-p.Exited():
default:
stuck = append(stuck, fmt.Sprintf("%s(pid=%d)", p.Name(), p.PID()))
go p.Kill()
}
}
if len(stuck) > 0 {
log.Printf("[proc] %v 内未优雅退出,强制结束: %v", timeout, stuck)
}
}
}

View File

@ -0,0 +1,71 @@
//go:build ignore
// forkplugin 在启动时 fork 一个存活时间比自己长的子进程(继承同一个 stdout),
// 然后在收到 die 工具调用时让自己退出。
//
// 用途:复现「EOF 不等于进程死亡」这一缺陷。
// 插件本体死后,孙子进程仍持有 stdout 写端,父进程(homed)的 readLoop
// 永远读不到 EOF——若内核只靠 EOF 判定死亡,就会完全感知不到插件已死:
// 工具调用一直超时、崩溃回调不触发、自动重启永不发生。
// 生产上 browser 拉 chromium、editdoc 拉 python 都是这个形状。
package main
import (
"bufio"
"encoding/json"
"os"
"os/exec"
)
type request struct {
ID uint64 `json:"id,omitempty"`
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
}
type response struct {
ID uint64 `json:"id"`
Result interface{} `json:"result,omitempty"`
Error string `json:"error,omitempty"`
}
func main() {
// 孙子进程**只**继承 stdout(本测试的要点),不给 stderr:
// 插件的 stderr 直通到 go test 的捕获管道,孙子抿着它不放会让
// go test 在测试全部通过后仍等 60s I/O。
//
// sleep 给 3s:只需在插件本体退出的那一瞬间它还持有写端即可(实际 <200ms),
// 不必拖得更久而拖慢测试。
child := exec.Command("sleep", "3")
child.Stdout = os.Stdout
_ = child.Start()
in := bufio.NewScanner(bufio.NewReader(os.Stdin))
out := bufio.NewWriter(os.Stdout)
send := func(v interface{}) {
b, _ := json.Marshal(v)
out.Write(b)
out.WriteByte('\n')
out.Flush()
}
for in.Scan() {
var req request
if err := json.Unmarshal(in.Bytes(), &req); err != nil {
continue
}
switch req.Method {
case "handshake":
send(response{ID: req.ID, Result: map[string]interface{}{
"protocol": 1, "sdk_version": "test", "plugin_name": "fork", "pid": os.Getpid(),
}})
case "tool.invoke":
// 不回应答,直接退出:模拟插件突然死亡(崩溃/被 kill)。
os.Exit(7)
default:
if req.ID != 0 {
send(response{ID: req.ID})
}
}
}
}