mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 01:18:08 +00:00
fix(proc): stage 协调器双重解锁——内核本体 fatal 崩溃的真因
## 现象
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。
This commit is contained in:
@ -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 返回共享段大小(供诊断/日志)。
|
||||
|
||||
220
internal/plugin/proc/host_stage_test.go
Normal file
220
internal/plugin/proc/host_stage_test.go
Normal file
@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user