diff --git a/internal/plugin/evtring.go b/internal/plugin/evtring.go index 49cc331..2ce9146 100644 --- a/internal/plugin/evtring.go +++ b/internal/plugin/evtring.go @@ -18,13 +18,45 @@ type EventRing struct { ring *proc.EvtRing bus *events.Bus efd int - mu sync.Mutex + + // unsubs 保存全部已注册订阅的取消函数。 + // + // ★ 为什么必须留着:这些 handler 会 ring.WritePush(写共享内存)。 + // 而 Host.Close() 会 freeShm 解除整块映射 —— 若那时 handler 还在 Bus 上, + // 一条事件就会让 handler 写已解除映射的内存:SIGSEGV。 + // 注意 Bus.safeCall 的 recover **捕不到** SIGSEGV(它是 runtime 致命错误, + // 不是 panic),所以这不是「最坏情况只丢一条事件」,而是整个内核进程被杀。 + // + // 此前 handleEvents 把 EvtRingSubscribe 返回的取消函数直接丢弃 + // (且 EventsUnsubscribe 是空实现),于是每个订阅过的插件都在 Bus 上 + // 永久留了一个写共享内存的 handler —— 内核关停时必炸。 + // 现在改为在这里登记,由 Close 统一退订(内核关停、以及插件自己的 + // events.unsubscribe 都走这里)。 + mu sync.Mutex + unsubs []func() } func NewEventRing(ring *proc.EvtRing, efd int, bus *events.Bus) *EventRing { return &EventRing{ring: ring, bus: bus, efd: efd} } +// Close 退订本适配层注册到 Bus 的全部 handler。 +// +// 必须在 Host.Close()(munmap 共享段)**之前**调用;见 unsubs 的说明。 +// 幂等:重复调用安全(退订函数本身在 Bus 侧是「找不到就什么都不做」)。 +func (er *EventRing) Close() { + er.mu.Lock() + unsubs := er.unsubs + er.unsubs = nil + er.mu.Unlock() + + for _, fn := range unsubs { + if fn != nil { + fn() + } + } +} + // Subscribe 在 Bus 上注册一个把事件分发到事件环的 handler,返回取消函数。 // // 不改 Bus 自身结构——handler 把事件序列化后写入环并 post eventfd, @@ -53,3 +85,16 @@ func (er *EventRing) EvtRingSubscribe(types []pubsdk.EventType) func() { } } } + +// EvtRingSubscribeTracked 与 EvtRingSubscribe 相同,但把取消函数登记到 +// unsubs,供 Close 统一退订。内核的 events.subscribe 走这条。 +func (er *EventRing) EvtRingSubscribeTracked(types []pubsdk.EventType) func() { + un := er.EvtRingSubscribe(types) + if un == nil { + return nil + } + er.mu.Lock() + er.unsubs = append(er.unsubs, un) + er.mu.Unlock() + return un +} diff --git a/internal/plugin/evtring_close_test.go b/internal/plugin/evtring_close_test.go new file mode 100644 index 0000000..748fb90 --- /dev/null +++ b/internal/plugin/evtring_close_test.go @@ -0,0 +1,53 @@ +package plugin + +import ( + "testing" + + "gitcode.com/JianFeeeee/HomeAgent/internal/events" + "gitcode.com/JianFeeeee/HomeAgent/internal/plugin/proc" + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// TestEventRing_CloseUnsubscribesFromBus 钉死:EventRing.Close 必须把 +// 自己注册到 Bus 的 handler 全部撤掉。 +// +// 为什么关键:那些 handler 会 ring.WritePush —— 也就是**写共享内存**。 +// Host.Close 会 munmap 整块区域;若 handler 还挂在 Bus 上,munmap 之后 +// 任意一条事件经过 Publish 都会让它写已解除映射的内存 ⇒ SIGSEGV。 +// Bus.safeCall 虽有 recover,但 SIGSEGV 是 runtime 致命错误、recover 捕不到, +// 后果是整个 homed 进程被杀。 +// +// 判据用「Publish 之后共享内存内容是否被改动」:这是端到端的可观察后果, +// 比断言内部计数器更接近真实危害。 +func TestEventRing_CloseUnsubscribesFromBus(t *testing.T) { + host, err := proc.NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + bus := events.NewBus() + er := NewEventRing(host.EvtRing(), int(host.Evtfd().Fd()), bus) + er.EvtRingSubscribeTracked([]pubsdk.EventType{pubsdk.EventSystem}) + + // 订阅生效:Publish 一条事件应写入事件环。 + before := host.EvtRing().Written() + bus.Publish(&events.Event{Type: events.EventSystem, Source: "test", Payload: map[string]interface{}{"a": 1}}) + if host.EvtRing().Written() == before { + t.Fatal("订阅后 Publish 未写入事件环(测试前提不成立)") + } + + // 关停:退订 + er.Close() + + // 退订后 Publish 不应再写入事件环(即不再触碰共享内存)。 + after := host.EvtRing().Written() + bus.Publish(&events.Event{Type: events.EventSystem, Source: "test", Payload: map[string]interface{}{"b": 2}}) + if host.EvtRing().Written() != after { + t.Fatal("EventRing.Close 未从 Bus 退订:\n" + + " munmap 后 handler 仍会写已解除映射的内存 ⇒ SIGSEGV。") + } + + // 幂等:重复 Close 不 panic + er.Close() +} diff --git a/internal/plugin/proc/corehandler.go b/internal/plugin/proc/corehandler.go index 53fe433..6047571 100644 --- a/internal/plugin/proc/corehandler.go +++ b/internal/plugin/proc/corehandler.go @@ -56,8 +56,21 @@ type coreHandler struct { // EvtRingSubscribe 返回一个取消函数(与 Bus.Subscribe 约定一致)。 type EvtRingSubscriber interface { EvtRingSubscribe(types []pubsdk.EventType) func() + + // EvtRingSubscribeTracked 与上面相同,但订阅会被登记、可在内核关停时统一退订。 + // + // ★ 为什么需要单独的 tracked 版本:这些 handler 会写共享内存,而 + // Host.Close() 会 munmap 整块区域。若订阅不在关停前撤销,一条事件就会让 + // handler 写已解除映射的内存 ⇒ SIGSEGV(Bus.safeCall 的 recover 捕不到 + // runtime 致命错误)。详见 internal/plugin/evtring.go 的 unsubs 说明。 + EvtRingSubscribeTracked(types []pubsdk.EventType) func() } +// evtCloser 是可关闭的事件环适配层(可选实现)。 +// +// Host.Close 在 munmap 前调用它,撤掉全部写共享内存的 Bus handler。 +type evtCloser interface{ Close() } + func (h *coreHandler) invokeStageWithCtx(ctx context.Context, stage string, seq uint64) error { if h.invokeStageFn == nil { return fmt.Errorf("插件 %s: stage 调用通道未就绪", h.name) diff --git a/internal/plugin/proc/corehandler_runtime.go b/internal/plugin/proc/corehandler_runtime.go index 36f56d1..0ad1efd 100644 --- a/internal/plugin/proc/corehandler_runtime.go +++ b/internal/plugin/proc/corehandler_runtime.go @@ -100,12 +100,20 @@ func (h *coreHandler) handleEvents(method string, params json.RawMessage) (inter } // 订阅请求来自子进程——handler 直接注册到 Bus, // 事件经 EventRing 写入环后由子进程消费。 - h.evtRing.EvtRingSubscribe(p.Types) + // + // ★ 用 tracked 版本:订阅会被登记,Host.Close 在内核关停时统一退订。 + // 必须如此——这些 handler 写共享内存,而 Host.Close 会 munmap 整块区域; + // 未退订的 handler 在关停后会写已解除映射的内存 ⇒ SIGSEGV。 + h.evtRing.EvtRingSubscribeTracked(p.Types) return nil, nil case MethodEventsUnsubscribe: - // 事件环的订阅没有持久化句柄(取消函数由 Subscribe 返回但子进程未保存)。 - // 当前设计:子进程 Stop 时由内核统一清理其订阅。 + // 事件环的订阅没有**按插件**持久化句柄(取消函数由 Subscribe 返回, + // 但子进程不保存,故无法精确撤销单个插件的订阅)。 + // 当前设计:子进程 Stop 时由内核统一清理——具体落点是 + // Host.Close → evtCloser.Close 退订全部 tracked 订阅。 + // 因此这里仍是 no-op;但「统一清理」现在是真的有实现, + // 不再是只写在注释里的承诺。 return nil, nil } diff --git a/internal/plugin/proc/evtring.go b/internal/plugin/proc/evtring.go index a0ba3be..5b2208c 100644 --- a/internal/plugin/proc/evtring.go +++ b/internal/plugin/proc/evtring.go @@ -130,6 +130,12 @@ func (r *EvtRing) Init() { r.writeSeq.Store(0) } +// Written 返回已写入的事件条数(含因环满而只标记未落盘的那些)。 +// +// 供诊断与测试观测「某次 Publish 是否真的通过了 EventRing」—— +// 这比读内部字段稳定,也是关停退订验证所需的可观察量。 +func (r *EvtRing) Written() uint64 { return r.writeSeq.Load() } + // WritePush post-and-forget,**绝不阻塞**(§3.6 约束 B)。 func (r *EvtRing) WritePush(evtType pubsdk.EventType, payload []byte) { seq := r.writeSeq.Add(1) - 1 diff --git a/internal/plugin/proc/exit_order_test.go b/internal/plugin/proc/exit_order_test.go new file mode 100644 index 0000000..6577864 --- /dev/null +++ b/internal/plugin/proc/exit_order_test.go @@ -0,0 +1,78 @@ +package proc + +import ( + "sync/atomic" + "testing" + "time" +) + +// TestProcess_OnExitCompletesBeforeExitedCloses 钉死一条**顺序不变量**: +// +// Exited() 通道关闭时,onExit 回调必须**已经返回**。 +// +// ============================ 为什么这是一条安全不变量 ============================ +// onExit(内核侧 Plugin.handleExit)会调 Host.ReclaimOwner 回收残留共享槽 +// —— 那要读共享内存区域。而任何等待者(Stop/Kill/CallContext/Alive)看到 +// Exited() 关闭就会认为「完全收尾」,进而释放资源(Host.Close 会 freeShm +// 解除整块 mmap)。 +// +// 若 Exited() 先于 onExit 返回而关闭,就会出现: +// 等待者 → 释放映射 → onExit 仍在读那块内存 → SIGSEGV +// 这正是 2026-09-25 全量测试偶发崩溃的根因(栈见 process.go 的 markExited 注释)。 +// +// ============================ 为什么用原子标志而非 channel ============================ +// 要断言的是「关闭**之前**回调已完成」这一 happened-before 关系。 +// 用一个在回调里置位的原子量 + 在收到关闭信号后立刻读它: +// - 修复前:关闭先发生,回调尚未跑 ⇒ 读到 false ⇒ 判红 +// - 修复后:回调先跑完再关闭 ⇒ 读到 true ⇒ 判绿 +// 用 atomic 而非普通 bool 是为了让「回调的写」与「测试的读」之间 +// 有明确的同步语义(否则是数据竞态,-race 下会报)。 +func TestProcess_OnExitCompletesBeforeExitedCloses(t *testing.T) { + bin := buildTestPlugin(t, "crashplugin.go") + + var onExitDone atomic.Bool + exitCh := make(chan struct{}) + + p, err := Spawn("exitorder", bin, Options{ + Handler: noopHandler, + OnExit: func(name string, err error) { + // 模拟 handleExit 里的 ReclaimOwner:真实实现要读共享内存, + // 这里用一个短暂延迟把「回调还在跑」这个窗口放大到可观测。 + // 关键:置位发生在**回调返回之前**。 + time.Sleep(50 * time.Millisecond) + onExitDone.Store(true) + close(exitCh) + }, + }) + if err != nil { + t.Fatalf("Spawn: %v", err) + } + defer p.Kill() + + // 触发插件 panic 自杀 + if _, err := p.Call(MethodToolInvoke, ToolInvokeParams{Name: "boom"}); err == nil { + t.Error("崩溃插件应返回错误") + } + + // 等 Exited() 关闭 —— 此后任何等待者都会认为「可以安全 unmap」 + select { + case <-p.Exited(): + case <-time.After(10 * time.Second): + t.Fatal("10s 内未观测到进程退出") + } + + // ★ 核心断言:Exited() 已关闭时,onExit 必须已经跑完。 + if !onExitDone.Load() { + t.Fatal("顺序违例:Exited() 已关闭,但 onExit 尚未返回。\n" + + " 后果:等待者(Stop/Kill/Host.Close)会立刻 freeShm 解除映射,\n" + + " 而 onExit 里的 ReclaimOwner 仍要读共享内存 ⇒ SIGSEGV。\n" + + " 修法:markExited 中 close(p.exited) 必须放在 onExit 之后。") + } + + // 顺带确认回调确实被调用过(而非因 bug 整个跳过) + select { + case <-exitCh: + default: + t.Fatal("onExit 未在 Exited() 关闭前完成") + } +} diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 2ccae7a..2d5239e 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -151,6 +151,19 @@ func (h *Host) Close() error { if h.sup != nil { h.sup.StopAll(0) } + + // ★ 必须在 unmap **之前**退掉事件环订阅。 + // + // 那些订阅的 handler 会 ring.WritePush(写共享内存)。若让它们留在 + // Bus 上,munmap 之后只要有一条事件经过 Publish,handler 就写已解除 + // 映射的内存 ⇒ SIGSEGV。注意 Bus.safeCall 的 recover **捕不到**它 + // (runtime 致命错误不是 panic),所以后果是整个 homed 被杀。 + // + // 顺序要求:StopAll 之后(不再有新订阅进来)、freeShm 之前。 + if c, ok := h.evtSubscriber.(evtCloser); ok && c != nil { + c.Close() + } + var firstErr error if h.data != nil { if err := freeShm(h.memfd, h.data); err != nil && firstErr == nil { diff --git a/internal/plugin/proc/host_evtunsub_test.go b/internal/plugin/proc/host_evtunsub_test.go new file mode 100644 index 0000000..3b8e2e1 --- /dev/null +++ b/internal/plugin/proc/host_evtunsub_test.go @@ -0,0 +1,84 @@ +package proc + +import ( + "sync/atomic" + "testing" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// stubSubscriber 是最小 EvtRingSubscriber 实现,用于验证「关停时统一退订」。 +// +// 为什么必须桩掉 Host.evtSubscriber:真实实现是 internal/plugin.EventRing, +// 它依赖 Bus(另一个包),在 proc 包内会造成循环依赖。这里只关心 +// Host.Close 是否**调用了** closer —— 真实实现拿到 Close 后的行为 +// 由 internal/plugin 侧测试覆盖。 +type stubSubscriber struct { + closed atomic.Bool + // subCount 记录被登记的订阅数(供断言 tracked 语义) + subCount atomic.Int32 +} + +func (s *stubSubscriber) EvtRingSubscribe(types []pubsdk.EventType) func() { + s.subCount.Add(1) + return func() {} +} + +func (s *stubSubscriber) EvtRingSubscribeTracked(types []pubsdk.EventType) func() { + s.subCount.Add(1) + return func() {} +} + +func (s *stubSubscriber) Close() { s.closed.Store(true) } + +// TestHost_CloseUnsubscribesEventRing 钉死不变量: +// +// Host.Close() 必须在 munmap 共享段**之前**退订事件环。 +// +// 为什么这是安全不变量:事件环订阅的 handler 会 ring.WritePush(写共享内存)。 +// Host.Close 会 unmap 那块内存;若订阅还在 Bus 上,munmap 后任意一条事件经过 +// Publish 都会让 handler 写已解除映射的内存 ⇒ SIGSEGV。 +// Bus.safeCall 虽有 recover,但 SIGSEGV 是 runtime 致命错误、recover 捕不到, +// 后果是整个内核进程被杀。 +// +// 修复前:handleEvents 丢弃取消函数、EventsUnsubscribe 是 no-op、 +// Host.Close 也从不停订阅 —— 每个订阅过的插件都在 Bus 上永久留了一个 +// 写共享内存的 handler,内核关停时必炸。 +func TestHost_CloseUnsubscribesEventRing(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + + sub := &stubSubscriber{} + host.SetEvtSubscriber(sub) + + // 模拟插件订阅(走 tracked 路径,与 handleEvents 一致) + host.evtSubscriber.EvtRingSubscribeTracked([]pubsdk.EventType{pubsdk.EventSystem}) + if sub.subCount.Load() != 1 { + t.Fatalf("订阅登记数 = %d, want 1", sub.subCount.Load()) + } + + if err := host.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + + if !sub.closed.Load() { + t.Fatal("Host.Close 未退订事件环:\n" + + " munmap 之后 handler 仍挂在 Bus 上,一条事件就会写已解除映射的内存\n" + + " ⇒ SIGSEGV(recover 捕不到,内核进程被杀)。\n" + + " 修法:Host.Close 在 freeShm 之前调用 evtCloser.Close。") + } +} + +// TestHost_CloseWithoutSubscriber 确认没设订阅时 Close 不 panic +// (evtSubscriber 为 nil 是合法状态:未注册任何事件的部署)。 +func TestHost_CloseWithoutSubscriber(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + if err := host.Close(); err != nil { + t.Fatalf("无订阅者时 Close 应成功: %v", err) + } +} diff --git a/internal/plugin/proc/process.go b/internal/plugin/proc/process.go index 3fee74b..af6559d 100644 --- a/internal/plugin/proc/process.go +++ b/internal/plugin/proc/process.go @@ -391,13 +391,30 @@ func (p *Process) markExited() { ch <- &Response{Error: ErrProcessExited.Error()} } - close(p.exited) + // ★ 顺序至关重要:onExit 必须在 close(p.exited) **之前**完成。 + // + // onExit(内核侧即 Plugin.handleExit)会调 Host.ReclaimOwner 回收该插件 + // 残留的共享槽——**那是要读共享内存区域的**。而 exited 一关闭, + // Stop()/Kill() 就返回,StopAll 随即返回,调用方(Host.Close)立刻 + // freeShm 解除映射;若此刻 onExit 还没跑完,ReclaimOwner 就成了读 + // 已 munmap 的内存 —— SIGSEGV(recover 捕不到,直接杀进程)。 + // + // 实测崩溃栈(2026-09-25,全量 go test 偶发): + // readLoop(process.go:334) → markExited → once.Do + // → onExit → handleExit → Host.ReclaimOwner + // → arenaRegion.ReclaimOwner → blockBase → getU32 → SIGSEGV + // + // 因此 exited 的语义是「**完全**收尾完毕」,而不是「进程已死」: + // 任何等待者(Stop/Kill/CallContext/Alive)在它关闭后都可以安全地 + // 释放共享内存、卸载资源。 if p.sup != nil { p.sup.untrack(p.name) } if p.onExit != nil { p.onExit(p.name, p.ExitError()) } + + close(p.exited) }) } diff --git a/internal/plugin/proc/supervisor.go b/internal/plugin/proc/supervisor.go index 9936cac..c697ede 100644 --- a/internal/plugin/proc/supervisor.go +++ b/internal/plugin/proc/supervisor.go @@ -148,16 +148,27 @@ func (s *Supervisor) StopAll(timeout time.Duration) { // 优雅停止没在预算内完成:剩下的直接 Kill。 // 不能无限等——homed 关停被单个卡住的插件拖住比杀掉它更糟。 var stuck []string + var killers sync.WaitGroup for _, p := range procs { select { case <-p.Exited(): default: stuck = append(stuck, fmt.Sprintf("%s(pid=%d)", p.Name(), p.PID())) - go p.Kill() + // ★ 必须等 Kill 完成,不能发射后不管。 + // 本函数返回后调用方(Host.Close)立刻 freeShm 解除映射, + // 而 Kill 内部要等 markExited 跑完(含 onExit → ReclaimOwner, + // 那是要读共享内存的)。不等就 unmap ⇒ SIGSEGV。 + // Kill 自带 killReapTimeout 上限,不会无限拖住关停。 + killers.Add(1) + go func(pr *Process) { + defer killers.Done() + _ = pr.Kill() + }(p) } } if len(stuck) > 0 { log.Printf("[proc] %v 内未优雅退出,强制结束: %v", timeout, stuck) } + killers.Wait() } }