diff --git a/internal/plugin/lua_plugin.go b/internal/plugin/lua_plugin.go index 6e60ced..df121c5 100644 --- a/internal/plugin/lua_plugin.go +++ b/internal/plugin/lua_plugin.go @@ -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") diff --git a/internal/plugin/lua_plugin_test.go b/internal/plugin/lua_plugin_test.go index 02fbf82..a2e2beb 100644 --- a/internal/plugin/lua_plugin_test.go +++ b/internal/plugin/lua_plugin_test.go @@ -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"}) +}