diff --git a/cmd/ohos/HomeAgent/entry/src/main/ets/pages/SettingsPage.ets b/cmd/ohos/HomeAgent/entry/src/main/ets/pages/SettingsPage.ets index 2a5470f..ec45d56 100644 --- a/cmd/ohos/HomeAgent/entry/src/main/ets/pages/SettingsPage.ets +++ b/cmd/ohos/HomeAgent/entry/src/main/ets/pages/SettingsPage.ets @@ -62,6 +62,14 @@ export struct SettingsPage { @State editName: string = ''; @State showAddForm: boolean = false; @State addFormVisible: boolean = false; + /** + * 表单当前在编辑哪条连接:空串表示新建。 + * + * 之前只有"添加"入口,ConnStore.updateConnection 写好了却没有任何调用者, + * 于是地址填错的连接只能删掉重建(API Key 也得重敲)。同一套表单 + * 靠这个 id 区分保存走 add 还是 update。 + */ + @State editingId: string = ''; @State themeMode: string = 'system'; @State lang: string = 'zh'; @State bgImage: string = ''; @@ -205,11 +213,12 @@ export struct SettingsPage { this.showToast('名称和地址不能为空', true); return; } + if (this.editingId.length > 0) { + this.updateConnection(this.editingId, name, url, apiKey); + return; + } connStore.addConnection(name, url, apiKey).then(() => { - this.editName = ''; - this.editUrl = ''; - this.editApiKey = ''; - this.showAddForm = false; + this.closeConnForm(); const cur = connStore.getCurrentConnection(); if (cur !== null) { apiClient.setConnection(cur); @@ -219,6 +228,53 @@ export struct SettingsPage { }); } + /** + * 保存对已有连接的修改。 + * + * 修改当前生效的连接后必须重新 setConnection:ApiClient 持有的是 + * ConnectionConfig 的引用快照,不刷新的话后续请求还会打到旧地址。 + */ + private updateConnection(id: string, name: string, url: string, apiKey: string): void { + connStore.updateConnection(id, name, url, apiKey).then(() => { + this.closeConnForm(); + const cur = connStore.getCurrentConnection(); + if (cur !== null) { + apiClient.setConnection(cur); + } + this.loadConnections(); + this.showToast('连接已更新', false); + }); + } + + /** 打开表单:id 为空是新建,非空是编辑并回填原值(API Key 一并带出,避免用户重敲)。 */ + private openConnForm(conn: ConnectionConfig | null): void { + this.showAddForm = true; + this.addFormVisible = false; + if (conn === null) { + this.editingId = ''; + this.editName = ''; + this.editUrl = ''; + this.editApiKey = ''; + } else { + this.editingId = conn.id; + this.editName = conn.name; + this.editUrl = conn.url; + this.editApiKey = conn.apiKey; + } + setTimeout(() => { + this.addFormVisible = true; + }, 30); + } + + private closeConnForm(): void { + this.showAddForm = false; + this.addFormVisible = false; + this.editingId = ''; + this.editName = ''; + this.editUrl = ''; + this.editApiKey = ''; + } + private deleteConnection(id: string): void { connStore.deleteConnection(id).then(() => { this.loadConnections(); @@ -840,14 +896,7 @@ export struct SettingsPage { .backgroundColor(this.palette().accent) .fontColor(Color.White) .onClick(() => { - this.showAddForm = true; - this.addFormVisible = false; - this.editName = ''; - this.editUrl = ''; - this.editApiKey = ''; - setTimeout(() => { - this.addFormVisible = true; - }, 30); + this.openConnForm(null); }) } .width('100%') @@ -855,6 +904,10 @@ export struct SettingsPage { if (this.showAddForm) { Column() { + Text(this.editingId.length > 0 ? '编辑连接' : '新建连接') + .fontSize(12) + .fontColor(this.palette().textSecondary) + .margin({ bottom: 10 }) TextInput({ placeholder: '名称 (如 HomeAgent)', text: this.editName }) .height(36).fontSize(13).fontColor(this.palette().textPrimary) .placeholderColor(this.palette().textMuted).backgroundColor(this.palette().bgInput) @@ -888,7 +941,7 @@ export struct SettingsPage { .border({ width: 1, color: this.palette().btnGhostBorder }) .fontColor(this.palette().textSecondary) .onClick(() => { - this.showAddForm = false; + this.closeConnForm(); }) Blank() Button('保存') @@ -961,6 +1014,16 @@ export struct SettingsPage { .backgroundColor(this.palette().accentBg) .margin({ right: 6 }) } + Button('编辑') + .height(26) + .fontSize(11) + .backgroundColor(Color.Transparent) + .border({ width: 1, color: this.palette().btnGhostBorder }) + .fontColor(this.palette().textSecondary) + .margin({ right: 6 }) + .onClick(() => { + this.openConnForm(conn); + }) Button('删除') .height(26) .fontSize(11) @@ -980,7 +1043,10 @@ export struct SettingsPage { color: conn.id === this.currentId ? this.palette().accent : Color.Transparent, }) .margin({ bottom: 6 }) - }, (conn: ConnectionConfig) => conn.id) + // 键里带上 name/url:ForEach 对相同键只更新绑定、不重跑 @Builder 体, + // 只用 id 做键时改完地址这一行还显示旧值。行内没有 TextInput, + // 因此把可变字段放进键不会有"编辑时焦点被销毁"的副作用。 + }, (conn: ConnectionConfig) => conn.id + '|' + conn.name + '|' + conn.url) } } diff --git a/internal/agent/core/stages.go b/internal/agent/core/stages.go index add7da1..36cb9ee 100644 --- a/internal/agent/core/stages.go +++ b/internal/agent/core/stages.go @@ -14,14 +14,25 @@ type StageHost struct { toolDefs []sdk.ToolDef tools map[string]sdk.ToolHandler toolPlugins map[string]string - stages map[sdk.Stage][]sdk.StageHandler + stages map[sdk.Stage][]stageEntry +} + +// stageEntry 把 stage handler 与它的归属插件绑定。 +// +// 为何需要归属:子进程插件崩溃后,它注册的 handler 闭包仍在这张表里, +// 每次 RunStage 都会经 RPC 打向已死进程并报 ErrProcessExited;重启后新 handler +// 又追加进来,旧的永不退场——错误与重复执行随重启次数线性累积。 +// 有了归属才能在卸载/崩溃时成组摘除。 +type stageEntry struct { + plugin string + fn sdk.StageHandler } func NewStageHost() *StageHost { return &StageHost{ tools: make(map[string]sdk.ToolHandler), toolPlugins: make(map[string]string), - stages: make(map[sdk.Stage][]sdk.StageHandler), + stages: make(map[sdk.Stage][]stageEntry), } } @@ -41,9 +52,15 @@ func (h *StageHost) RegisterTool(name string, def sdk.ToolDef, handler sdk.ToolH } func (h *StageHost) RegisterStage(stage sdk.Stage, handler sdk.StageHandler) { + h.RegisterStageFor("", stage, handler) +} + +// RegisterStageFor 注册带归属插件名的 stage handler。 +// plugin 为空时等同 RegisterStage(内核自身注册的 handler,不参与成组摘除)。 +func (h *StageHost) RegisterStageFor(plugin string, stage sdk.Stage, handler sdk.StageHandler) { h.mu.Lock() defer h.mu.Unlock() - h.stages[stage] = append(h.stages[stage], handler) + h.stages[stage] = append(h.stages[stage], stageEntry{plugin: plugin, fn: handler}) } func (h *StageHost) GetToolDefs() []sdk.ToolDef { @@ -106,6 +123,36 @@ func (h *StageHost) UnregisterPluginTools(pluginName string) { h.toolDefs = keepDefs } +// UnregisterPluginStages 摘除某插件注册的全部 stage handler,返回摘除数量。 +// +// 与 UnregisterPluginTools 成对:卸载/重载/崩溃时两者都得做, +// 否则插件的工具没了但 stage handler 还在,继续打向不存在的插件。 +func (h *StageHost) UnregisterPluginStages(pluginName string) int { + if pluginName == "" { + return 0 + } + h.mu.Lock() + defer h.mu.Unlock() + + removed := 0 + for stage, entries := range h.stages { + keep := entries[:0:0] + for _, e := range entries { + if e.plugin == pluginName { + removed++ + continue + } + keep = append(keep, e) + } + if len(keep) == 0 { + delete(h.stages, stage) + continue + } + h.stages[stage] = keep + } + return removed +} + func inferToolPlugin(name string) string { for i := 0; i < len(name); i++ { if name[i] == '_' { @@ -123,14 +170,15 @@ func inferToolPlugin(name string) string { // handler 返回的 error 会被收集到 ctx.Errors 中并记录日志,不会中断其他 handler 的执行。 func (h *StageHost) RunStage(stage sdk.Stage, ctx *sdk.StageContext) { h.mu.RLock() - handlers := h.stages[stage] + entries := make([]stageEntry, len(h.stages[stage])) + copy(entries, h.stages[stage]) h.mu.RUnlock() - if len(handlers) == 0 { + if len(entries) == 0 { return } var wg sync.WaitGroup - errCh := make(chan error, len(handlers)) - for _, handler := range handlers { + errCh := make(chan error, len(entries)) + for _, entry := range entries { wg.Add(1) go func(fn sdk.StageHandler) { defer wg.Done() @@ -142,7 +190,7 @@ func (h *StageHost) RunStage(stage sdk.Stage, ctx *sdk.StageContext) { if err := fn(ctx); err != nil { errCh <- err } - }(handler) + }(entry.fn) } wg.Wait() close(errCh) diff --git a/internal/agent/core/stages_plugin_test.go b/internal/agent/core/stages_plugin_test.go new file mode 100644 index 0000000..f67ddab --- /dev/null +++ b/internal/agent/core/stages_plugin_test.go @@ -0,0 +1,96 @@ +package core + +import ( + "testing" + + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// stage handler 必须能按插件成组摘除。 +// +// 修复前 StageHost.stages 只存匿名函数,没有归属信息: +// 插件崩溃/卸载后它的 handler 永远留在表里,每轮 RunStage 都被并发调起并 +// 打向已死进程;重启后新 handler 追加进来,旧的仍不退场—— +// 错误与重复执行随重启次数线性累积。 +func TestUnregisterPluginStages_RemovesOnlyThatPlugin(t *testing.T) { + h := NewStageHost() + + var aRan, bRan, coreRan int + h.RegisterStageFor("a", sdk.StagePreAction, func(*sdk.StageContext) error { aRan++; return nil }) + h.RegisterStageFor("b", sdk.StagePreAction, func(*sdk.StageContext) error { bRan++; return nil }) + // 内核自身注册的 handler(无归属)不该被插件摘除波及 + h.RegisterStage(sdk.StagePreAction, func(*sdk.StageContext) error { coreRan++; return nil }) + + h.RunStage(sdk.StagePreAction, &sdk.StageContext{}) + if aRan != 1 || bRan != 1 || coreRan != 1 { + t.Fatalf("首轮应全部执行,a=%d b=%d core=%d", aRan, bRan, coreRan) + } + + if n := h.UnregisterPluginStages("a"); n != 1 { + t.Errorf("应摘除 1 个 handler,实际 %d", n) + } + + h.RunStage(sdk.StagePreAction, &sdk.StageContext{}) + if aRan != 1 { + t.Errorf("已摘除的插件 handler 不该再被调用,实际执行 %d 次", aRan) + } + if bRan != 2 || coreRan != 2 { + t.Errorf("其他 handler 应照常执行,b=%d core=%d", bRan, coreRan) + } +} + +// 摘除某插件的最后一个 handler 后,该 stage 应从表中消失(RunStage 直接短路)。 +func TestUnregisterPluginStages_DropsEmptyStage(t *testing.T) { + h := NewStageHost() + h.RegisterStageFor("solo", sdk.StageAfterToolcall, func(*sdk.StageContext) error { return nil }) + + if n := h.UnregisterPluginStages("solo"); n != 1 { + t.Fatalf("应摘除 1 个,实际 %d", n) + } + h.mu.RLock() + _, exists := h.stages[sdk.StageAfterToolcall] + h.mu.RUnlock() + if exists { + t.Error("stage 已无 handler 时应从表中删除") + } +} + +// 空插件名不得误摘内核自身注册的 handler。 +func TestUnregisterPluginStages_EmptyNameIsNoop(t *testing.T) { + h := NewStageHost() + ran := 0 + h.RegisterStage(sdk.StagePreAction, func(*sdk.StageContext) error { ran++; return nil }) + + if n := h.UnregisterPluginStages(""); n != 0 { + t.Errorf("空插件名应是 no-op,实际摘除 %d", n) + } + h.RunStage(sdk.StagePreAction, &sdk.StageContext{}) + if ran != 1 { + t.Errorf("内核 handler 应保留并执行,实际 %d 次", ran) + } +} + +// 工具与 stage 的摘除互不干扰:都摘完后两者皆空。 +func TestUnregisterPluginToolsAndStages_Together(t *testing.T) { + h := NewStageHost() + if err := h.RegisterTool("demo_run", sdk.ToolDef{Name: "demo_run", Plugin: "demo"}, + func(map[string]interface{}) (interface{}, error) { return nil, nil }); err != nil { + t.Fatal(err) + } + h.RegisterStageFor("demo", sdk.StagePreAction, func(*sdk.StageContext) error { return nil }) + + h.UnregisterPluginTools("demo") + h.UnregisterPluginStages("demo") + + if h.ToolCount() != 0 { + t.Errorf("工具应已摘除,实际 %d", h.ToolCount()) + } + if h.ToolPlugin("demo_run") != "" { + t.Error("工具→插件映射应清空") + } + // 摘除后可重新注册同名工具(重启路径的前提) + if err := h.RegisterTool("demo_run", sdk.ToolDef{Name: "demo_run", Plugin: "demo"}, + func(map[string]interface{}) (interface{}, error) { return nil, nil }); err != nil { + t.Errorf("摘除后应可重新注册同名工具,实际: %v", err) + } +} diff --git a/internal/plugin/crash_recovery_test.go b/internal/plugin/crash_recovery_test.go new file mode 100644 index 0000000..a3930ed --- /dev/null +++ b/internal/plugin/crash_recovery_test.go @@ -0,0 +1,249 @@ +package plugin + +import ( + "errors" + "os" + "path/filepath" + "sync" + "testing" + "time" + + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// fakeCleaner 记录内核侧摘除动作,用于断言「插件死后注册面被摘干净」。 +// +// 同时实现 PluginToolCleaner 与 PluginStageCleaner——生产里 StageHost 两者都实现。 +type fakeCleaner struct { + mu sync.Mutex + tools []string + stages []string +} + +func (f *fakeCleaner) UnregisterPluginTools(name string) { + f.mu.Lock() + defer f.mu.Unlock() + f.tools = append(f.tools, name) +} + +func (f *fakeCleaner) UnregisterPluginStages(name string) int { + f.mu.Lock() + defer f.mu.Unlock() + f.stages = append(f.stages, name) + return 1 +} + +func (f *fakeCleaner) toolCalls() []string { + f.mu.Lock() + defer f.mu.Unlock() + return append([]string(nil), f.tools...) +} + +func (f *fakeCleaner) stageCalls() []string { + f.mu.Lock() + defer f.mu.Unlock() + return append([]string(nil), f.stages...) +} + +func contains(list []string, want string) bool { + for _, v := range list { + if v == want { + return true + } + } + return false +} + +// newTestRegistry 造一个可用于 detach/崩溃路径测试的最小 Registry。 +func newTestRegistry(t *testing.T) (*Registry, *fakeCleaner, *agentIO.IOManager) { + t.Helper() + cleaner := &fakeCleaner{} + iom := agentIO.NewIOManager() + r := &Registry{ + plugins: make(map[string]sdk.Plugin), + factories: make(map[string]NativeFactory), + pluginAutoRestart: make(map[string]bool), + sdkRefs: make(map[string]*sdk.PluginSDK), + knownDisabled: make(map[string]bool), + pluginHashes: make(map[string]string), + pluginChannels: make(map[string]*pluginChannelSet), + plgDir: t.TempDir(), + toolCleaner: cleaner, + iom: iom, + } + return r, cleaner, iom +} + +// detachPlugin 必须同时摘工具、stage handler、IO 通道。 +// +// 此前各卸载路径只调 UnregisterPluginTools,漏了后两项: +// 插件的工具没了但 stage handler 还在每轮 RunStage 里被调用并失败, +// output device 还留在 IOManager 里让模型看到一个永远发不出去的通道。 +func TestDetachPlugin_RemovesToolsStagesAndChannels(t *testing.T) { + r, cleaner, iom := newTestRegistry(t) + + // 模拟插件注册过通道 + if err := iom.RegisterDevice(&channelDevice{name: "demo_out"}); err != nil { + t.Fatalf("RegisterDevice: %v", err) + } + iom.RegisterInputChannel("demo_in", agentIO.ChannelDef{}) + r.noteChannel("demo", "demo_out", true) + r.noteChannel("demo", "demo_in", false) + + r.detachPlugin("demo") + + if !contains(cleaner.toolCalls(), "demo") { + t.Error("应摘除插件工具") + } + if !contains(cleaner.stageCalls(), "demo") { + t.Error("应摘除插件 stage handler(否则每轮 RunStage 都会打向已死插件)") + } + if iom.GetDevice("demo_out") != nil { + t.Error("output device 应被摘除,否则模型仍看到一个必然失败的通道") + } + if _, ok := iom.GetInputChannelDef("demo_in"); ok { + t.Error("input channel 定义应被摘除") + } +} + +// 通道台账在 detach 后清空,使插件重启时能重新注册同名通道。 +// +// 不清空的后果:RegisterDevice 撞上同名旧 device 直接报 already registered, +// 新进程的通道注册不上——插件“重启成功”了但通道永久指向已死进程。 +func TestReleasePluginChannels_AllowsReRegistrationAfterRestart(t *testing.T) { + r, _, iom := newTestRegistry(t) + + if err := iom.RegisterDevice(&channelDevice{name: "qq"}); err != nil { + t.Fatalf("首次注册: %v", err) + } + r.noteChannel("qq", "qq", true) + + r.detachPlugin("qq") + + // 重启后同名通道必须能重新注册 + if err := iom.RegisterDevice(&channelDevice{name: "qq"}); err != nil { + t.Fatalf("摘除后应可重新注册同名通道,实际: %v", err) + } + // 台账已清空,重复 detach 不应再摘掉新注册的那个 + r.releasePluginChannels("qq") + if iom.GetDevice("qq") == nil { + t.Error("台账已清空,重复 detach 不应摘掉重启后新注册的通道") + } +} + +// 崩溃回调必须摘注册面 + 排重启,且不阻塞调用方(它跑在 readLoop 的 goroutine 里)。 +func TestOnProcCrash_DetachesImmediately(t *testing.T) { + r, cleaner, _ := newTestRegistry(t) + // 没有插件目录 → ReloadOne 必然失败,但摘除动作应已完成 + r.pluginAutoRestart["ghost"] = false // 关掉自动重启,只验摘除 + + done := make(chan struct{}) + go func() { + r.onProcCrash("ghost", errors.New("signal: killed")) + close(done) + }() + + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("onProcCrash 不应阻塞(它在 readLoop 的 goroutine 上)") + } + + if !contains(cleaner.toolCalls(), "ghost") { + t.Error("崩溃后应立即摘除工具,否则模型继续调用一个必然失败的工具") + } + if !contains(cleaner.stageCalls(), "ghost") { + t.Error("崩溃后应摘除 stage handler") + } +} + +// 声明了不自动重启的插件,崩溃后不得被拉起。 +func TestScheduleProcRestart_RespectsAutoRestartOff(t *testing.T) { + r, _, _ := newTestRegistry(t) + r.pluginAutoRestart["noauto"] = false + + r.scheduleProcRestart("noauto", errors.New("boom")) + + if n := r.crashCount("noauto"); n != 0 { + t.Errorf("禁用自动重启时不该记崩溃计数,实际 %d", n) + } +} + +// 窗口内连续崩溃超过上限后停止自动重启,避免崩溃循环打满 CPU。 +func TestScheduleProcRestart_StopsAfterThreshold(t *testing.T) { + r, _, _ := newTestRegistry(t) + + for i := 0; i < procMaxRestarts+2; i++ { + r.noteCrash("loopy") + } + if got := r.crashCount("loopy"); got != procMaxRestarts+2 { + t.Fatalf("崩溃计数应累计,实际 %d", got) + } + + // 超阈值后再调不应尝试重启(无插件目录时重启必然失败并留日志, + // 这里只验它提前返回:计数不再增长)。 + before := r.crashCount("loopy") + r.scheduleProcRestart("loopy", errors.New("again")) + if after := r.crashCount("loopy"); after != before+1 { + t.Errorf("应只记一次计数即返回,before=%d after=%d", before, after) + } +} + +// 崩溃计数在窗口外自动归零,避免偶发崩溃永久累积成“不可重启”。 +func TestNoteCrash_WindowExpiry(t *testing.T) { + r, _, _ := newTestRegistry(t) + + r.noteCrash("old") + r.crashMu.Lock() + r.procCrashes["old"].last = time.Now().Add(-procCrashWindow - time.Second) + r.crashMu.Unlock() + + if got := r.crashCount("old"); got != 0 { + t.Errorf("窗口外计数应归零,实际 %d", got) + } + if got := r.noteCrash("old"); got != 1 { + t.Errorf("窗口外应重新从 1 计,实际 %d", got) + } +} + +// 关停途中不得再拉起插件:段已拆而进程还在会直接 SIGBUS。 +func TestScheduleProcRestart_SkippedDuringShutdown(t *testing.T) { + r, _, _ := newTestRegistry(t) + r.shuttingDown.Store(true) + + r.scheduleProcRestart("any", errors.New("boom")) + + if n := r.crashCount("any"); n != 0 { + t.Errorf("关停中应直接返回,不记计数,实际 %d", n) + } +} + +// PluginRuntime 对未安装插件返回 false,对有目录的插件报告加载通道。 +func TestPluginRuntime_ReportsChannelAndInstallState(t *testing.T) { + r, _, _ := newTestRegistry(t) + + if _, ok := r.PluginRuntime("nope"); ok { + t.Error("未安装插件应返回 false") + } + + dir := filepath.Join(r.plgDir, "procplug") + if err := os.MkdirAll(dir, 0o755); err != nil { + t.Fatal(err) + } + bin := filepath.Join(dir, binEntry) + if err := os.WriteFile(bin, []byte("#!/bin/true\n"), 0o755); err != nil { + t.Fatal(err) + } + + info, ok := r.PluginRuntime("procplug") + if !ok { + t.Fatal("有插件目录应视为已安装") + } + if info.Channel != "proc" { + t.Errorf("应识别为 proc 通道,实际 %q", info.Channel) + } + if info.Loaded || info.Alive { + t.Error("未加载的插件不应报告 loaded/alive") + } +} diff --git a/internal/plugin/dynamic_proc_unix.go b/internal/plugin/dynamic_proc_unix.go index 8e62471..20ed3f1 100644 --- a/internal/plugin/dynamic_proc_unix.go +++ b/internal/plugin/dynamic_proc_unix.go @@ -129,23 +129,126 @@ func (r *Registry) closeProcHost() { // **崩溃隔离**:子进程死亡只影响自己,homed 继续服务——对比 C ABI 下 // 插件 panic 直接带崩整个进程(§1.2,现网已发生)。 // -// 崩溃计数/冷却/自愈复用既有 plugin_health(§2.3),本函数只负责把 -// 进程退出这一事实转成事件通知;具体重载策略由 agent 侧决定。 +// 但「homed 没崩」不等于「内核状态干净」。此前本函数只发了一个事件, +// 而全仓没有任何订阅者,于是生产上出现过 editdoc 被 kill 后: +// - `edit_document` 仍留在 StageHost 的工具表里,模型照旧看得到、照旧调用, +// 每次都吃到 `proc: 插件进程已退出`; +// - 该插件的 stage handler 仍在每轮 RunStage 里被并发调起并失败; +// - 没有任何路径把它拉回来,插件永久缺席直到重启 homed。 +// +// 所以崩溃回调必须做三件事:摘注册面、喂健康计数、排一次重启。 func (r *Registry) onProcCrash(name string, err error) { log.Printf("[plugin] 子进程插件 %s 异常退出: %v(homed 未受影响)", name, err) - if r.evBus == nil { + + // 1) 摘掉工具/stage/通道。**必须先做**:从这一刻起模型就不该再看到这些工具, + // 否则在重启完成前的窗口里每次调用都是确定的失败。 + r.detachPlugin(name) + + // 2) 从注册表移除。不做的后果:scheduleProcRestart 里的“已被其他路径重新加载” + // 复核会误判(旧条目还在,See plugins[name] != nil),跳过真正的自动重启。 + // Plugin 对象本身仍被 proc 持有,Kill/回收不受影响。 + r.mu.Lock() + delete(r.plugins, name) + delete(r.sdkRefs, name) + for i, inst := range r.instances { + if inst.Name() == name { + r.instances = append(r.instances[:i], r.instances[i+1:]...) + break + } + } + r.mu.Unlock() + + // 2) 事件通知(webui/诊断插件可订阅)。 + if r.evBus != nil { + r.evBus.Publish(&events.Event{ + Type: events.EventSystem, + Source: "plugin", + Payload: map[string]interface{}{ + "event": "plugin_crashed", + "plugin": name, + "error": err.Error(), + }, + Timestamp: time.Now().Unix(), + }) + } + + // 3) 排一次重启。**必须异步**:本回调由 proc.markExited 在 readLoop 的 + // goroutine 里触发,而 ReloadOne 要拿 registry 锁、还要 Kill 并 join 同一个 + // readLoop(Process.Kill 里 readerWG.Wait),同步调用会自锁死。 + go r.scheduleProcRestart(name, err) +} + +// scheduleProcRestart 在崩溃后按退避重启子进程插件。 +// +// 退避与阈值语义与 agent 侧 plugin_health 对齐(窗口内 3 次即判定不健康), +// 但重启动作落在 registry:崩溃事实产生于此,agent 的 distillLoop 默认 30 分钟 +// 才转一次(生产实配 2d),靠它兜底等于插件缺席数小时。 +func (r *Registry) scheduleProcRestart(name string, cause error) { + if r.shuttingDown.Load() { + return // 内核正在关停,不再拉起 + } + if !r.AutoRestartEnabled(name) { + log.Printf("[plugin] %s 声明了不自动重启,保持缺席状态", name) return } - // 不在此处直接重载:重载需要 registry 锁,而本回调可能在 - // 持锁路径的 goroutine 中触发,直接调用会死锁。 - r.evBus.Publish(&events.Event{ - Type: events.EventSystem, - Source: "plugin", - Payload: map[string]interface{}{ - "event": "plugin_crashed", - "plugin": name, - "error": err.Error(), - }, - Timestamp: time.Now().Unix(), - }) + + n := r.noteCrash(name) + if n > procMaxRestarts { + log.Printf("[plugin] %s 在 %v 内崩溃 %d 次,停止自动重启(需人工介入)", + name, procCrashWindow, n) + return + } + + // 线性退避:1 次→1s,2 次→2s,3 次→3s。崩溃循环时不至于打满 CPU, + // 又足够快到用户感知不到工具缺席。 + delay := time.Duration(n) * procRestartBackoff + time.Sleep(delay) + + // 期间可能已被 Disable/Remove/手工 plgreload 处理掉,重启前复核。 + if r.shuttingDown.Load() { + return + } + if r.isDisabled(name) { + log.Printf("[plugin] %s 已被禁用,取消自动重启", name) + return + } + r.mu.RLock() + already := r.plugins[name] != nil + r.mu.RUnlock() + if already { + log.Printf("[plugin] %s 已被其他路径重新加载,取消自动重启", name) + return + } + + log.Printf("[plugin] 自动重启 %s(第 %d 次,退避 %v,起因: %v)", name, n, delay, cause) + if err := r.ReloadOne(name); err != nil { + log.Printf("[plugin] %s 自动重启失败: %v", name, err) + return + } + log.Printf("[plugin] %s 自动重启成功", name) +} + +// noteCrash 记录一次崩溃并返回窗口内的累计次数。 +func (r *Registry) noteCrash(name string) int { + now := time.Now() + r.crashMu.Lock() + defer r.crashMu.Unlock() + if r.procCrashes == nil { + r.procCrashes = make(map[string]*procCrashRecord) + } + rec := r.procCrashes[name] + if rec == nil || now.Sub(rec.last) > procCrashWindow { + rec = &procCrashRecord{} + r.procCrashes[name] = rec + } + rec.count++ + rec.last = now + return rec.count +} + +// ResetProcCrashCount 清空某插件的崩溃计数(人工 plgreload / 重新启用后调用)。 +func (r *Registry) ResetProcCrashCount(name string) { + r.crashMu.Lock() + defer r.crashMu.Unlock() + delete(r.procCrashes, name) } diff --git a/internal/plugin/proc/host.go b/internal/plugin/proc/host.go index 5c50955..218a922 100644 --- a/internal/plugin/proc/host.go +++ b/internal/plugin/proc/host.go @@ -45,6 +45,12 @@ type Host struct { stageMu sync.Mutex coordMu sync.Mutex coord *stageCoordinator + + // sup 是内核侧唯一的子进程台账,与共享段同生命周期。 + // + // 放在 Host 而不是 registry 的理由:能拿到 Host 的地方就能拿到台账, + // 而 Host 本就是「全部子进程插件共享的那一份内核侧状态」。 + sup *Supervisor } // NewHost 创建共享段(平台层 allocShm + 布局初始化)。 @@ -81,6 +87,7 @@ func NewHost() (*Host, error) { evtRing.Init() return &Host{ + sup: NewSupervisor(), memfd: memfd, data: data, seg: seg, @@ -102,7 +109,15 @@ func NewHost() (*Host, error) { const shmDefaultSize = 256 * 1024 // Close 释放共享段(StageContext + 事件环)。 +// Supervisor 返回子进程台账(供 registry 查询/关停)。 +func (h *Host) Supervisor() *Supervisor { return h.sup } + func (h *Host) Close() error { + // 先停全部子进程再拆段:插件还持有映射时 unmap, + // 它们下一次访问共享段就是 SIGBUS。 + if h.sup != nil { + h.sup.StopAll(0) + } var firstErr error if h.data != nil { if err := freeShm(h.memfd, h.data); err != nil && firstErr == nil { diff --git a/internal/plugin/proc/plugin.go b/internal/plugin/proc/plugin.go index 3c99cd3..70a4be9 100644 --- a/internal/plugin/proc/plugin.go +++ b/internal/plugin/proc/plugin.go @@ -6,6 +6,7 @@ import ( "fmt" "log" "sync" + "sync/atomic" pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" ) @@ -43,6 +44,14 @@ type Plugin struct { caps *capabilitySet stopOnce sync.Once + + // stopping 标记「本次退出是内核主动发起的」,用于压掉 onCrash。 + // + // 必要性:Stop() 宽限期超时与 Close() 都走 Process.Kill(), + // 而 Kill 产生的 `signal: killed` 是非 nil 的 waitErr——若不区分, + // 重载/禁用/卸载这些**内核自己发起**的停止会被 handleExit 当成崩溃上报, + // 触发一轮多余的自动重启(重载路径下等于把刚装好的插件又推倒一次)。 + stopping atomic.Bool } // New 创建子进程插件(不启动进程)。 @@ -67,6 +76,31 @@ func New(name, bin, dir string, config map[string]interface{}, host *Host, onCra // Name 实现 sdk.Plugin。 func (p *Plugin) Name() string { return p.name } +// PID 返回子进程号;未启动或已退出返回 0。 +// 供 pluginmgr 呈现「插件实际在跑哪个进程」。 +func (p *Plugin) PID() int { + if p.proc == nil { + return 0 + } + if !p.Alive() { + return 0 + } + return p.proc.PID() +} + +// Alive 报告子进程是否仍存活。 +func (p *Plugin) Alive() bool { + if p.proc == nil { + return false + } + select { + case <-p.proc.Exited(): + return false + default: + return true + } +} + // Start 启动子进程并完成注册。 // // core 是内核为该插件构建的能力面(internal/sdk.PluginSDK 天然满足 CoreSDK)。 @@ -99,6 +133,7 @@ func (p *Plugin) Start(core CoreSDK) error { EvtRingSize: evtTotalSize, Handler: p.handler.Handle, OnExit: p.handleExit, + Supervisor: p.host.Supervisor(), }) if err != nil { return err @@ -124,6 +159,7 @@ func (p *Plugin) Start(core CoreSDK) error { // Stop 优雅停止(实现 sdk.Plugin)。 func (p *Plugin) Stop() error { + p.stopping.Store(true) var err error p.stopOnce.Do(func() { if p.proc != nil { @@ -138,6 +174,7 @@ func (p *Plugin) Stop() error { // **这里是真 kill + wait**——对比 cabi 路径的 Close 只做 dlclose, // 而 dlclose 对 Go c-shared 是 no-op(§1.1,热重载静默失效的根因)。 func (p *Plugin) Close() error { + p.stopping.Store(true) var err error p.stopOnce.Do(func() { if p.proc != nil { @@ -156,6 +193,11 @@ func (p *Plugin) handleExit(name string, err error) { if p.host != nil && p.host.ForceReleaseLock(name) { log.Printf("[proc] %s 退出,内核已释放其持有的 stage 锁", name) } + // 内核主动停止(Stop/Close,含宽限期超时后的 Kill)不算崩溃: + // 否则重载/禁用/卸载都会误触发自动重启。 + if p.stopping.Load() { + return + } if err != nil && p.onCrash != nil { p.onCrash(name, err) } diff --git a/internal/plugin/proc/procattr_linux.go b/internal/plugin/proc/procattr_linux.go new file mode 100644 index 0000000..fff185e --- /dev/null +++ b/internal/plugin/proc/procattr_linux.go @@ -0,0 +1,30 @@ +//go:build linux + +package proc + +import ( + "os/exec" + "syscall" +) + +// applyProcAttr 让子进程在父进程(homed)死亡时收到 SIGKILL。 +// +// 这是**最后一道兜底**,不是主路径:正常关停走 Supervisor.StopAll。 +// 它兜的是内核自身异常终止的场景——homed 被 SIGKILL、段错误、OOM—— +// 此时没有任何 Go 代码有机会运行,Supervisor 也来不及 StopAll, +// 子进程会被 init 收养成孤儿: +// - 继续持有已被 unmap 的共享段映射,下次访问即 SIGBUS; +// - 与新启动的 homed 抢同一份外部资源(qq 的 WS 会话、browser 的 +// chromium profile 锁),表现为"重启后插件时好时坏"。 +// +// Pdeathsig 由内核在父进程退出时投递,不依赖任何用户态代码, +// 因此在 homed 被 SIGKILL 的情况下依然生效。 +// +// 仅 Linux 有此机制。macOS/Windows 无等价物,回退为空实现(procattr_other.go): +// 那两个平台上孤儿风险依旧存在,靠 StopAll 覆盖正常关停路径。 +func applyProcAttr(cmd *exec.Cmd) { + if cmd.SysProcAttr == nil { + cmd.SysProcAttr = &syscall.SysProcAttr{} + } + cmd.SysProcAttr.Pdeathsig = syscall.SIGKILL +} diff --git a/internal/plugin/proc/procattr_other.go b/internal/plugin/proc/procattr_other.go new file mode 100644 index 0000000..c128d2b --- /dev/null +++ b/internal/plugin/proc/procattr_other.go @@ -0,0 +1,15 @@ +//go:build !linux + +package proc + +import "os/exec" + +// applyProcAttr 在非 Linux 平台是空实现。 +// +// macOS 没有 Pdeathsig(kqueue 的 NOTE_EXIT 要求父进程存活才能监听, +// 恰好在父进程被 SIGKILL 时失效);Windows 的 Job Object 可做到类似效果, +// 但需要额外的句柄管理,且 Windows 侧尚未真机验证(§12.5),不在此引入。 +// +// 后果:这两个平台上 homed 被强杀时子进程会成为孤儿。 +// 正常关停路径(Supervisor.StopAll)不受影响。 +func applyProcAttr(cmd *exec.Cmd) {} diff --git a/internal/plugin/proc/process.go b/internal/plugin/proc/process.go index 0c04f50..0c5b3d0 100644 --- a/internal/plugin/proc/process.go +++ b/internal/plugin/proc/process.go @@ -36,6 +36,11 @@ type Process struct { stdin *bufio.Writer stdout io.ReadCloser + // stdinFile / stdoutFile 是父进程侧的管道端(手工 os.Pipe,非 cmd.StdinPipe)。 + // 持有它们才能在退出时主动 Close,逼 readLoop 从 Scan 里出来。 + stdinFile *os.File + stdoutFile *os.File + // writeMu 串行化 stdin 写入:NDJSON 帧不能交错,否则对端解析错乱。 writeMu sync.Mutex @@ -48,18 +53,27 @@ type Process struct { // handler 处理插件反向发起的调用(51 个 core.* method)。 handler RequestHandler - // exited 在 readLoop 检测到 EOF/进程退出后关闭,用于唤醒所有等待者。 + // exited 在进程被收割后关闭,用于唤醒所有等待者。 exited chan struct{} exitOnce sync.Once exitErr atomic.Pointer[error] readerWG sync.WaitGroup + waiterWG sync.WaitGroup readyOnce sync.Once ready chan struct{} + // waitErr 由**唯一的** waitLoop 写入:cmd.Wait() 的返回值。 + // waitDone 关闭后 waitErr 才可读。 + waitErr error + waitDone chan struct{} + // onExit 在进程退出时回调(内核用它喂 plugin_health.recordCrash, // 以及 ForceRelease 释放该插件持有的 stage 锁)。 onExit func(name string, err error) + // sup 是内核的集中进程表(可为 nil,单测直接 Spawn 时)。 + sup *Supervisor + // shmSize 是握手时告知插件的共享段大小(0 表示本插件不用共享段)。 shmSize int // evtRingSize 是事件环段大小(0 表示不支持事件环)。 @@ -87,6 +101,8 @@ type Options struct { Handler RequestHandler // OnExit 进程退出回调。 OnExit func(name string, err error) + // Supervisor 是内核的集中进程表;为 nil 时不纳管(单测路径)。 + Supervisor *Supervisor // HandshakeTimeout 建链超时,默认 10s。 HandshakeTimeout time.Duration } @@ -97,6 +113,9 @@ const ( // stopGracePeriod 是发出 plugin.stop 后等待进程自行退出的时间。 // 超时则 Kill——**这是"真正的取消"**,对比 cgo 路径超时后线程永久泄漏。 stopGracePeriod = 5 * time.Second + // killReapTimeout 是 SIGKILL 后等待 waitLoop 收割的上限。 + // 正常情况 wait4 微秒级返回;超过说明卡在不可中断的内核态。 + killReapTimeout = 2 * time.Second ) // ErrProcessExited 表示子进程已退出,调用无法完成。 @@ -120,35 +139,68 @@ func Spawn(name, bin string, opts Options) (*Process, error) { cmd.Env = append(os.Environ(), opts.Env...) } cmd.ExtraFiles = opts.ExtraFiles + applyProcAttr(cmd) - stdinPipe, err := cmd.StdinPipe() + // 管道手工创建而非用 cmd.StdinPipe/StdoutPipe。 + // + // 原因:cmd.Wait() 会等待并**关闭** StdinPipe/StdoutPipe 创建的管道, + // 且文档明确要求“读完再 Wait”。既然现在有一根专职的 waitLoop 立即 + // Wait(不等 readLoop),就必须自己控制管道生命期,否则会与 + // os/exec 的内部关闭竞争,在 readLoop 里读到 "file already closed"。 + stdinR, stdinW, err := os.Pipe() if err != nil { return nil, fmt.Errorf("proc: %s stdin 管道: %w", name, err) } - stdoutPipe, err := cmd.StdoutPipe() + stdoutR, stdoutW, err := os.Pipe() if err != nil { + stdinR.Close() + stdinW.Close() return nil, fmt.Errorf("proc: %s stdout 管道: %w", name, err) } + cmd.Stdin = stdinR + cmd.Stdout = stdoutW p := &Process{ name: name, bin: bin, dir: opts.Dir, cmd: cmd, - stdin: bufio.NewWriter(stdinPipe), - stdout: stdoutPipe, + stdin: bufio.NewWriter(stdinW), + stdout: stdoutR, + stdinFile: stdinW, + stdoutFile: stdoutR, pending: make(map[uint64]chan *Response), handler: opts.Handler, exited: make(chan struct{}), ready: make(chan struct{}), + waitDone: make(chan struct{}), onExit: opts.OnExit, + sup: opts.Supervisor, shmSize: opts.ShmSize, evtRingSize: opts.EvtRingSize, } if err := cmd.Start(); err != nil { + stdinR.Close() + stdinW.Close() + stdoutR.Close() + stdoutW.Close() return nil, fmt.Errorf("proc: 启动 %s (%s): %w", name, bin, err) } + // 子进程已继承它们,父进程侧关掉对端。 + // stdoutW 必须关:否则子进程死后写端仍被父进程持有,readLoop 永不到 EOF。 + stdinR.Close() + stdoutW.Close() + + // 专职收割协程:这是 cmd.Wait() 的**唯一**调用点。 + // + // 为何不能靠 readLoop 的 EOF:EOF 只说明 stdout 写端全部关闭,而插件 + // fork 出去的孙子进程(browser 拉 chromium、editdoc 拉 python)继承着 + // 同一个 stdout:插件本体死了但孙子还持有写端,EOF 就不来, + // 内核完全感知不到插件已死(进程表里是僵尸,注册表里一切正常)。 + // wait 直接盯进程本身,不受 fd 继承影响。 + p.waiterWG.Add(1) + go p.waitLoop() p.readerWG.Add(1) go p.readLoop() @@ -160,6 +212,9 @@ func Spawn(name, bin string, opts Options) (*Process, error) { p.Kill() return nil, err } + if p.sup != nil { + p.sup.track(p) + } return p, nil } @@ -270,22 +325,54 @@ func (p *Process) readLoop() { log.Printf("[proc] %s 读取 stdout 出错: %v", p.name, err) } - // stdout 关闭(EOF)意味着进程结束——2.5ms 内即可感知(实验 6)。 + // stdout 关闭(EOF)通常意味着进程结束——2.5ms 内即可感知(实验 6)。 + // + // 但 EOF **不是**权威信号:插件 fork 的孙子进程继承同一 stdout 写端时, + // 插件本体死了 EOF 也不会到。真正的死亡判定在 waitLoop。 + // 这里只等 waitLoop 的结果(若进程确实已退,它立即就给)。 + <-p.waitDone p.markExited() } -// markExited 回收进程、唤醒所有等待者、触发 onExit 回调。 +// waitLoop 是内核侧**唯一**的 cmd.Wait() 调用点,每个子进程一根。 // -// 这是「把 panic 捕获换成进程退出检测」的落点(§2.3): -// plugin_health 的 recordCrash / 冷却 / 自愈 / pendingReloads 全部逻辑复用, -// 只是信号源从 recover() 变成进程退出。 +// 为何需要专职协程而不是靠 readLoop 的 EOF: +// 1. **EOF 不等于进程死**。插件用 exec.Command 拉起的孙子进程(browser 拉 +// chromium、editdoc 拉 python)默认继承插件的 stdout。插件被 kill 后 +// 孙子还活着持有写端,readLoop 就永远阻在 Scan 上——内核根本不知道 +// 插件已经死了,工具调用一直超时,自愈也永不触发。 +// 2. **不收割就是僵尸进程**。不调 Wait 的已退出子进程以 Z 状态占着 PID 槽位。 +// 3. **反应速度**。Wait 底层是 wait4(2),内核侧退出即返回(微秒级), +// 比任何轮询健康检查都快,也不消耗 CPU。 +func (p *Process) waitLoop() { + defer p.waiterWG.Done() + p.waitErr = p.cmd.Wait() + close(p.waitDone) + + // 主动拆管道:若孙子进程仍持有 stdout 写端,readLoop 不会自己退, + // 关掉读端逼它从 Scan 里出来(报 file already closed,已预期)。 + if p.stdoutFile != nil { + _ = p.stdoutFile.Close() + } + if p.stdinFile != nil { + _ = p.stdinFile.Close() + } + + p.markExited() +} + +// markExited 唤醒所有等待者、触发 onExit 回调(幂等,两条路径可并发调用)。 +// +// 这是「把 panic 捕获换成进程退出检测」的落点(§2.3)。 +// 注意:不在此处调 cmd.Wait()——它属于 waitLoop,Wait 并非并发安全, +// 两处调会报 "wait: no child processes" 或丢失真实退出码。 func (p *Process) markExited() { p.exitOnce.Do(func() { - waitErr := p.cmd.Wait() - if waitErr != nil { - e := fmt.Errorf("插件进程 %s 异常退出: %w", p.name, waitErr) + <-p.waitDone // 保证 waitErr 可读 + if p.waitErr != nil { + e := fmt.Errorf("插件进程 %s 异常退出: %w", p.name, p.waitErr) p.exitErr.Store(&e) - log.Printf("[proc] %s 退出: %v", p.name, waitErr) + log.Printf("[proc] %s 退出: %v", p.name, p.waitErr) } else { log.Printf("[proc] %s 正常退出", p.name) } @@ -305,6 +392,9 @@ func (p *Process) markExited() { } close(p.exited) + if p.sup != nil { + p.sup.untrack(p.name) + } if p.onExit != nil { p.onExit(p.name, p.ExitError()) } @@ -491,11 +581,14 @@ func (p *Process) Kill() error { return nil } err := p.cmd.Process.Kill() - // 等 readLoop 观察到 EOF 并完成 Wait/清理 + // 等 waitLoop 收割完成。不再在此兜底调 markExited: + // cmd.Wait 只能由 waitLoop 调一次,两处调会报 "wait: no child processes"。 select { case <-p.exited: - case <-time.After(2 * time.Second): - p.markExited() // 兜底:极端情况下强制走清理 + case <-time.After(killReapTimeout): + // SIGKILL 后仍未收割:进程卡在不可中断的内核态(D 状态,如 NFS I/O)。 + // 不能无限等,否则重载路径整体挂死;留日志供定位。 + log.Printf("[proc] %s SIGKILL 后 %v 仍未被收割(进程可能卡在内核态)", p.name, killReapTimeout) } p.readerWG.Wait() if err != nil && !errors.Is(err, os.ErrProcessDone) { diff --git a/internal/plugin/proc/process_test.go b/internal/plugin/proc/process_test.go index 50c2893..d931656 100644 --- a/internal/plugin/proc/process_test.go +++ b/internal/plugin/proc/process_test.go @@ -332,3 +332,150 @@ func TestProcess_SpawnRequiresHandler(t *testing.T) { t.Fatal("缺少 Handler 应报错(插件无法回调内核)") } } + +// 插件死亡但孙子进程仍持有 stdout 写端时,内核必须仍能感知退出。 +// +// 这是「EOF 不等于进程死亡」的回归测试。旧实现只在 readLoop 读到 EOF 后 +// 才 markExited,而 exec.Command 起的孙子进程默认继承插件的 stdout: +// 插件本体退出后写端仍被孙子持有,EOF 永不到来,于是 +// - 在途调用挂到自己的超时; +// - OnExit 不触发 → 崩溃计数、工具摘除、自动重启全都不发生; +// - 进程表里插件已是僵尸,注册表里却一切正常。 +// 生产上 browser 拉 chromium、editdoc 拉 python 正是这个形状。 +// 现在由专职 waitLoop 直接 wait4(2) 判定,不再依赖 fd 生命周期。 +func TestProcess_ExitDetectedDespiteInheritedStdout(t *testing.T) { + if _, err := exec.LookPath("sleep"); err != nil { + t.Skip("环境无 sleep,跳过") + } + bin := buildTestPlugin(t, "forkplugin.go") + + exitCh := make(chan error, 1) + p, err := Spawn("fork", bin, Options{ + Handler: noopHandler, + OnExit: func(name string, err error) { exitCh <- err }, + }) + if err != nil { + t.Fatalf("Spawn: %v", err) + } + defer p.Kill() + + // 让插件本体退出(孙子 sleep 300 仍活着,继续持有 stdout 写端) + if _, callErr := p.Call(MethodToolInvoke, ToolInvokeParams{Name: "die"}); callErr == nil { + t.Error("插件退出时在途调用应返回错误") + } + + select { + case exitErr := <-exitCh: + if exitErr == nil { + t.Error("非零退出码应报告为错误(供崩溃计数使用)") + } + case <-time.After(5 * time.Second): + t.Fatal("孙子进程持有 stdout 时未能感知插件退出——退化回只靠 EOF 判定") + } + + if _, err := p.Call(MethodToolInvoke, ToolInvokeParams{Name: "x"}); !errors.Is(err, ErrProcessExited) { + t.Errorf("退出后调用应返回 ErrProcessExited,实际 %v", err) + } +} + +// Supervisor 台账:握手成功即在册,进程退出即注销。 +func TestSupervisor_TrackAndUntrack(t *testing.T) { + bin := buildTestPlugin(t, "echoplugin.go") + sup := NewSupervisor() + + p, err := Spawn("echo", bin, Options{Handler: noopHandler, Supervisor: sup}) + if err != nil { + t.Fatalf("Spawn: %v", err) + } + if sup.Count() != 1 { + t.Fatalf("握手成功后应在册,实际 %d", sup.Count()) + } + got, ok := sup.Get("echo") + if !ok || got.PID() != p.PID() { + t.Errorf("台账里的进程应是刚 spawn 的那个") + } + list := sup.List() + if len(list) != 1 || !list[0].Alive || list[0].PID != p.PID() { + t.Errorf("List 应报告存活与 PID,实际 %+v", list) + } + + if err := p.Stop(); err != nil { + t.Fatalf("Stop: %v", err) + } + // 退出回调在 markExited 里注销,等它落地 + deadline := time.Now().Add(3 * time.Second) + for sup.Count() != 0 && time.Now().Before(deadline) { + time.Sleep(10 * time.Millisecond) + } + if sup.Count() != 0 { + t.Errorf("进程退出后应注销,实际仍有 %d 个在册", sup.Count()) + } +} + +// StopAll 必须停掉全部在册子进程——内核关停时不留孤儿。 +func TestSupervisor_StopAllLeavesNoSurvivor(t *testing.T) { + bin := buildTestPlugin(t, "echoplugin.go") + sup := NewSupervisor() + + var procs []*Process + for i := 0; i < 3; i++ { + p, err := Spawn(fmt.Sprintf("echo%d", i), bin, Options{Handler: noopHandler, Supervisor: sup}) + if err != nil { + t.Fatalf("Spawn %d: %v", i, err) + } + procs = append(procs, p) + } + if sup.Count() != 3 { + t.Fatalf("应有 3 个在册,实际 %d", sup.Count()) + } + + sup.StopAll(5 * time.Second) + + for _, p := range procs { + select { + case <-p.Exited(): + case <-time.After(2 * time.Second): + t.Errorf("%s 未被 StopAll 停掉(会成为孤儿进程)", p.Name()) + } + } +} + +// 卡死插件(不响应 plugin.stop)必须在 StopAll 的预算内被强杀。 +func TestSupervisor_StopAllKillsUnresponsive(t *testing.T) { + bin := buildTestPlugin(t, "hangplugin.go") + sup := NewSupervisor() + + p, err := Spawn("hang", bin, Options{Handler: noopHandler, Supervisor: sup}) + if err != nil { + t.Fatalf("Spawn: %v", err) + } + + // 预算给足以覆盖 stopGracePeriod,之后剩下的一律 Kill + sup.StopAll(500 * time.Millisecond) + + select { + case <-p.Exited(): + case <-time.After(10 * time.Second): + t.Error("不响应 plugin.stop 的插件应被强制结束,否则 homed 关停会被它拖住") + } +} + +// 关停后完成握手的进程不得留存:立即被结束,不能活过内核。 +func TestSupervisor_TrackAfterCloseKillsProcess(t *testing.T) { + bin := buildTestPlugin(t, "echoplugin.go") + sup := NewSupervisor() + sup.StopAll(time.Second) // 置 closed + + p, err := Spawn("late", bin, Options{Handler: noopHandler, Supervisor: sup}) + if err != nil { + t.Fatalf("Spawn: %v", err) + } + if sup.Count() != 0 { + t.Errorf("关停后不应再纳管新进程,实际在册 %d", sup.Count()) + } + select { + case <-p.Exited(): + case <-time.After(3 * time.Second): + t.Error("关停后冒出的进程应被立即结束") + } +} diff --git a/internal/plugin/proc/supervisor.go b/internal/plugin/proc/supervisor.go new file mode 100644 index 0000000..9936cac --- /dev/null +++ b/internal/plugin/proc/supervisor.go @@ -0,0 +1,163 @@ +package proc + +import ( + "fmt" + "log" + "sort" + "sync" + "time" +) + +// Supervisor 是内核侧**唯一**的子进程台账。 +// +// 为什么必须有它,而不是让每个 Plugin 各自管好自己的 Process: +// +// 1. **没有台账就没有"全部子进程"这个概念**。内核关停时只能遍历 registry 的 +// 插件表逐个 Stop,而 registry 表是按插件名索引的——握手失败、Start 中途 +// 出错、或刚 spawn 还没进表就崩了的进程,registry 根本不知道它们存在, +// 那些进程会变成孤儿(ppid=1)继续跑,还持有共享段映射。 +// 2. **诊断面缺失**。此前 `/api/manager/status` 之类的接口拿不到"实跑几个子进程、 +// 各自 PID 多少、活了多久、崩过几次",运维只能 ps | grep。 +// 3. **收割保证**。每个 Process 自带一根 waitLoop 立即 wait4(2),Supervisor +// 只负责登记/注销与聚合视图;两者配合才能做到"进程一死内核立刻知道"。 +// +// 生命周期:Spawn 成功握手后 track,Process.markExited 里 untrack。 +type Supervisor struct { + mu sync.RWMutex + procs map[string]*Process + // closed 后拒绝新的 track,防止关停竞态里又冒出新进程。 + closed bool +} + +// NewSupervisor 创建空台账。 +func NewSupervisor() *Supervisor { + return &Supervisor{procs: make(map[string]*Process)} +} + +// track 登记一个已握手成功的子进程。 +// +// 同名覆盖是正常情况(重载:旧进程 untrack 早于或晚于新进程 track 都可能, +// 取决于 Kill 与 Spawn 的交错),故不报错,只在真覆盖时留日志。 +func (s *Supervisor) track(p *Process) { + if p == nil { + return + } + s.mu.Lock() + defer s.mu.Unlock() + if s.closed { + // 关停途中还有进程完成握手:立即结束它,不让它活过内核。 + go p.Kill() + return + } + if old, ok := s.procs[p.name]; ok && old != p { + log.Printf("[proc] 台账中 %s 已有 pid=%d,被 pid=%d 覆盖", p.name, old.PID(), p.PID()) + } + s.procs[p.name] = p +} + +// untrack 注销(进程已退出)。只有当表里那一项确实是它时才删, +// 避免重载时新进程被旧进程的退出回调误删。 +func (s *Supervisor) untrack(name string) { + s.mu.Lock() + defer s.mu.Unlock() + delete(s.procs, name) +} + +// Get 按插件名取子进程句柄。 +func (s *Supervisor) Get(name string) (*Process, bool) { + s.mu.RLock() + defer s.mu.RUnlock() + p, ok := s.procs[name] + return p, ok +} + +// Count 返回在册子进程数。 +func (s *Supervisor) Count() int { + s.mu.RLock() + defer s.mu.RUnlock() + return len(s.procs) +} + +// ProcInfo 是单个子进程的运行期快照。 +type ProcInfo struct { + Name string `json:"name"` + PID int `json:"pid"` + Alive bool `json:"alive"` + Bin string `json:"bin"` +} + +// List 返回全部在册子进程的快照(按插件名排序,便于稳定展示)。 +func (s *Supervisor) List() []ProcInfo { + s.mu.RLock() + out := make([]ProcInfo, 0, len(s.procs)) + for name, p := range s.procs { + alive := true + select { + case <-p.Exited(): + alive = false + default: + } + out = append(out, ProcInfo{Name: name, PID: p.PID(), Alive: alive, Bin: p.bin}) + } + s.mu.RUnlock() + sort.Slice(out, func(i, j int) bool { return out[i].Name < out[j].Name }) + return out +} + +// StopAll 停止全部在册子进程:先并发发 plugin.stop 走优雅路径, +// 到期仍在的一律 Kill。 +// +// 这是内核关停时**必须**调的:不调则子进程被 init 收养成孤儿, +// 继续持有共享段映射(段已被内核 unmap,它们下次访问就是 SIGBUS), +// 并且下次 homed 启动时同名插件会与残留进程抢同一份外部资源 +// (qq 的 WS 连接、browser 的 chromium profile 锁)。 +func (s *Supervisor) StopAll(timeout time.Duration) { + s.mu.Lock() + s.closed = true + procs := make([]*Process, 0, len(s.procs)) + for _, p := range s.procs { + procs = append(procs, p) + } + s.mu.Unlock() + + if len(procs) == 0 { + return + } + log.Printf("[proc] 关停 %d 个子进程插件", len(procs)) + + var wg sync.WaitGroup + for _, p := range procs { + wg.Add(1) + go func(pr *Process) { + defer wg.Done() + if err := pr.Stop(); err != nil { + log.Printf("[proc] 停止 %s: %v", pr.Name(), err) + } + }(p) + } + + done := make(chan struct{}) + go func() { wg.Wait(); close(done) }() + + if timeout <= 0 { + timeout = stopGracePeriod * 2 + } + select { + case <-done: + case <-time.After(timeout): + // 优雅停止没在预算内完成:剩下的直接 Kill。 + // 不能无限等——homed 关停被单个卡住的插件拖住比杀掉它更糟。 + var stuck []string + for _, p := range procs { + select { + case <-p.Exited(): + default: + stuck = append(stuck, fmt.Sprintf("%s(pid=%d)", p.Name(), p.PID())) + go p.Kill() + } + } + if len(stuck) > 0 { + log.Printf("[proc] %v 内未优雅退出,强制结束: %v", timeout, stuck) + } + } +} diff --git a/internal/plugin/proc/testdata/forkplugin.go b/internal/plugin/proc/testdata/forkplugin.go new file mode 100644 index 0000000..aa768cc --- /dev/null +++ b/internal/plugin/proc/testdata/forkplugin.go @@ -0,0 +1,71 @@ +//go:build ignore + +// forkplugin 在启动时 fork 一个存活时间比自己长的子进程(继承同一个 stdout), +// 然后在收到 die 工具调用时让自己退出。 +// +// 用途:复现「EOF 不等于进程死亡」这一缺陷。 +// 插件本体死后,孙子进程仍持有 stdout 写端,父进程(homed)的 readLoop +// 永远读不到 EOF——若内核只靠 EOF 判定死亡,就会完全感知不到插件已死: +// 工具调用一直超时、崩溃回调不触发、自动重启永不发生。 +// 生产上 browser 拉 chromium、editdoc 拉 python 都是这个形状。 +package main + +import ( + "bufio" + "encoding/json" + "os" + "os/exec" +) + +type request struct { + ID uint64 `json:"id,omitempty"` + Method string `json:"method"` + Params json.RawMessage `json:"params,omitempty"` +} + +type response struct { + ID uint64 `json:"id"` + Result interface{} `json:"result,omitempty"` + Error string `json:"error,omitempty"` +} + +func main() { + // 孙子进程**只**继承 stdout(本测试的要点),不给 stderr: + // 插件的 stderr 直通到 go test 的捕获管道,孙子抿着它不放会让 + // go test 在测试全部通过后仍等 60s I/O。 + // + // sleep 给 3s:只需在插件本体退出的那一瞬间它还持有写端即可(实际 <200ms), + // 不必拖得更久而拖慢测试。 + child := exec.Command("sleep", "3") + child.Stdout = os.Stdout + _ = child.Start() + + in := bufio.NewScanner(bufio.NewReader(os.Stdin)) + out := bufio.NewWriter(os.Stdout) + send := func(v interface{}) { + b, _ := json.Marshal(v) + out.Write(b) + out.WriteByte('\n') + out.Flush() + } + + for in.Scan() { + var req request + if err := json.Unmarshal(in.Bytes(), &req); err != nil { + continue + } + switch req.Method { + case "handshake": + send(response{ID: req.ID, Result: map[string]interface{}{ + "protocol": 1, "sdk_version": "test", "plugin_name": "fork", "pid": os.Getpid(), + }}) + case "tool.invoke": + // 不回应答,直接退出:模拟插件突然死亡(崩溃/被 kill)。 + os.Exit(7) + default: + if req.ID != 0 { + send(response{ID: req.ID}) + } + } + } +} diff --git a/internal/plugin/registry.go b/internal/plugin/registry.go index b67adfc..2db1830 100644 --- a/internal/plugin/registry.go +++ b/internal/plugin/registry.go @@ -11,6 +11,8 @@ import ( "sort" "strings" "sync" + "sync/atomic" + "time" agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" @@ -59,11 +61,19 @@ func RegisterFactory(name string, factory NativeFactory) { globalFactories.Store(name, factory) } -// PluginToolCleaner 定义插件工具注销接口,由 StageHost 实现。 +// PluginToolCleaner 定义插件注销接口,由 StageHost 实现。 type PluginToolCleaner interface { UnregisterPluginTools(pluginName string) } +// PluginStageCleaner 摘除插件注册的 stage handler,由 StageHost 实现。 +// +// 与 PluginToolCleaner 分开是为了向后兼容:旧 toolCleaner 实现(测试替身) +// 只有 UnregisterPluginTools,经类型断言取 stage 能力,取不到则跳过。 +type PluginStageCleaner interface { + UnregisterPluginStages(pluginName string) int +} + type Registry struct { mu sync.RWMutex plugins map[string]sdk.Plugin @@ -113,6 +123,48 @@ type Registry struct { // lost update 原样复现(§8.4 实测 35.8~36.8%)。 procHostMu sync.Mutex procHost *proc.Host + + // pluginChannels 记录每个插件注册过哪些 IO 通道(输出 device + 输入通道)。 + // + // 不记的后果:子进程插件崩溃后它的 output device 仍在 IOManager 里, + // 模型依旧看到 output_send__ 并调用,只能拿到 ErrProcessExited; + // 重启时 RegisterDevice 又因同名已存在而报 already registered, + // 插件回来了但通道永久指向旧进程的死闭包。 + channelsMu sync.Mutex + pluginChannels map[string]*pluginChannelSet + + // procCrashes 记录子进程插件的崩溃频次,防止崩溃循环无休止重启。 + // + // 与 agent 侧 plugin_health 并存而非重复:后者只能看到工具调用路径上的 + // panic,进程级退出(signal: killed / OOM / 自身 exit)根本不经那里。 + crashMu sync.Mutex + procCrashes map[string]*procCrashRecord + + // shuttingDown 在 StopAll 起始置位,用于冻结自动重启。 + shuttingDown atomic.Bool +} + +// 子进程插件自动重启策略。 +const ( + // procMaxRestarts 是窗口内允许的自动重启次数上限。 + // 超过则停手:再重启也只是重复同一个崩溃,得让人看日志。 + procMaxRestarts = 3 + // procCrashWindow 内无新崩溃则计数归零。 + procCrashWindow = 5 * time.Minute + // procRestartBackoff 是线性退避步长(第 n 次重启前等 n × 此值)。 + procRestartBackoff = time.Second +) + +// procCrashRecord 是单插件的崩溃计数。 +type procCrashRecord struct { + count int + last time.Time +} + +// pluginChannelSet 是单个插件注册过的通道名集合。 +type pluginChannelSet struct { + outputs map[string]bool + inputs map[string]bool } func NewRegistry() *Registry { @@ -123,6 +175,7 @@ func NewRegistry() *Registry { sdkRefs: make(map[string]*sdk.PluginSDK), knownDisabled: make(map[string]bool), pluginHashes: make(map[string]string), + pluginChannels: make(map[string]*pluginChannelSet), } } @@ -219,6 +272,13 @@ func (r *Registry) buildSDK(name string) *sdk.PluginSDK { if regStage == nil { regStage = func(stage sdk.Stage, handler sdk.StageHandler) {} } + // stage handler 注册时带上归属插件名,使卸载/崩溃时能成组摘除。 + // StageHost 实现了 RegisterStageFor;其他实现(测试替身)退回无归属注册。 + if h, ok := r.stageRegistrarFor(); ok { + regStage = func(stage sdk.Stage, handler sdk.StageHandler) { + h(name, stage, handler) + } + } regAPI := r.regAPI if regAPI == nil { regAPI = func(name string) error { return nil } @@ -228,20 +288,25 @@ func (r *Registry) buildSDK(name string) *sdk.PluginSDK { if r.iom == nil { return nil } - return r.iom.RegisterDevice(&channelDevice{ + if err := r.iom.RegisterDevice(&channelDevice{ name: chName, caps: agentIO.OutputCapability(caps), desc: desc, handler: handler, chDef: agentIO.ChannelDef(def), - }) + }); err != nil { + return err + } + r.noteChannel(name, chName, true) + return nil } - regInput := func(name string, def sdk.ChannelDef) error { + regInput := func(chName string, def sdk.ChannelDef) error { if r.iom == nil { return nil } - r.iom.RegisterInputChannel(name, agentIO.ChannelDef(def)) + r.iom.RegisterInputChannel(chName, agentIO.ChannelDef(def)) + r.noteChannel(name, chName, false) return nil } @@ -466,6 +531,92 @@ func (r *Registry) runStopHandlers(name string) { } } +// stageRegistrarFor 取带归属的 stage 注册入口。 +// +// r.regStage 是 cmd/homed 注入的闭包(无插件名参数),而 stageHost 本体同时 +// 以 sdk.ToolSource 存在 r.stageHost 上。能取到 RegisterStageFor 时就直接用它, +// 否则退回无归属注册(测试替身、旧集成方)。 +func (r *Registry) stageRegistrarFor() (func(plugin string, stage sdk.Stage, handler sdk.StageHandler), bool) { + if r.stageHost == nil { + return nil, false + } + if h, ok := r.stageHost.(interface { + RegisterStageFor(plugin string, stage sdk.Stage, handler sdk.StageHandler) + }); ok { + return h.RegisterStageFor, true + } + return nil, false +} + +// noteChannel 记住插件注册了哪个通道,供卸载/崩溃时摘除。 +func (r *Registry) noteChannel(plugin, channel string, output bool) { + if plugin == "" || channel == "" { + return + } + r.channelsMu.Lock() + defer r.channelsMu.Unlock() + set := r.pluginChannels[plugin] + if set == nil { + set = &pluginChannelSet{outputs: map[string]bool{}, inputs: map[string]bool{}} + r.pluginChannels[plugin] = set + } + if output { + set.outputs[channel] = true + } else { + set.inputs[channel] = true + } +} + +// releasePluginChannels 摘除插件注册过的全部 IO 通道,返回摘除的通道名。 +// +// 必须做:不摘除则 ① 模型仍看得到 output_send__ 却永远失败; +// ② 插件重启时 RegisterDevice 报 already registered,新进程的通道注不上, +// 通道永久指向已死进程的闭包。 +func (r *Registry) releasePluginChannels(plugin string) []string { + if plugin == "" { + return nil + } + r.channelsMu.Lock() + set := r.pluginChannels[plugin] + delete(r.pluginChannels, plugin) + r.channelsMu.Unlock() + if set == nil || r.iom == nil { + return nil + } + var released []string + for ch := range set.outputs { + r.iom.UnregisterDevice(ch) + released = append(released, ch) + } + for ch := range set.inputs { + r.iom.UnregisterInputChannel(ch) + if !set.outputs[ch] { + released = append(released, ch) + } + } + sort.Strings(released) + return released +} + +// detachPlugin 把插件在内核侧的全部注册面摸干净:工具 + stage handler + IO 通道。 +// +// 这是「卸载一个插件」的完整含义。之前各路径(Disable/Reload/Remove/ +// StopAndUnload)只调 UnregisterPluginTools,漏了 stage 与通道两项, +// 子进程崩溃路径更是三项都没做。 +func (r *Registry) detachPlugin(name string) { + if r.toolCleaner != nil { + r.toolCleaner.UnregisterPluginTools(name) + if sc, ok := r.toolCleaner.(PluginStageCleaner); ok { + if n := sc.UnregisterPluginStages(name); n > 0 { + log.Printf("[plugin] %s: 摘除 %d 个 stage handler", name, n) + } + } + } + if chans := r.releasePluginChannels(name); len(chans) > 0 { + log.Printf("[plugin] %s: 摘除 IO 通道 %v", name, chans) + } +} + // runOnRemoveHandlers 执行插件注册的删除清理回调(SDK 层),插件 Stop() 之后、从注册表移除前执行。 func (r *Registry) runOnRemoveHandlers(name string) { if sdk, ok := r.sdkRefs[name]; ok { @@ -474,6 +625,10 @@ func (r *Registry) runOnRemoveHandlers(name string) { } func (r *Registry) StopAll() { + // 关停开始即冻结自动重启:否则「Stop 触发退出 → 崩溃判定 → 重新 spawn」 + // 会在内核正在关停时把子进程又拉起来,段已拆而进程还在,直接 SIGBUS。 + r.shuttingDown.Store(true) + r.mu.Lock() for _, p := range r.instances { r.runStopHandlers(p.Name()) @@ -567,6 +722,11 @@ func (r *Registry) ReloadOne(name string) error { } r.mu.Unlock() + // 重载前必须把旧注册面摸干净。不做的后果:loadOne 重新 Start 时 + // RegisterTool 碰上同名旧工具直接报 already registered,新实例的工具一个都注不上; + // stage handler 与 output device 同理——旧闭包指向已死进程,永不退场。 + // (此前只有 agent 的 autoReloadPlugins 在外层手动摸工具,plgreload 路径漏了。) + r.detachPlugin(name) r.closeDynamic(removed) ok := r.loadOne(plgDir, name) @@ -657,9 +817,7 @@ func (r *Registry) Disable(name string) error { r.knownDisabled[name] = true r.mu.Unlock() - if r.toolCleaner != nil { - r.toolCleaner.UnregisterPluginTools(name) - } + r.detachPlugin(name) if r.cfgReg != nil { r.cfgReg.AddDisabledPlugin(name, "system") @@ -767,9 +925,7 @@ func (r *Registry) DisablePlugin(name, by string) error { r.knownDisabled[name] = true r.mu.Unlock() - if r.toolCleaner != nil { - r.toolCleaner.UnregisterPluginTools(name) - } + r.detachPlugin(name) if r.cfgReg != nil { r.cfgReg.AddDisabledPlugin(name, by) @@ -804,9 +960,7 @@ func (r *Registry) StopAndUnload(name string) error { } r.mu.Unlock() - if r.toolCleaner != nil { - r.toolCleaner.UnregisterPluginTools(name) - } + r.detachPlugin(name) r.closeDynamic(unloaded) log.Printf("[plugin] unloaded (config kept): %s", name) return nil @@ -837,9 +991,7 @@ func (r *Registry) RemovePlugin(name string) error { r.runOnRemoveHandlers(name) r.mu.Unlock() - if r.toolCleaner != nil { - r.toolCleaner.UnregisterPluginTools(name) - } + r.detachPlugin(name) if r.cfgReg != nil { r.cfgReg.RemoveDisabledPlugin(name) r.cfgReg.RemovePlugin(name) @@ -851,6 +1003,98 @@ func (r *Registry) RemovePlugin(name string) error { func (r *Registry) ReloadPlugins() (string, error) { return r.Reload(r.plgDir) } +// PluginRuntime 返回单个插件的运行期状态。 +// +// 这是「插件管理器能看到真实死活」的数据源。子进程模型下, +// 「注册表里有条目」不等于「进程还活着」;只读 plugin.json 的旧实现 +// 无法区分两者,插件被 kill 后 WebUI 仍显示“正常”。 +func (r *Registry) PluginRuntime(name string) (sdk.PluginRuntimeInfo, bool) { + if name == "" { + return sdk.PluginRuntimeInfo{}, false + } + + r.mu.RLock() + plg, loaded := r.plugins[name] + _, hasFactory := r.factories[name] + r.mu.RUnlock() + if !hasFactory { + _, hasFactory = globalFactories.Load(name) + } + + installed := loaded || hasFactory + var dirExists bool + if r.plgDir != "" { + if st, err := os.Stat(filepath.Join(r.plgDir, name)); err == nil && st.IsDir() { + dirExists = true + installed = true + } + } + if !installed { + return sdk.PluginRuntimeInfo{}, false + } + + info := sdk.PluginRuntimeInfo{ + Name: name, + Loaded: loaded, + Disabled: r.isDisabled(name), + Builtin: hasFactory, + AutoRestart: r.AutoRestartEnabled(name), + CrashCount: r.crashCount(name), + } + + switch { + case hasFactory: + info.Channel = "builtin" + case dirExists: + info.Channel = detectEntryKind(filepath.Join(r.plgDir, name)).String() + } + + // 子进程插件报真实 PID 与存活;其余形态与 Loaded 同值(无独立进程)。 + type procStatus interface { + PID() int + Alive() bool + } + if ps, ok := plg.(procStatus); ok && loaded { + info.PID = ps.PID() + info.Alive = ps.Alive() + } else { + info.Alive = loaded + } + + if r.stageHost != nil { + for _, def := range r.stageHost.GetToolDefs() { + if def.Plugin == name { + info.Tools = append(info.Tools, def.Name) + } + } + sort.Strings(info.Tools) + } + return info, true +} + +// ListPluginRuntimes 返回全部已知插件的运行期状态(含已安装未加载者)。 +func (r *Registry) ListPluginRuntimes() []sdk.PluginRuntimeInfo { + names := r.ListKnown() + out := make([]sdk.PluginRuntimeInfo, 0, len(names)) + for _, name := range names { + if info, ok := r.PluginRuntime(name); ok { + out = append(out, info) + } + } + return out +} + +// crashCount 读取窗口内的崩溃计数(过期视为 0)。 +func (r *Registry) crashCount(name string) int { + r.crashMu.Lock() + defer r.crashMu.Unlock() + rec := r.procCrashes[name] + if rec == nil || time.Since(rec.last) > procCrashWindow { + return 0 + } + return rec.count +} + // ListKnown 返回所有已知插件(已加载 + 已禁用 + 已安装但未加载)。 func (r *Registry) ListKnown() []string { r.mu.RLock() diff --git a/internal/plugins/pluginmgr/plugin.go b/internal/plugins/pluginmgr/plugin.go index 9ecf0a0..6ad8ff4 100644 --- a/internal/plugins/pluginmgr/plugin.go +++ b/internal/plugins/pluginmgr/plugin.go @@ -14,6 +14,7 @@ import ( "os" "path/filepath" "runtime" + "sort" "strconv" "strings" "sync" @@ -163,7 +164,7 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { s.RegisterTool("plugin_list", sdk.ToolDef{ Name: "plugin_list", - Description: "列出已安装的所有外部插件及其版本", + Description: "列出已安装的所有外部插件及其版本。同时返回运行状态(loaded/alive/pid/崩溃次数),子进程插件死了在此体现为 alive=false。", Parameters: map[string]interface{}{ "type": "object", "properties": map[string]interface{}{}, @@ -172,6 +173,46 @@ func (p *Plugin) registerTools(s *sdk.PluginSDK) { return p.listPlugins() }) + s.RegisterTool("plugin_status", sdk.ToolDef{ + Name: "plugin_status", + Description: "查看插件运行状态:进程是否存活、PID、加载通道、最近崩溃次数、当前注册的工具。" + + "不传 name 则返回全部插件概览。工具调不通时先用它确认插件是否还活着。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "name": map[string]interface{}{ + "type": "string", + "description": "插件名称;缺省返回全部", + }, + }, + }, + }, func(args map[string]interface{}) (interface{}, error) { + name, _ := args["name"].(string) + return p.pluginStatus(name) + }) + + s.RegisterTool("plugin_restart", sdk.ToolDef{ + Name: "plugin_restart", + Description: "重启单个插件(停止后重新加载,保留配置)。适用于:插件进程已死但自动重启被用尽" + + "(plugin_status 的 crash_count 达上限),或换了 plugin.bin 需立即生效。不需重启 homed。", + Parameters: map[string]interface{}{ + "type": "object", + "properties": map[string]interface{}{ + "name": map[string]interface{}{ + "type": "string", + "description": "插件名称", + }, + }, + "required": []string{"name"}, + }, + }, func(args map[string]interface{}) (interface{}, error) { + name, _ := args["name"].(string) + if name == "" { + return map[string]interface{}{"error": "name is required"}, nil + } + return p.restartPlugin(name) + }) + s.RegisterTool("plugin_remove", sdk.ToolDef{ Name: "plugin_remove", Description: "卸载一个已安装的外部插件", @@ -540,6 +581,18 @@ func (p *Plugin) listPlugins() (interface{}, error) { return nil, err } + // 运行期状态一次取齐,避免逐个插件回内核查。 + // + // 为何要带运行期:只读 plugin.json 的旧实现无法区分「已安装」与「正在跑」。 + // 生产上 editdoc 子进程被 kill 后,plugin_list 依旧把它列为正常插件, + // 模型与 WebUI 都看不出异常,只能在调工具时吃一个“进程已退出”。 + runtimes := map[string]sdk.PluginRuntimeInfo{} + if mgr := p.pluginMgr(); mgr != nil { + for _, rt := range mgr.ListPluginRuntimes() { + runtimes[rt.Name] = rt + } + } + var plugins []map[string]interface{} for _, entry := range entries { if !entry.IsDir() { @@ -549,14 +602,27 @@ func (p *Plugin) listPlugins() (interface{}, error) { if err != nil { continue } - plugins = append(plugins, map[string]interface{}{ + item := map[string]interface{}{ "name": m.Name, "version": m.Version, "description": m.Description, "author": m.Author, "entry": m.Entry, "deprecated": m.Deprecated, - }) + } + if rt, ok := runtimes[m.Name]; ok { + item["loaded"] = rt.Loaded + item["alive"] = rt.Alive + item["disabled"] = rt.Disabled + item["channel"] = rt.Channel + if rt.PID > 0 { + item["pid"] = rt.PID + } + if rt.CrashCount > 0 { + item["crash_count"] = rt.CrashCount + } + } + plugins = append(plugins, item) } if plugins == nil { plugins = []map[string]interface{}{} @@ -564,6 +630,84 @@ func (p *Plugin) listPlugins() (interface{}, error) { return plugins, nil } +// pluginMgr 取内核插件管理面(可能为 nil:单测/未注入)。 +func (p *Plugin) pluginMgr() sdk.PluginManager { + if p.sdk == nil { + return nil + } + return p.sdk.PluginMgr() +} + +// pluginStatus 返回插件运行期状态(进程存活/PID/崩溃计数/工具清单)。 +// +// 这是子进程化后插件管理器必须补上的一块:以前插件与内核同进程, +// “加载了”就等于“能用”;现在插件是独立进程,两者不再等价。 +func (p *Plugin) pluginStatus(name string) (interface{}, error) { + mgr := p.pluginMgr() + if mgr == nil { + return nil, fmt.Errorf("内核插件管理面不可用") + } + if name != "" { + info, ok := mgr.PluginRuntime(name) + if !ok { + return nil, fmt.Errorf("plugin %q not found", name) + } + return info, nil + } + + all := mgr.ListPluginRuntimes() + // 汇总一行:让模型不用自己数就能看出“有东西挂了”。 + var loaded, dead int + var unhealthy []string + for _, rt := range all { + if rt.Loaded { + loaded++ + } + if rt.Loaded && !rt.Alive { + dead++ + unhealthy = append(unhealthy, rt.Name) + continue + } + if rt.CrashCount > 0 { + unhealthy = append(unhealthy, fmt.Sprintf("%s(崩溃%d次)", rt.Name, rt.CrashCount)) + } + } + sort.Strings(unhealthy) + return map[string]interface{}{ + "total": len(all), + "loaded": loaded, + "dead": dead, + "unhealthy": unhealthy, + "plugins": all, + }, nil +} + +// restartPlugin 重启单个插件(保留配置)。 +// +// 与 plgreload 的区别:后者按入口文件 hash 增量重载,二进制没改就不动; +// 而进程被 kill 时二进制正是没改的,所以一定要有一个无条件重启的入口。 +func (p *Plugin) restartPlugin(name string) (interface{}, error) { + mgr := p.pluginMgr() + if mgr == nil { + return nil, fmt.Errorf("内核插件管理面不可用") + } + if _, ok := mgr.PluginRuntime(name); !ok { + return nil, fmt.Errorf("plugin %q not found", name) + } + if mgr.IsPluginDisabled(name) { + return nil, fmt.Errorf("plugin %s 已被禁用,请先启用再重启", name) + } + if err := mgr.ReloadOne(name); err != nil { + return nil, fmt.Errorf("restart %s: %w", name, err) + } + info, _ := mgr.PluginRuntime(name) + return map[string]interface{}{ + "status": "restarted", + "name": name, + "runtime": info, + }, nil +} + func (p *Plugin) removePlugin(name string) (interface{}, error) { // 内置插件只能禁用不能卸载:目录下无产物,且从注册表删除会破坏内核依赖。 if p.sdk != nil && p.sdk.PluginMgr() != nil && p.sdk.PluginMgr().IsBuiltinPlugin(name) { diff --git a/internal/plugins/pluginmgr/upgrade_test.go b/internal/plugins/pluginmgr/upgrade_test.go index 063ac4f..2fafc1c 100644 --- a/internal/plugins/pluginmgr/upgrade_test.go +++ b/internal/plugins/pluginmgr/upgrade_test.go @@ -81,6 +81,10 @@ func (f *fakePluginMgr) PluginMetas() map[string]sdk.PluginMeta { return map[string]sdk.PluginMeta{} } func (f *fakePluginMgr) PluginDir() string { return "" } +func (f *fakePluginMgr) PluginRuntime(string) (sdk.PluginRuntimeInfo, bool) { + return sdk.PluginRuntimeInfo{}, false +} +func (f *fakePluginMgr) ListPluginRuntimes() []sdk.PluginRuntimeInfo { return nil } // buildHmap 构造一个最小 .hmap 包。 func buildHmap(t *testing.T, name, version string) []byte { diff --git a/internal/plugins/webui/handler_plugin_test.go b/internal/plugins/webui/handler_plugin_test.go index e768ddd..d911e4f 100644 --- a/internal/plugins/webui/handler_plugin_test.go +++ b/internal/plugins/webui/handler_plugin_test.go @@ -13,6 +13,7 @@ type mockPluginMgr struct { builtins map[string]bool disabled []sdk.DisabledPluginInfo isDisabled map[string]bool + runtimes map[string]sdk.PluginRuntimeInfo removed []string reloadN int @@ -57,6 +58,17 @@ func (m *mockPluginMgr) PluginMetas() map[string]sdk.PluginMeta { return nil } func (m *mockPluginMgr) PluginDir() string { return "" } +func (m *mockPluginMgr) PluginRuntime(name string) (sdk.PluginRuntimeInfo, bool) { + info, ok := m.runtimes[name] + return info, ok +} +func (m *mockPluginMgr) ListPluginRuntimes() []sdk.PluginRuntimeInfo { + out := make([]sdk.PluginRuntimeInfo, 0, len(m.runtimes)) + for _, v := range m.runtimes { + out = append(out, v) + } + return out +} // newHandlerWithMock 构造带 mock PluginManager 的 Handler(绕过 SDK 组装)。 func newHandlerWithMock(m *mockPluginMgr) *Handler { diff --git a/internal/sdk/plugin.go b/internal/sdk/plugin.go index a76789a..7396f6d 100644 --- a/internal/sdk/plugin.go +++ b/internal/sdk/plugin.go @@ -65,6 +65,34 @@ type PluginMeta struct { NameEn string `json:"name_en"` } +// PluginRuntimeInfo 是插件的**运行期**状态,与 plugin.json 里的静态元数据相对。 +// +// 为何需要:子进程插件的进程可能已经死了而注册表里还有条目(或反过来, +// 崩溃摘除后注册表已无条目但目录还在)。此前 plugin_list / GET /plugins +// 只读 plugin.json,无论插件死活都返回同一份内容——WebUI 与模型都看不出 +// 「已安装」与「正在运行」的区别,插件被 kill 后只表现为工具静默失败。 +type PluginRuntimeInfo struct { + Name string `json:"name"` + // Loaded 表示注册表中存在该插件实例。 + Loaded bool `json:"loaded"` + // Disabled 表示插件被显式禁用(不该运行)。 + Disabled bool `json:"disabled"` + // Builtin 表示编译期内置插件(无独立进程)。 + Builtin bool `json:"builtin"` + // Channel 是加载通道:proc(子进程)/ lua / builtin。 + Channel string `json:"channel"` + // PID 是子进程插件的进程号;非子进程或已退出为 0。 + PID int `json:"pid"` + // Alive 表示子进程仍存活;非子进程插件与 Loaded 同值。 + Alive bool `json:"alive"` + // CrashCount 是最近窗口内的崩溃次数(0 表示健康)。 + CrashCount int `json:"crash_count"` + // AutoRestart 表示崩溃后内核是否会自动拉起。 + AutoRestart bool `json:"auto_restart"` + // Tools 是该插件当前注册在内核里的工具名。 + Tools []string `json:"tools,omitempty"` +} + type PluginManager interface { ListLoadedPlugins() []string ListDisabledPlugins() []DisabledPluginInfo @@ -86,6 +114,11 @@ type PluginManager interface { ReloadOne(name string) error PluginMetas() map[string]PluginMeta PluginDir() string + // PluginRuntime 返回单个插件的运行期状态(进程存活 / PID / 崩溃计数)。 + // 未安装的插件返回零值 + false。 + PluginRuntime(name string) (PluginRuntimeInfo, bool) + // ListPluginRuntimes 返回全部已加载插件的运行期状态。 + ListPluginRuntimes() []PluginRuntimeInfo } type PluginSDK struct {