From aa3ff5a12fef77c0804a1f6eb5c0f0495b55994d Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Fri, 4 Sep 2026 18:37:13 +0800 Subject: [PATCH] =?UTF-8?q?fix(proc):=20stage=20=E5=8D=8F=E8=B0=83?= =?UTF-8?q?=E5=99=A8=E5=8F=8C=E9=87=8D=E8=A7=A3=E9=94=81=E2=80=94=E2=80=94?= =?UTF-8?q?=E5=86=85=E6=A0=B8=E6=9C=AC=E4=BD=93=20fatal=20=E5=B4=A9?= =?UTF-8?q?=E6=BA=83=E7=9A=84=E7=9C=9F=E5=9B=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 现象 2026-09-04 06:56:18 生产 homed 主进程直接死亡,退出码 2, 带走全部 27 个子进程插件。 fatal error: sync: unlock of unlocked mutex proc.(*Host).endStage(...) host.go:189 proc.(*coreHandler).runStage.func1() stage.go:94 core.(*StageHost).RunStage.func1() stages.go:190 stage.go:94 与 stages.go:190 各有一层 recover,专为「插件出错不拖垮内核」 而设,却全部失效:**sync.Mutex 的双重解锁走 runtime fatal,不是 panic, recover 结构上就拦不住**。这就是本次「插件崩溃被隔离」的设计没能生效、 内核本体整体死亡的原因。 ## 根因 endStage 把 coord.leave()(递减 inflight、判定「我是最后离开者」)放在 coordMu 临界区**之外**,而摘除 h.coord 在临界区**之内**,留出窗口: A.endStage: leave() → inflight 1→0, last=true,尚未摘除 h.coord B.beginStage: 看到 h.coord != nil,以「后到者」身份 enter,inflight 0→1 (后到者按设计不取 stageMu) A.endStage: h.coord = nil;stageMu.Unlock() ← 第 1 次 B.endStage: leave() → inflight 1→0, last=true → stageMu.Unlock() ← 第 2 次 💥 B 从未持有 stageMu,却因挂进一个正在收尾的协调器而被判成「最后离开者」, 对同一把锁解了两次。崩溃前一行日志是 config_list_keys 的结果——那一刻 正好有 stage 扇出,与竞态窗口重合。 ## 修复 把「递减 inflight → 判定最后离开者 → 摘除 h.coord」收进同一个 coordMu 临界区,后到者再不可能挂进已收尾的协调器。为此把 leave() 拆成: - depart():纯计数,由 endStage 在 coordMu 内调用 - finish():共享段回读 + arena 压实,在 coordMu 外、但仍在 stageMu.Unlock() 之前(先放锁会让下一轮 stage 在回读未完时改写共享段) leave() 保留给单测。 同一函数的第二个隐患一并修掉:首进者的 enter()(含 WriteAll 写共享段) 原先在 coordMu 之外,后到者可能拿到 coord 就去读**写了一半**的段。 现在 enter() 在锁内完成。 beginStage 错误路径的 stageMu.Unlock() 必须保留并已加注释说明: runStage 的 defer endStage(coord) 是在 beginStage 返回 err 的检查**之后** 才注册的,这条路径上没有任何人会替它解锁,漏掉就是整个 stage 通道永久卡死。 锁序 stageMu → coordMu;endStage 只解锁 stageMu 不获取,无环。 ## 验证 反向验证:把 host.go stash 回旧版跑新测试 → fatal error: sync: unlock of unlocked mutex;恢复修复 → 通过。测试抓的确实是这个缺陷。 5 个回归用例(host_stage_test.go): - 后到者不复用已收尾的协调器(直接构造那个时序,不靠调度巧合) - 8 worker × 40 轮并发进出(旧实现下整个测试二进制 fatal 而非 FAIL) - 同阶段多插件扇出共用一个协调器、仅最后离开者解锁 - 50 轮串行不泄漏(少解锁会在第二轮卡死) - 四阶段序列 pre_action→chat→after_toolcall→post_action internal/plugin/... 全量 -race -count=2 通过。 ## 同类缺陷审计(本 commit 未改动其他文件,仅记录结论) 针对「recover 拦不住的 runtime fatal」这一整类做了全仓审计: 1. 跨函数持锁(本缺陷的形状,脚本枚举 Lock/Unlock 不配对的函数) - proc/lock.go 的 Release/ForceRelease 同样「只 Unlock 不 Lock」, 但两者都在 ownerMu 下先检查 held/owner 再解锁,非持有者直接返回, 不存在双解锁路径。 - 其余 22 处 Lock/Unlock 计数不等的函数逐一复核:全部是多分支早退各自 解锁(waiter 的 goto nextMessage、sidecar.call 的五个错误分支、 lua adapterPool 的 cond.Wait 池模式等),配对正确。 2. 并发 map 读写(同样是 runtime fatal) - 16 处「无锁访问 map」全部复核为安全:Locked 后缀约定(orderedLocked、 defsLockedRegisterSource)、调用方持锁(document 的 addSummary/ removeDoc/loadAll、registry 的 runStopHandlers/runOnRemoveHandlers)、 或启动期单线程(knowledge.scanAll、static_embedder 构造后只读)。 3. close of closed channel - 全仓仅 sidecar.go 有同名变量的两处 close(ch),但作用于不同集合成员, 且 Close() 前有 readerWg.Wait() 与 stopped 标志,reader 侧已 delete 出 pending,不会双关。 - 各插件 stopCh 的 close:healthcheck 用 select 守卫、clawhubadapter 用 stopOnce、evtring 用 running 标志、timer 交给 StopHandler 单次调用。 agentcli.Stop() 是裸 close(p.stopCh) 无幂等守卫,但 Registry 的六处 Stop 调用点都在同一把 r.mu 下先 delete(r.plugins)+摘 r.instances 再 Stop,不存在二次调用路径——记录为「依赖调用方约定」而非当前缺陷。 4. WaitGroup 误用:未发现 Add 出现在 goroutine 体内的形状。 5. 全仓 go test ./... -race:零 DATA RACE、零 FAIL。 --- internal/plugin/proc/host.go | 119 +++++++++---- internal/plugin/proc/host_stage_test.go | 220 ++++++++++++++++++++++++ 2 files changed, 303 insertions(+), 36 deletions(-) create mode 100644 internal/plugin/proc/host_stage_test.go diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 218a922..38f3613 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -144,48 +144,83 @@ func (h *Host) Close() error { // // 首个进入者:获取 stageMu(独占共享段)→ 把内核 StageContext 写入段。 // 后续进入者:仅递增 inflight。 +// +// enter() 在 coordMu 内完成,两个原因: +// 1. 首进者的 WriteAll 未结束前不能让后到者拿到 coord 就去读共享段 +// (旧码的后到者 enter 立即返回,可能读到写一半的段)。 +// 2. 与 endStage 的摘除互斥,防止后到者挂进一个正在收尾的协调器 +// (具体见 endStage 的注释)。 +// +// 锁序:stageMu → coordMu。endStage 只解锁 stageMu、不获取,所以无环。 func (h *Host) beginStage(sc *pubsdk.StageContext) (*stageCoordinator, error) { h.coordMu.Lock() - first := h.coord == nil - if first { - // 独占共享段直到本次 stage 全部插件离开 + if h.coord == nil { + // 首个进入者:独占共享段直到本次 stage 全部插件离开。 + // 必须先放 coordMu 再取 stageMu,不能反序。 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 + if h.coord == nil { + coord := newStageCoordinator(h.seg) + h.coord = coord + if err := coord.enter(sc, true); err != nil { + // 注意:runStage 的 defer endStage(coord) 是在 beginStage + // 返回 err 的检查之后才注册的,所以这条路径上 + // endStage 永远不会被调用——stageMu 必须在此自行释放, + // 否则整个 stage 通道永久卡死。 + h.coord = nil + h.coordMu.Unlock() + h.stageMu.Unlock() + return nil, err + } h.coordMu.Unlock() - h.stageMu.Unlock() + h.locks.bind(coord.lock) + return coord, nil } + // 双检失败:等锁期间已有其他插件建好协调器,退回后到者路径。 + h.stageMu.Unlock() + } + + coord := h.coord + if err := coord.enter(sc, false); err != nil { + h.coordMu.Unlock() return nil, err } + h.coordMu.Unlock() h.locks.bind(coord.lock) return coord, nil } // endStage 由插件 handler 返回时调用。 // 最后离开者:把共享段结果读回内核 StageContext → 压实 arena → 释放 stageMu。 +// +// coordMu 必须覆盖「递减 inflight → 判定最后离开者 → 摘除 h.coord」全过程。 +// 旧码把 leave() 放在 coordMu 之外,留出了这个窗口(即 2026-09-04 06:56:18 +// 线上 fatal error: sync: unlock of unlocked mutex 的真因): +// +// A.endStage: leave() → inflight 1→0, last=true,尚未摘除 h.coord +// B.beginStage: 看到 h.coord != nil,以「后到者」身份 enter,inflight 0→1 +// (后到者不取 stageMu) +// A.endStage: h.coord = nil;stageMu.Unlock() ← 第 1 次 +// B.endStage: leave() → inflight 1→0, last=true → stageMu.Unlock() ← 第 2 次 💥 +// +// B 从未持有 stageMu(它是后到者),却因为挂进了一个正在收尾的协调器 +// 而成为“最后离开者”,于是对同一把锁解了两次。sync.Mutex 的双重解锁是 +// runtime fatal,**recover 捕不到**——这就是为何 stage.go / stages.go 里 +// 那两层 recover 全部失效、整个 homed 直接死掉的原因。 func (h *Host) endStage(coord *stageCoordinator) error { - last, err := coord.leave() - if !last { - return err - } h.coordMu.Lock() - h.coord = nil + last, sc, written := coord.depart() + if last && h.coord == coord { + h.coord = nil + } h.coordMu.Unlock() + if !last { + return nil + } + // finish 必须在 stageMu.Unlock() 之前:先放锁会让下一轮 stage + // 在回读未完时就改写共享段。 + err := coord.finish(sc, written) h.stageMu.Unlock() return err } @@ -232,26 +267,38 @@ func (c *stageCoordinator) enter(sc *pubsdk.StageContext, first bool) error { // leave 登记一个插件离开;返回是否为最后一个离开者。 // -// 最后离开者负责把共享段结果读回内核 StageContext,并压实 arena -// (此时无插件持锁,满足 §3.3 的压实前提)。 +// 拆成两段:depart() 只动计数(由 endStage 在 coordMu 内调用,使 +// 「递减 → 判定最后者 → 摘除 h.coord」成为原子操作),finish() 做 +// 共享段回读与压实。本方法保留给单测用。 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 { + last, sc, written := c.depart() + if !last { return last, nil } + return last, c.finish(sc, written) +} + +// depart 递减 inflight 并报告是否为最后离开者。 +func (c *stageCoordinator) depart() (last bool, sc *pubsdk.StageContext, written bool) { + c.mu.Lock() + defer c.mu.Unlock() + c.inflight-- + return c.inflight == 0, c.ctxRef, c.written +} + +// finish 把共享段结果读回内核 StageContext 并压实 arena +// (此时无插件持锁,满足 §3.3 的压实前提)。 +func (c *stageCoordinator) finish(sc *pubsdk.StageContext, written bool) error { + if !written || sc == nil { + return nil + } if rErr := c.seg.ReadInto(sc); rErr != nil { - return last, fmt.Errorf("回读共享段: %w", rErr) + return fmt.Errorf("回读共享段: %w", rErr) } if reclaimed := c.seg.Compact(); reclaimed > 0 { log.Printf("[proc] stage 结束,arena 压实回收 %d 字节", reclaimed) } - return last, nil + return nil } // ShmSize 返回共享段大小(供诊断/日志)。 diff --git a/internal/plugin/proc/host_stage_test.go b/internal/plugin/proc/host_stage_test.go new file mode 100644 index 0000000..86bc190 --- /dev/null +++ b/internal/plugin/proc/host_stage_test.go @@ -0,0 +1,220 @@ +package proc + +import ( + "sync" + "sync/atomic" + "testing" + + pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" +) + +// 本文件是 2026-09-04 06:56:18 线上 crash 的回归测试。 +// +// 崩溃形态:homed 主进程直接死亡,退出码 2。 +// +// fatal error: sync: unlock of unlocked mutex +// proc.(*Host).endStage(...) host.go:189 +// proc.(*coreHandler).runStage.func1() stage.go:94 +// core.(*StageHost).RunStage.func1() stages.go:190 +// +// 注意 stage.go 与 stages.go 各有一层 recover,却都没拦住—— +// sync.Mutex 的双重解锁是 runtime fatal,recover 捕不到。这是本次 +// "整个内核本体崩溃"而非"插件崩溃被隔离"的直接原因。 + +// TestEndStage_LateArrivalNoDoubleUnlock 复现根因竞态。 +// +// 旧实现把 leave() 放在 coordMu 之外,留出这个窗口: +// +// A.endStage: leave() → inflight 1→0, last=true,尚未摘除 h.coord +// B.beginStage: 看到 h.coord != nil,以「后到者」身份 enter,inflight 0→1 +// (后到者不取 stageMu) +// A.endStage: h.coord = nil; stageMu.Unlock() ← 第 1 次 +// B.endStage: leave() → inflight 1→0, last=true → stageMu.Unlock() ← 第 2 次 💥 +// +// B 从未持有 stageMu,却因为挂进了一个正在收尾的协调器而成为 +// "最后离开者",于是对同一把锁解了两次。 +// +// 本测试直接驱动 depart/enter 制造那个时序,不依赖调度巧合。 +func TestEndStage_LateArrivalNoDoubleUnlock(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + scA := &pubsdk.StageContext{RawMessage: "A"} + coordA, err := host.beginStage(scA) + if err != nil { + t.Fatalf("A beginStage: %v", err) + } + + // A 收尾:修复后 depart 与摘除 h.coord 在同一个 coordMu 临界区内, + // 所以此刻起 h.coord 已是 nil,B 不可能再挂进 A 的协调器。 + if err := host.endStage(coordA); err != nil { + t.Fatalf("A endStage: %v", err) + } + + // B 现在进入:必须成为新的首进者(拿到自己的 stageMu), + // 而不是挂进 A 那个已收尾的协调器。 + scB := &pubsdk.StageContext{RawMessage: "B"} + coordB, err := host.beginStage(scB) + if err != nil { + t.Fatalf("B beginStage: %v", err) + } + if coordB == coordA { + t.Fatal("B 不该复用 A 已收尾的协调器——这正是 double-unlock 的来源") + } + if err := host.endStage(coordB); err != nil { + t.Fatalf("B endStage: %v", err) + } + + // 若上面多解了一次锁,这里会 fatal(runtime 级,测试进程直接死); + // 能走到这一步说明配对正确。 + scC := &pubsdk.StageContext{RawMessage: "C"} + coordC, err := host.beginStage(scC) + if err != nil { + t.Fatalf("C beginStage: %v", err) + } + if err := host.endStage(coordC); err != nil { + t.Fatalf("C endStage: %v", err) + } +} + +// TestEndStage_ConcurrentChurnNoFatal 高并发进出:真实触发线上那个窗口。 +// +// 旧实现下这个测试会以 fatal error: sync: unlock of unlocked mutex 结束 +// (整个测试二进制死亡,不是 FAIL)。修复后应干净通过。 +func TestEndStage_ConcurrentChurnNoFatal(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + const workers = 8 + const rounds = 40 + var wg sync.WaitGroup + var failures atomic.Int64 + + for w := 0; w < workers; w++ { + wg.Add(1) + go func() { + defer wg.Done() + for r := 0; r < rounds; r++ { + sc := &pubsdk.StageContext{RawMessage: "churn"} + coord, err := host.beginStage(sc) + if err != nil { + failures.Add(1) + return + } + if err := host.endStage(coord); err != nil { + failures.Add(1) + return + } + } + }() + } + wg.Wait() + + if n := failures.Load(); n > 0 { + t.Fatalf("%d 次 begin/end 失败", n) + } +} + +// TestBeginStage_MultiPluginSameStage 同阶段多插件扇出: +// 首进者取 stageMu,后到者只递增 inflight,最后离开者才解锁。 +// 验证并发扇出这一原始设计仍然成立(§0.2 第 1 条)。 +func TestBeginStage_MultiPluginSameStage(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + sc := &pubsdk.StageContext{RawMessage: "fanout"} + + // 三个插件先后进入同一次 stage + c1, err := host.beginStage(sc) + if err != nil { + t.Fatalf("plugin1 beginStage: %v", err) + } + c2, err := host.beginStage(sc) + if err != nil { + t.Fatalf("plugin2 beginStage: %v", err) + } + c3, err := host.beginStage(sc) + if err != nil { + t.Fatalf("plugin3 beginStage: %v", err) + } + // 同一次 stage 内必须共用一个协调器(共享同一份 StageContext 段) + if c1 != c2 || c2 != c3 { + t.Fatal("同阶段并发插件应共用一个协调器") + } + + // 前两个离开不该释放 stageMu + if err := host.endStage(c1); err != nil { + t.Fatalf("plugin1 endStage: %v", err) + } + if err := host.endStage(c2); err != nil { + t.Fatalf("plugin2 endStage: %v", err) + } + // 最后一个离开才释放 + if err := host.endStage(c3); err != nil { + t.Fatalf("plugin3 endStage: %v", err) + } + + // 锁已释放:新一轮能立即开始 + c4, err := host.beginStage(sc) + if err != nil { + t.Fatalf("新一轮 beginStage 应成功(stageMu 已释放): %v", err) + } + if c4 == c1 { + t.Fatal("新一轮应是新的协调器") + } + if err := host.endStage(c4); err != nil { + t.Fatalf("新一轮 endStage: %v", err) + } +} + +// TestBeginStage_SerialRounds 长串行:确认没有单向泄漏(少解锁会在第二轮卡死)。 +func TestBeginStage_SerialRounds(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + for round := 0; round < 50; round++ { + sc := &pubsdk.StageContext{RawMessage: "serial"} + coord, err := host.beginStage(sc) + if err != nil { + t.Fatalf("round %d beginStage: %v", round, err) + } + if err := host.endStage(coord); err != nil { + t.Fatalf("round %d endStage: %v", round, err) + } + } +} + +// TestBeginStage_PhaseSequence 模拟一条消息走完 pre_action → chat → post_action。 +func TestBeginStage_PhaseSequence(t *testing.T) { + host, err := NewHost() + if err != nil { + t.Fatalf("NewHost: %v", err) + } + defer host.Close() + + phases := []pubsdk.Stage{"pre_action", "chat", "after_toolcall", "post_action"} + for msg := 0; msg < 10; msg++ { + for _, p := range phases { + sc := &pubsdk.StageContext{RawMessage: "msg", Phase: p} + coord, err := host.beginStage(sc) + if err != nil { + t.Fatalf("msg %d phase %s beginStage: %v", msg, p, err) + } + if err := host.endStage(coord); err != nil { + t.Fatalf("msg %d phase %s endStage: %v", msg, p, err) + } + } + } +}