mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 09:28:14 +00:00
fix(lua): events.subscribe 改用内部 Subscribe + 订阅生命周期(修死锁/use-after-close)
上一版 Lua 对齐引入的 sdk.events.subscribe 有两个真问题,本提交修掉: 1) 用了公共 SDK 的 Events(),但本内核从未注入 event subscriber (SetEventSubscriber 全仓无调用点),拿到永远是 nil ⇒ subscribe 只会 返回 "events unavailable"。改用内部 SDK 的 s.Subscribe——内置插件走的就是 这条路径(cli/webui/skillmgr 全用它)。 2) 自死锁:subscribe 会在 Lua 的 plugin.start(sdk) 回调里被调用,而 luaPlugin.Start 正持有 p.mu;原实现在 subscribe 里再 lock p.mu 追加 subs, 不可重入 ⇒ 测试实测 30s 超时。改用独立的 subsMu。 3) use-after-close:Stop 会 Close LState,但事件订阅此前无人取消,残留回调 再触发就会碰已关的 L。现在:Stop 先(不持 p.mu,避免与 Bus.Publish 锁序反转)取 subsMu 取消全部订阅,再置 closed 并关 L;事件回调持 p.mu 后 先查 closed,已进入等锁的旧回调会直接返回。 4) plugin_mgr 访问补 nil 保护(部分单测构造的 SDK 不含 pluginMgr)。 回归:TestLuaEventsSubscribeAndStopCleanup——订阅后 Publish 命中、Stop 后 再 Publish 不 panic。全套 Lua 测试在 -race 下通过。
This commit is contained in:
@ -10,9 +10,9 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
agentEvents "gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
||||
luaSDK "gitcode.com/JianFeeeee/HomeAgent/internal/lua/sdk"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
|
||||
lua "github.com/yuin/gopher-lua"
|
||||
)
|
||||
|
||||
@ -42,7 +42,16 @@ type luaPlugin struct {
|
||||
stages map[sdk.Stage]*stageReg
|
||||
outputChs map[string]*outputChReg
|
||||
inputDefs map[string]sdk.ChannelDef
|
||||
mu sync.Mutex
|
||||
// subs 是本插件注册的事件订阅取消函数;Stop 时兜底取消,
|
||||
// 避免 L 已 Close 后残留回调被触发(use-after-close)。
|
||||
// 用独立的 subsMu 而非 mu:subscribe 会在 Lua 的 start 回调里被调,
|
||||
// 而 Start 正持着 mu —— 用 mu 就是不可重入的自死锁。
|
||||
subs []func()
|
||||
subsMu sync.Mutex
|
||||
// closed 在 Stop 里置位(持 mu);事件回调持 mu 后先查它,
|
||||
// 防止“回调已通过取消订阅检查、但等锁期间 L 被 Close”的竞态。
|
||||
closed bool
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
func newLuaPlugin(luaPath, name string) (*luaPlugin, error) {
|
||||
@ -791,20 +800,25 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
|
||||
return pushVal([]interface{}{})
|
||||
}))
|
||||
|
||||
// ---- sdk.events.*(只读事件订阅,与外部插件的 Events() 对齐)----
|
||||
// 回调在内核事件发布 goroutine 上执行,必须只做轻量转发(Lua 单状态 + 互斥锁);
|
||||
// 阻塞会卡死本插件的全部调用。返回一个取消订阅函数。
|
||||
// ---- sdk.events.*(只读事件订阅)----
|
||||
//
|
||||
// 用内部 SDK 的 Subscribe(内置插件用的是同一条路径);
|
||||
// 不用公共 SDK 的 Events()——那个 subscriber 在本内核里从未被注入
|
||||
// (SetEventSubscriber 无调用点),拿到的永远是 nil。
|
||||
//
|
||||
// 回调用内核事件发布 goroutine 上执行,必须只做轻量转发(Lua 单状态 + 互斥锁);
|
||||
// 阻塞会卡死本插件的全部调用。返回一个取消订阅函数,并在 Stop 时兜底取消
|
||||
// (否则插件停掉/重载后 L 已 Close,残留回调再触发就是 use-after-close)。
|
||||
evTbl := subTable("events")
|
||||
evTbl.RawSetString("subscribe", L.NewFunction(func(L *lua.LState) int {
|
||||
eventType := L.CheckString(1)
|
||||
fn := L.CheckFunction(2)
|
||||
sub := s.Events()
|
||||
if sub == nil {
|
||||
return pushErr(fmt.Errorf("events unavailable"))
|
||||
}
|
||||
unsub := sub.Subscribe(pubsdk.EventType(eventType), func(evt *pubsdk.Event) {
|
||||
unsub := s.Subscribe(agentEvents.EventType(eventType), func(evt *agentEvents.Event) {
|
||||
plg.mu.Lock()
|
||||
defer plg.mu.Unlock()
|
||||
if plg.closed {
|
||||
return
|
||||
}
|
||||
L2 := plg.L
|
||||
tbl := L2.NewTable()
|
||||
tbl.RawSetString("type", lua.LString(string(evt.Type)))
|
||||
@ -817,8 +831,11 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
|
||||
fmt.Printf("[lua-plugin/%s] event handler error: %v\n", plg.name, err)
|
||||
}
|
||||
})
|
||||
plg.subsMu.Lock()
|
||||
plg.subs = append(plg.subs, unsub)
|
||||
plg.subsMu.Unlock()
|
||||
L.Push(L.NewFunction(func(L *lua.LState) int {
|
||||
unsub()
|
||||
unsub() // 事件总线的取消订阅是幂等的(重复调用只会匹配不到)
|
||||
return 0
|
||||
}))
|
||||
L.Push(lua.LNil)
|
||||
@ -826,18 +843,31 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
|
||||
}))
|
||||
|
||||
// ---- sdk.plugin_mgr.*(插件管理,与外部插件的 PluginMgrAPI 对齐)----
|
||||
// PluginMgr 可能未装配(如部分单测的 SDK 构造),此时返回"不可用"而不是 panic。
|
||||
pmTbl := subTable("plugin_mgr")
|
||||
pmTbl.RawSetString("reload_one", L.NewFunction(func(L *lua.LState) int {
|
||||
if err := s.PluginMgr().ReloadOne(L.CheckString(1)); err != nil {
|
||||
pm := s.PluginMgr()
|
||||
if pm == nil {
|
||||
return pushErr(fmt.Errorf("plugin manager unavailable"))
|
||||
}
|
||||
if err := pm.ReloadOne(L.CheckString(1)); err != nil {
|
||||
return pushErr(err)
|
||||
}
|
||||
return pushNil()
|
||||
}))
|
||||
pmTbl.RawSetString("list_loaded", L.NewFunction(func(L *lua.LState) int {
|
||||
return pushList(s.PluginMgr().ListLoadedPlugins())
|
||||
pm := s.PluginMgr()
|
||||
if pm == nil {
|
||||
return pushList([]interface{}{})
|
||||
}
|
||||
return pushList(pm.ListLoadedPlugins())
|
||||
}))
|
||||
pmTbl.RawSetString("is_disabled", L.NewFunction(func(L *lua.LState) int {
|
||||
return pushVal(s.PluginMgr().IsPluginDisabled(L.CheckString(1)))
|
||||
pm := s.PluginMgr()
|
||||
if pm == nil {
|
||||
return pushVal(false)
|
||||
}
|
||||
return pushVal(pm.IsPluginDisabled(L.CheckString(1)))
|
||||
}))
|
||||
}
|
||||
|
||||
@ -1207,8 +1237,21 @@ func (p *luaPlugin) Start(s *sdk.PluginSDK) error {
|
||||
}
|
||||
|
||||
func (p *luaPlugin) Stop() error {
|
||||
// ① 先取消事件订阅。**不持 p.mu**:Bus.Publish 持总线锁回调 handler,
|
||||
// 而 handler 要 p.mu;若此处持 p.mu 再取总线锁,就是锁序反转死锁。
|
||||
p.subsMu.Lock()
|
||||
subs := p.subs
|
||||
p.subs = nil
|
||||
p.subsMu.Unlock()
|
||||
for _, unsub := range subs {
|
||||
unsub()
|
||||
}
|
||||
|
||||
// ② 置 closed 并关 L。置位在持锁下完成:已进入但等锁的 event 回调
|
||||
// 拿到锁后会先看到 closed 而直接返回,不会碰已关的 L。
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
p.closed = true
|
||||
|
||||
if p.tbl != nil {
|
||||
fn := p.tbl.RawGetString("stop")
|
||||
|
||||
@ -6,6 +6,7 @@ import (
|
||||
"testing"
|
||||
|
||||
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
|
||||
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
||||
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
|
||||
lua "github.com/yuin/gopher-lua"
|
||||
)
|
||||
@ -597,3 +598,66 @@ return plugin
|
||||
t.Errorf("llm_text writeback: got %q, want %q", sc2.LLMText, "模型输出[尾部标记]")
|
||||
}
|
||||
}
|
||||
|
||||
// TestLuaEventsSubscribeAndStopCleanup 覆盖 sdk.events.subscribe:
|
||||
// 1. 订阅真的能收到内核事件(走内部 SDK 的 Subscribe,不是永远为 nil 的公共 Events());
|
||||
// 2. Stop 会取消订阅,之后 Publish 不得再触碰已 Close 的 LState。
|
||||
func TestLuaEventsSubscribeAndStopCleanup(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
|
||||
os.WriteFile(filepath.Join(dir, "plugin.json"), []byte(`{"name":"evlua","entry":"main.lua"}`), 0644)
|
||||
os.WriteFile(filepath.Join(dir, "main.lua"), []byte(`
|
||||
local plugin = { name = "evlua" }
|
||||
|
||||
function plugin.start(sdk)
|
||||
_G.hits = 0
|
||||
local unsub, err = sdk.events.subscribe("agent_output", function(evt)
|
||||
_G.hits = _G.hits + 1
|
||||
_G.last_type = evt.type
|
||||
_G.last_source = evt.source
|
||||
end)
|
||||
_G.sub_err = err
|
||||
_G.unsub_type = type(unsub)
|
||||
end
|
||||
|
||||
function plugin.stop() end
|
||||
return plugin
|
||||
`), 0644)
|
||||
|
||||
plg, err := tryLoadLua(dir, "evlua", nil)
|
||||
if err != nil {
|
||||
t.Fatalf("tryLoadLua failed: %v", err)
|
||||
}
|
||||
lp := plg.(*luaPlugin)
|
||||
|
||||
bus := events.NewBus()
|
||||
reg := internalConfig.NewConfigRegistry("")
|
||||
sett := sdk.NewSettings("evlua", reg)
|
||||
s := sdk.New("evlua", sdk.SDKConfig{EventBus: bus, Settings: sett})
|
||||
|
||||
if err := plg.Start(s); err != nil {
|
||||
t.Fatalf("Start failed: %v", err)
|
||||
}
|
||||
|
||||
L := lp.L
|
||||
if errStr := L.GetGlobal("sub_err").String(); errStr != "nil" {
|
||||
t.Fatalf("subscribe returned error: %s", errStr)
|
||||
}
|
||||
if got := L.GetGlobal("unsub_type").String(); got != "function" {
|
||||
t.Fatalf("subscribe should return an unsubscribe function, got %s", got)
|
||||
}
|
||||
|
||||
bus.Publish(&events.Event{Type: events.EventAgentOutput, Source: "test-src"})
|
||||
if hits := int(lua.LVAsNumber(L.GetGlobal("hits"))); hits != 1 {
|
||||
t.Fatalf("event handler hits = %d, want 1", hits)
|
||||
}
|
||||
if got := L.GetGlobal("last_source").String(); got != "test-src" {
|
||||
t.Fatalf("event source = %q, want test-src", got)
|
||||
}
|
||||
|
||||
// Stop 取消订阅 + 关 L;此后再 Publish 不得 panic / use-after-close。
|
||||
if err := plg.Stop(); err != nil {
|
||||
t.Fatalf("Stop failed: %v", err)
|
||||
}
|
||||
bus.Publish(&events.Event{Type: events.EventAgentOutput, Source: "after-stop"})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user