mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-22 18:08:04 +00:00
根因:子进程插件被 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 双摘后可重新注册同名)
281 lines
8.8 KiB
Go
281 lines
8.8 KiB
Go
package proc
|
||
|
||
import (
|
||
"fmt"
|
||
"log"
|
||
"os"
|
||
"sync"
|
||
|
||
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
|
||
)
|
||
|
||
// Host 持有**被全部子进程插件共享的一块 StageContext 段**,是共享内存数据面的
|
||
// 所有权中心(§3.3/§3.4)。
|
||
//
|
||
// ❗ 为什么必须共享一块段(这是一个容易走错的关键点):
|
||
// 若每个插件各持一块段,则「内核 ctx → 段 → 插件改 → 回读 ctx」在多插件下退化成
|
||
// 副本模型——两个插件各写各的段、各自回读,最后回读者覆盖前者,
|
||
// lost update 原样复现(§8.4 实测 35.8~36.8%)。
|
||
// 实验 8 的做法是 5 个 worker 进程 mmap **同一个 memfd**,本实现与之一致。
|
||
//
|
||
// 生命周期:Host 由 registry 创建一次,随内核存活;每个插件 spawn 时经
|
||
// ExtraFiles 拿到同一 memfd(fd 3),mmap 后即看到同一份物理页。
|
||
//
|
||
// 另外持有事件环段(§3.6):独立于 StageContext 的事件通知通道,
|
||
// 子进程从 eventfd 感知新事件并从 mmap 读 slot。
|
||
// fd 分配:fd 3 = StageContext,fd 4 = 事件环,fd 5 = eventfd。
|
||
type Host struct {
|
||
memfd *os.File
|
||
data []byte
|
||
seg *Segment
|
||
shmSize int
|
||
|
||
// 事件环段(独立于 StageContext)
|
||
evtfd *os.File // Unix:eventfd/pipe 读端(fd 5)。Windows 为 nil,用 evtNotifyFd 。
|
||
evtNotifyFd int // 通知句柄的平台无关标识(Unix 是真 fd,Windows 是伪 fd)
|
||
evtRing *EvtRing // 内核侧事件环句柄
|
||
evtRingFd *os.File // Unix:事件环段 memfd(fd 4)。Windows 为 nil(命名段)。
|
||
evtData []byte // 事件环段 mmap 数据
|
||
|
||
// evtSubscriber 由 internal/plugin 注入,coreHandler 用它接子进程的 events.subscribe 请求。
|
||
// proc 包不依赖 internal/plugin(循环依赖),故用接口类型存储。
|
||
evtSubscriber EvtRingSubscriber
|
||
|
||
locks *lockRegistry
|
||
stageMu sync.Mutex
|
||
coordMu sync.Mutex
|
||
coord *stageCoordinator
|
||
|
||
// sup 是内核侧唯一的子进程台账,与共享段同生命周期。
|
||
//
|
||
// 放在 Host 而不是 registry 的理由:能拿到 Host 的地方就能拿到台账,
|
||
// 而 Host 本就是「全部子进程插件共享的那一份内核侧状态」。
|
||
sup *Supervisor
|
||
}
|
||
|
||
// NewHost 创建共享段(平台层 allocShm + 布局初始化)。
|
||
//
|
||
// 段的**传递机制**按平台分开(shmalloc_*.go),但**布局**完全一致:
|
||
// - Linux:memfd,经 ExtraFiles 传继承 fd
|
||
// - macOS:立即 unlink 的临时文件(无 memfd_create),同样走 fd 继承
|
||
// - Windows:命名 FileMapping(无 fd 继承语义),插件按名字打开
|
||
//
|
||
// 三者共同点:全部插件看到同一份物理页,段内一律用相对偏移而非指针
|
||
// (实验 2 已验证各进程 mmap 到不同虚拟地址时偏移解引用仍正确)。
|
||
func NewHost() (*Host, error) {
|
||
memfd, data, err := allocShm(shmDefaultSize)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
seg, err := NewSegment(data)
|
||
if err != nil {
|
||
freeShm(memfd, data)
|
||
return nil, err
|
||
}
|
||
|
||
// 创建事件环段(独立于 StageContext)
|
||
evtRingFd, evtData, efd, err := allocEvtRing()
|
||
if err != nil {
|
||
freeShm(memfd, data)
|
||
return nil, fmt.Errorf("事件环: %w", err)
|
||
}
|
||
evtRing, err := NewEvtRing(evtData)
|
||
if err != nil {
|
||
freeShm(memfd, data)
|
||
return nil, fmt.Errorf("事件环初始化: %w", err)
|
||
}
|
||
evtRing.Init()
|
||
|
||
return &Host{
|
||
sup: NewSupervisor(),
|
||
memfd: memfd,
|
||
data: data,
|
||
seg: seg,
|
||
shmSize: shmDefaultSize,
|
||
evtfd: evtfdReadFile(efd),
|
||
evtNotifyFd: efd,
|
||
evtRing: evtRing,
|
||
evtRingFd: evtRingFd,
|
||
evtData: evtData,
|
||
locks: &lockRegistry{},
|
||
}, nil
|
||
}
|
||
|
||
// shmDefaultSize 是共享 StageContext 段的大小。
|
||
//
|
||
// 取 256KB:StageContext 全字段 JSON 化后典型 < 4KB(工具结果中位 93B,§2.5),
|
||
// append-only 中间垃圾由 stage 结束时 Compact 回收,256KB 给足余量。
|
||
// 全部插件共享一块,总开销恒定,不随插件数增长。
|
||
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 {
|
||
firstErr = err
|
||
}
|
||
h.data, h.memfd = nil, nil
|
||
}
|
||
if h.evtData != nil {
|
||
if h.evtRingFd != nil {
|
||
h.evtRingFd.Close()
|
||
h.evtRingFd = nil
|
||
}
|
||
h.evtData = nil
|
||
}
|
||
if h.evtfd != nil {
|
||
h.evtfd.Close()
|
||
h.evtfd = nil
|
||
}
|
||
evtfdClose(h.evtNotifyFd)
|
||
return firstErr
|
||
}
|
||
|
||
// beginStage 由插件 handler 进入时调用。
|
||
//
|
||
// 首个进入者:获取 stageMu(独占共享段)→ 把内核 StageContext 写入段。
|
||
// 后续进入者:仅递增 inflight。
|
||
func (h *Host) beginStage(sc *pubsdk.StageContext) (*stageCoordinator, error) {
|
||
h.coordMu.Lock()
|
||
first := h.coord == nil
|
||
if first {
|
||
// 独占共享段直到本次 stage 全部插件离开
|
||
h.coordMu.Unlock()
|
||
h.stageMu.Lock()
|
||
h.coordMu.Lock()
|
||
// 双检:等锁期间可能已有其他插件建好协调器(它们会先拿到 stageMu)
|
||
if h.coord != nil {
|
||
first = false
|
||
h.stageMu.Unlock()
|
||
} else {
|
||
h.coord = newStageCoordinator(h.seg)
|
||
}
|
||
}
|
||
coord := h.coord
|
||
h.coordMu.Unlock()
|
||
|
||
if err := coord.enter(sc, first); err != nil {
|
||
if first {
|
||
h.coordMu.Lock()
|
||
h.coord = nil
|
||
h.coordMu.Unlock()
|
||
h.stageMu.Unlock()
|
||
}
|
||
return nil, err
|
||
}
|
||
h.locks.bind(coord.lock)
|
||
return coord, nil
|
||
}
|
||
|
||
// endStage 由插件 handler 返回时调用。
|
||
// 最后离开者:把共享段结果读回内核 StageContext → 压实 arena → 释放 stageMu。
|
||
func (h *Host) endStage(coord *stageCoordinator) error {
|
||
last, err := coord.leave()
|
||
if !last {
|
||
return err
|
||
}
|
||
h.coordMu.Lock()
|
||
h.coord = nil
|
||
h.coordMu.Unlock()
|
||
h.stageMu.Unlock()
|
||
return err
|
||
}
|
||
|
||
// ForceReleaseLock 在插件进程崩溃时释放其可能持有的 stage 锁(实验 9 自愈机制)。
|
||
func (h *Host) ForceReleaseLock(plugin string) bool {
|
||
return h.locks.forceRelease(plugin)
|
||
}
|
||
|
||
// Segment 暴露共享段(供诊断与测试)。
|
||
func (h *Host) Segment() *Segment { return h.seg }
|
||
|
||
// stageCoordinator 跟踪一次 stage 执行中参与插件的进出。
|
||
type stageCoordinator struct {
|
||
seg *Segment
|
||
lock *stageLock
|
||
|
||
mu sync.Mutex
|
||
inflight int
|
||
written bool
|
||
ctxRef *pubsdk.StageContext
|
||
}
|
||
|
||
func newStageCoordinator(seg *Segment) *stageCoordinator {
|
||
return &stageCoordinator{seg: seg, lock: newStageLock()}
|
||
}
|
||
|
||
// enter 登记一个插件进入本次 stage;first 为真时把内核状态写入共享段。
|
||
func (c *stageCoordinator) enter(sc *pubsdk.StageContext, first bool) error {
|
||
c.mu.Lock()
|
||
defer c.mu.Unlock()
|
||
c.inflight++
|
||
if !first || c.written {
|
||
return nil
|
||
}
|
||
c.ctxRef = sc
|
||
if err := c.seg.WriteAll(sc); err != nil {
|
||
c.inflight--
|
||
return fmt.Errorf("写入共享段: %w", err)
|
||
}
|
||
c.written = true
|
||
return nil
|
||
}
|
||
|
||
// leave 登记一个插件离开;返回是否为最后一个离开者。
|
||
//
|
||
// 最后离开者负责把共享段结果读回内核 StageContext,并压实 arena
|
||
// (此时无插件持锁,满足 §3.3 的压实前提)。
|
||
func (c *stageCoordinator) leave() (last bool, err error) {
|
||
c.mu.Lock()
|
||
c.inflight--
|
||
last = c.inflight == 0
|
||
sc := c.ctxRef
|
||
written := c.written
|
||
c.mu.Unlock()
|
||
|
||
if !last || !written || sc == nil {
|
||
return last, nil
|
||
}
|
||
if rErr := c.seg.ReadInto(sc); rErr != nil {
|
||
return last, fmt.Errorf("回读共享段: %w", rErr)
|
||
}
|
||
if reclaimed := c.seg.Compact(); reclaimed > 0 {
|
||
log.Printf("[proc] stage 结束,arena 压实回收 %d 字节", reclaimed)
|
||
}
|
||
return last, nil
|
||
}
|
||
|
||
// ShmSize 返回共享段大小(供诊断/日志)。
|
||
func (h *Host) ShmSize() int { return h.shmSize }
|
||
|
||
// EvtRing 返回内核侧事件环句柄。
|
||
func (h *Host) EvtRing() *EvtRing { return h.evtRing }
|
||
|
||
// Evtfd 返回通知读端的 *os.File(Unix;eventfd/pipe)。
|
||
// Windows 返回 nil——命名 Event 不是文件句柄,用 EvtNotifyFd 代替。
|
||
func (h *Host) Evtfd() *os.File { return h.evtfd }
|
||
|
||
// EvtNotifyFd 返回通知句柄的平台无关标识,供 EventRing 写通知。
|
||
//
|
||
// Unix 是真 fd;Windows 是映射到命名 Event 句柄的伪 fd。
|
||
// EvtfdNotify 接受这个值并按平台分派。
|
||
func (h *Host) EvtNotifyFd() int { return h.evtNotifyFd }
|
||
|
||
// SetEvtSubscriber 注入事件环订阅接口(由 Registry 在创建 Host 后设置)。
|
||
func (h *Host) SetEvtSubscriber(sub EvtRingSubscriber) { h.evtSubscriber = sub }
|
||
|
||
// EvtData 返回事件环段 mmap 数据(子进程消费者用)。
|
||
func (h *Host) EvtData() []byte { return h.evtData }
|
||
|
||
// EvtfdReadFile 返回 eventfd 的 *os.File(供子进程读取消费)。
|
||
func (h *Host) EvtfdReadFile() *os.File { return h.evtfd }
|