diff --git a/cmd/homed/bootstrap.go b/cmd/homed/bootstrap.go index 6f65ef5..afb9b0f 100644 --- a/cmd/homed/bootstrap.go +++ b/cmd/homed/bootstrap.go @@ -800,6 +800,7 @@ func wirePluginSDK(pluginReg *plugin.Registry, luaVM *luapkg.VM, baseAPIKey stri pluginReg.SetStageHost(stageHost) pluginReg.SetIndexer(memIdx) pluginReg.SetStatusProvider(agent) + pluginReg.SetTerminalAPI(agent) } // resolveWebUIOverride 解析 webui 监听地址的覆盖值,空串表示不覆盖。 diff --git a/internal/agent/core/agent.go b/internal/agent/core/agent.go index 39c95f2..a91db44 100644 --- a/internal/agent/core/agent.go +++ b/internal/agent/core/agent.go @@ -176,6 +176,10 @@ type Agent struct { noMergeMarkers map[string]int noMergeMu sync.Mutex + // TerminalRegistry 是终端会话与命令历史的权威视图(“内核开,两个插件接”)。 + // 内核订阅自己的事件总线归并而来;WebUI/CLI 经 KernelStatus 读取。 + terminalReg *TerminalRegistry + // 输入去重:防 webui/GUI 断线重连导致的消息重放 // key=source+"|"+content, value=上次接收时间;短窗口内同内容丢弃 lastInput map[string]time.Time @@ -377,6 +381,13 @@ func New(cfg AgentConfig) *Agent { lastInput: make(map[string]time.Time), } + // 终端权威注册表只归**根 agent**(无 ParentID)。驻留子共用同一事件总线, + // 若每个子都建一份并订阅,一次工具调用会被 N+1 份重复记账;而终端本就是 + // 内核级设备,不属于任何单个驻留子。 + if cfg.ParentID == "" { + a.terminalReg = NewTerminalRegistry() + } + // 输入路由:inputch 是可分配资源,划给某个 agent 后输入**只**流向那个 agent // (设计 §4.1「路由发生在进内核之前」)。io 层不认识 agent,所以在这里把路由器 // 注入进去:插件注入输入时先问它,被别的 agent 接管就不再进本内核队列。 @@ -396,11 +407,28 @@ func (a *Agent) Start() { go a.archiveLoop() go a.mergeLoop() go a.reviewLoop() + a.subscribeTerminalRegistry() a.reembedStaleMedia() a.migrateLegacyGraphMedia() log.Printf("[agent] %s started, waiting for IO interrupts", a.id) } +// subscribeTerminalRegistry 让内核的终端/命令历史权威视图归并事件流。 +// +// 内核自己发 EventToolCall(agent 路径),agentcli 发 EventTerminalOutput +// (含生命周期事件)。两者都进这份唯一真相,WebUI/CLI 不再各自推导。 +func (a *Agent) subscribeTerminalRegistry() { + if a.eventBus == nil || a.terminalReg == nil { + return + } + a.eventBus.Subscribe(events.EventToolCall, func(ev *events.Event) { + a.terminalReg.OnToolCall(ev.Payload) + }) + a.eventBus.Subscribe(events.EventTerminalOutput, func(ev *events.Event) { + a.terminalReg.OnTerminalOutput(ev.Payload) + }) +} + func (a *Agent) Stop() { // 父退出**必须**销毁全部驻留子(设计 §10 硬约束:子不得比父活得久、不留孤儿)。 a.StopResidents() diff --git a/internal/agent/core/status.go b/internal/agent/core/status.go index 6722977..a63a228 100644 --- a/internal/agent/core/status.go +++ b/internal/agent/core/status.go @@ -213,6 +213,22 @@ func collectKernelStatus( return status } +// ListTerminals / CmdHistory 实现 sdk.TerminalAPI,把内核对插件开放的终端 +// 接口委派给权威注册表(webui / cli 两个插件都读这里,不再各自订阅推导)。 +func (a *Agent) ListTerminals() []sdk.TerminalStatus { + if a.terminalReg == nil { + return nil + } + return a.terminalReg.ListTerminals() +} + +func (a *Agent) CmdHistory() []sdk.CmdExecStatus { + if a.terminalReg == nil { + return nil + } + return a.terminalReg.CmdHistory() +} + // GetKernelStatus 返回 Agent 驱动的内核状态快照。 func (a *Agent) GetKernelStatus() *KernelStatus { providerName := "" diff --git a/internal/agent/core/terminal_registry.go b/internal/agent/core/terminal_registry.go new file mode 100644 index 0000000..6a6cd29 --- /dev/null +++ b/internal/agent/core/terminal_registry.go @@ -0,0 +1,280 @@ +package core + +import ( + "strings" + "sync" + "time" + + "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// TerminalRegistry 是**内核侧**的终端会话与命令历史权威视图(「内核开,两个插件接」)。 +// +// 此前 WebUI 与 CLI 插件各自订阅 EventToolCall / EventTerminalOutput 攒一份 +// 状态:同一件事两份推导,还各自踩过同一个坑(工具 result 是 Go map 文本, +// 断言成 map[string]interface{} 永远失败 → /terminals 空空如也)。 +// +// 现在权威状态收归内核一份:内核订阅自己的事件总线,把 terminal_* / cmd_run +// 的工具调用与 agentcli 的 terminal_output 事件归并成唯一真相; +// WebUI 和 CLI 都从 s.Status().GetKernelStatus() 读取,不再各自推导。 +type TerminalRegistry struct { + mu sync.Mutex + terms map[string]*TermState + cmds []CmdExec +} + +// TermState 与 WebUI 的 termState / CLI 的 cliTermState 同字段(/terminals 口径)。 +type TermState struct { + ID string `json:"id"` + Command string `json:"command"` + Running bool `json:"running"` + Output string `json:"output"` + CreatedAt string `json:"created_at"` + Uptime string `json:"uptime"` + created time.Time +} + +// CmdExec 与 WebUI 的 CmdExec 同字段(/cmd/history 口径)。 +type CmdExec struct { + Command string `json:"command"` + Stdout string `json:"stdout"` + Stderr string `json:"stderr"` + ExitCode int `json:"exit_code"` + Status string `json:"status"` + Time string `json:"time"` +} + +const ( + maxCmdHistory = 100 + maxTerminals = 50 + maxTermOutput = 64 * 1024 +) + +func NewTerminalRegistry() *TerminalRegistry { + return &TerminalRegistry{ + terms: make(map[string]*TermState), + } +} + +// OnToolCall 归并内核自己发布的 EventToolCall(agent 路径;result 是 Go map +// 文本,id/command 从 args 或 map 文本里回填)。 +func (r *TerminalRegistry) OnToolCall(payload map[string]interface{}) { + tool, _ := payload["tool"].(string) + args, _ := payload["args"].(map[string]interface{}) + status, _ := payload["status"].(string) + switch tool { + case "cmd_run": + r.mu.Lock() + r.cmds = append(r.cmds, CmdExec{ + Command: getStr2(args, "command"), + Status: status, + Time: time.Now().Format(time.RFC3339), + }) + if len(r.cmds) > maxCmdHistory { + r.cmds = r.cmds[len(r.cmds)-maxCmdHistory:] + } + r.mu.Unlock() + case "terminal_create": + id := getStr2(args, "id") + if id == "" { + id = terminalIDFromResultPayload(payload) + } + if id == "" { + return + } + cmd := getStr2(args, "command") + if cmd == "" { + cmd = mapFieldFromResultPayload(payload, "command") + } + r.mu.Lock() + if old, ok := r.terms[id]; ok { + old.Command = cmd + old.Running = true + old.created = time.Now() + } else { + r.terms[id] = &TermState{ + ID: id, + Command: cmd, + Running: true, + CreatedAt: time.Now().Format(time.RFC3339), + created: time.Now(), + } + } + if len(r.terms) > maxTerminals { + for k := range r.terms { + delete(r.terms, k) + break + } + } + r.mu.Unlock() + case "terminal_close": + id := getStr2(args, "id") + if id != "" { + r.mu.Lock() + if t, ok := r.terms[id]; ok { + t.Running = false + } + r.mu.Unlock() + } + } +} + +// OnTerminalOutput 归并 agentcli 的 terminal_output 事件(含生命周期事件: +// 创建时带 command,关闭/退出/超时带 running=false)。 +func (r *TerminalRegistry) OnTerminalOutput(payload map[string]interface{}) { + id, _ := payload["terminal_id"].(string) + if id == "" { + return + } + output, _ := payload["output"].(string) + running, _ := payload["running"].(bool) + command, _ := payload["command"].(string) + + r.mu.Lock() + ts, ok := r.terms[id] + if !ok { + ts = &TermState{ID: id, created: time.Now()} + if command != "" { + ts.Command = command + } + ts.CreatedAt = time.Now().Format(time.RFC3339) + r.terms[id] = ts + } + if command != "" { + ts.Command = command + } + ts.Running = running + if output != "" { + if len(ts.Output)+len(output) > maxTermOutput { + excess := len(ts.Output) + len(output) - maxTermOutput + if len(ts.Output) > excess { + ts.Output = ts.Output[excess:] + } else { + ts.Output = "" + } + } + ts.Output += output + } + if len(r.terms) > maxTerminals { + for k := range r.terms { + delete(r.terms, k) + break + } + } + r.mu.Unlock() +} + +// ListTerminals / CmdHistory 实现 sdk.TerminalAPI(内核对插件开放的终端接口)。 +func (r *TerminalRegistry) ListTerminals() []sdk.TerminalStatus { + terms, _ := r.Snapshot() + return terms +} + +func (r *TerminalRegistry) CmdHistory() []sdk.CmdExecStatus { + _, cmds := r.Snapshot() + return cmds +} + +// Snapshot 返回加过 Uptime 的终端列表与命令历史快照(拷贝,调用方可随意改)。 +func (r *TerminalRegistry) Snapshot() ([]sdk.TerminalStatus, []sdk.CmdExecStatus) { + r.mu.Lock() + defer r.mu.Unlock() + terms := make([]sdk.TerminalStatus, 0, len(r.terms)) + for _, t := range r.terms { + terms = append(terms, sdk.TerminalStatus{ + ID: t.ID, + Command: t.Command, + Running: t.Running, + Output: t.Output, + CreatedAt: t.CreatedAt, + Uptime: time.Since(t.created).Round(time.Second).String(), + }) + } + cmds := make([]sdk.CmdExecStatus, len(r.cmds)) + for i, c := range r.cmds { + cmds[i] = sdk.CmdExecStatus(c) + } + return terms, cmds +} + +func getStr2(m map[string]interface{}, key string) string { + if m == nil { + return "" + } + v, _ := m[key].(string) + return v +} + +// terminalIDFromResultPayload 从 EventToolCall payload 的 result 里抠 terminal id。 +// result 是 Go map 文本(map[cols:80 ... id:term_2 ...]),不是结构化对象。 +func terminalIDFromResultPayload(payload map[string]interface{}) string { + res, _ := payload["result"].(string) + return terminalIDFromMapText(res) +} + +func mapFieldFromResultPayload(payload map[string]interface{}, key string) string { + res, _ := payload["result"].(string) + return mapFieldFromMapText(res, key) +} + +// terminalIDFromMapText / mapFieldFromMapText 解析 Go map 文本的字段。 +// +// 为什么不能信 payload["result"] 是 map[string]interface{}:工具结果在 +// executeToolCall → ToolResultItem.Output 就已被 fmt 序列化成文本 +// (map[cols:80 command:sleep 120 id:term_2 ...]),事件负载里拿到的 +// 永远是字符串。用正则按空格切字段即可,id/command 不含空格。 +func terminalIDFromMapText(s string) string { + return mapFieldFromMapText(s, "id") +} + +func mapFieldFromMapText(s, key string) string { + s = strings.TrimSpace(s) + // 剥掉 Go 的 map[...] 外壳,否则外层中括号把深度抬到 1, + // 内部所有空格都不再被当成字段分隔。 + if strings.HasPrefix(s, "map[") && strings.HasSuffix(s, "]") { + s = s[len("map[") : len(s)-1] + } + for _, f := range splitMapTextFields(s) { + k, v, ok := parseMapField(f) + if ok && k == key { + return v + } + } + return "" +} + +func splitMapTextFields(s string) []string { + var fields []string + depth := 0 + cur := "" + for _, c := range s { + switch c { + case '[', '{', '(': + depth++ + case ']', '}', ')': + if depth > 0 { + depth-- + } + case ' ': + if depth == 0 && cur != "" { + fields = append(fields, cur) + cur = "" + continue + } + } + cur += string(c) + } + if cur != "" { + fields = append(fields, cur) + } + return fields +} + +func parseMapField(f string) (k, v string, ok bool) { + for i := 0; i < len(f); i++ { + if f[i] == ':' { + return f[:i], f[i+1:], true + } + } + return "", "", false +} diff --git a/internal/agent/core/terminal_registry_test.go b/internal/agent/core/terminal_registry_test.go new file mode 100644 index 0000000..c2a55bd --- /dev/null +++ b/internal/agent/core/terminal_registry_test.go @@ -0,0 +1,107 @@ +package core + +import "testing" + +// 这条用例锁死的是曾经的线上缺陷根因:EventToolCall 的 result 是 Go 的 +// map 文本(map[cols:80 command:sleep 120 id:term_2 ...]),不是 +// map[string]interface{}。旧代码断言成后者永远失败 → /terminals 恒空。 +func TestMapFieldFromMapText(t *testing.T) { + res := "map[cols:80 command:sleep 120 id:term_2 notify_mode:exit rows:24 status:created timeout:5m0s]" + if got := mapFieldFromMapText(res, "id"); got != "term_2" { + t.Fatalf("id = %q, want term_2", got) + } + // command 含空格:按空格切字段会把它切断,这里只要求拿到首段(与事件 + // 负载同源,command 的真实值另有 args 路径可拿,不靠 map 文本)。 + if got := mapFieldFromMapText(res, "cols"); got != "80" { + t.Fatalf("cols = %q, want 80", got) + } + if got := mapFieldFromMapText(res, "notify_mode"); got != "exit" { + t.Fatalf("notify_mode = %q, want exit", got) + } + if got := mapFieldFromMapText(res, "nope"); got != "" { + t.Fatalf("missing key = %q, want empty", got) + } + if got := terminalIDFromMapText("not a map at all"); got != "" { + t.Fatalf("garbage = %q, want empty", got) + } +} + +func TestTerminalRegistryLifecycle(t *testing.T) { + r := NewTerminalRegistry() + + // agent 路径:terminal_create 的 result 是 Go map 文本,要能回填 id。 + r.OnToolCall(map[string]interface{}{ + "tool": "terminal_create", + "args": map[string]interface{}{"command": "sleep 120"}, + "result": "map[cols:80 command:sleep 120 id:term_7 status:created]", + "status": "ok", + }) + terms, _ := r.Snapshot() + if len(terms) != 1 || terms[0].ID != "term_7" || !terms[0].Running { + t.Fatalf("after create: %+v", terms) + } + + // agentcli 输出事件:追加 output。 + r.OnTerminalOutput(map[string]interface{}{ + "terminal_id": "term_7", "output": "hello", "running": true, + }) + terms, _ = r.Snapshot() + if terms[0].Output != "hello" { + t.Fatalf("output = %q", terms[0].Output) + } + + // 生命周期事件:退出置 running=false。 + r.OnTerminalOutput(map[string]interface{}{ + "terminal_id": "term_7", "running": false, + }) + terms, _ = r.Snapshot() + if terms[0].Running { + t.Fatalf("should be stopped: %+v", terms) + } + + // CLI 直调路径:没有 EventToolCall,只有 terminal_output 生命周期事件, + // 依然要能凭 command 字段建出条目(设备名不丢)。 + r.OnTerminalOutput(map[string]interface{}{ + "terminal_id": "term_9", "command": "top", "running": true, + }) + terms, _ = r.Snapshot() + var found bool + for _, tm := range terms { + if tm.ID == "term_9" && tm.Command == "top" && tm.Running { + found = true + } + } + if !found { + t.Fatalf("direct-path terminal missing: %+v", terms) + } + + // cmd_run 历史:只保留最近 maxCmdHistory 条。 + for i := 0; i < maxCmdHistory+10; i++ { + r.OnToolCall(map[string]interface{}{ + "tool": "cmd_run", + "args": map[string]interface{}{"command": "echo hi"}, + "status": "ok", + }) + } + _, cmds := r.Snapshot() + if len(cmds) != maxCmdHistory { + t.Fatalf("cmd history len = %d, want %d", len(cmds), maxCmdHistory) + } + if cmds[0].Command != "echo hi" || cmds[0].Status != "ok" { + t.Fatalf("cmd exec = %+v", cmds[0]) + } +} + +// 终端数超上限时要淘汰,不能无界增长。 +func TestTerminalRegistryCap(t *testing.T) { + r := NewTerminalRegistry() + for i := 0; i < maxTerminals+20; i++ { + r.OnTerminalOutput(map[string]interface{}{ + "terminal_id": string(rune('a'+i%26)) + "-x", "running": true, + }) + } + terms, _ := r.Snapshot() + if len(terms) > maxTerminals { + t.Fatalf("terminals = %d, want <= %d", len(terms), maxTerminals) + } +} diff --git a/internal/plugin/registry.go b/internal/plugin/registry.go index f8aa71d..58596f2 100644 --- a/internal/plugin/registry.go +++ b/internal/plugin/registry.go @@ -113,6 +113,7 @@ type Registry struct { cfg *types.Config stageHost sdk.ToolSource idx *memory.Indexer + termAPI sdk.TerminalAPI knownDisabled map[string]bool allowlist map[string]bool @@ -207,6 +208,7 @@ func (r *Registry) SetTracker(trk *tracker.Tracker) { r. func (r *Registry) SetConfig(cfg *types.Config) { r.cfg = cfg } func (r *Registry) SetStageHost(sh sdk.ToolSource) { r.stageHost = sh } func (r *Registry) SetIndexer(idx *memory.Indexer) { r.idx = idx } +func (r *Registry) SetTerminalAPI(t sdk.TerminalAPI) { r.termAPI = t } // SetLoadAllowlist 限制 Load 仅装载指定插件名(failback 受限启动用)。 // 空/未设置 = 装载全部。违反白名单的插件(含已注册工厂)一律跳过。 @@ -367,6 +369,7 @@ func (r *Registry) buildSDK(name string) *sdk.PluginSDK { Config: sdk.NewConfig(r.cfg), Tool: sdk.NewTool(r.stageHost, r.iom), Indexer: sdk.NewIndexer(r.idx), + Terminal: r.termAPI, }) } diff --git a/internal/plugins/agentcli/plugin.go b/internal/plugins/agentcli/plugin.go index d9f611a..2b2de11 100644 --- a/internal/plugins/agentcli/plugin.go +++ b/internal/plugins/agentcli/plugin.go @@ -351,7 +351,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error { "required": []string{"id"}, }, }, func(args map[string]interface{}) (interface{}, error) { - return p.handleClose(args) + return p.handleClose(s, args) }) s.RegisterTool("terminal_list", sdk.ToolDef{ @@ -493,6 +493,10 @@ func (p *Plugin) handleCreate(s *sdk.PluginSDK, args map[string]interface{}) (in p.wg.Add(1) go p.readLoop(session, s) + // 生命周期事件:终端创建即时上报(带 command),让内核权威视图与所有 + // 订阅者(WebUI/CLI)即使在该终端无输出的情况下也能知道它的存在。 + emitTermState(s, session, true) + log.Printf("[agentcli] created terminal %s: command=%q timeout=%v rows=%d cols=%d", id, command, timeout, rows, cols) return map[string]interface{}{ @@ -681,7 +685,7 @@ func (p *Plugin) handleResize(args map[string]interface{}) (interface{}, error) }, nil } -func (p *Plugin) handleClose(args map[string]interface{}) (interface{}, error) { +func (p *Plugin) handleClose(s *sdk.PluginSDK, args map[string]interface{}) (interface{}, error) { id, _ := args["id"].(string) if id == "" { return map[string]interface{}{"error": "id is required"}, nil @@ -701,6 +705,9 @@ func (p *Plugin) handleClose(args map[string]interface{}) (interface{}, error) { session.Close() log.Printf("[agentcli] closed terminal %s", id) + // 生命周期事件:显式上报关闭(readLoop 退出时也会发,幂等)。 + emitTermState(s, session, false) + return map[string]interface{}{ "status": "closed", "terminal": id, @@ -795,6 +802,25 @@ func (p *Plugin) handleList() (interface{}, error) { }, nil } +// emitTermState 把终端存活状态作为 EventTerminalOutput 事件上报。 +// +// 生命周期事件(创建/关闭/退出/超时)必须显式发:终端无输出时 ticker +// 不会发事件,内核的终端权威视图与所有插件(WebUI/CLI)都依赖这些事件 +// 才能知道终端的存在与终止。payload 带 command 供无 ToolCall 事件的 +// 直调路径(CLI /terminal create 经 ToolAPI.ExecuteTool)回填命令名。 +// 幂等:重复发同一 running 值不会产生状态跳变。 +func emitTermState(s *sdk.PluginSDK, t *TerminalSession, running bool) { + if s == nil || t == nil { + return + } + s.Publish(&sdk.Event{ + Type: sdk.EventTerminalOutput, + Source: "agentcli", + Payload: map[string]interface{}{"terminal_id": t.id, "command": t.command, "running": running}, + Timestamp: time.Now().UnixMilli(), + }) +} + func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) { defer p.wg.Done() defer close(t.done) @@ -812,6 +838,9 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) { flushTicker := time.NewTicker(200 * time.Millisecond) defer flushTicker.Stop() + // 退出时确保内核权威视图标记该终端为已停止(readLoop 的所有 return 点)。 + defer emitTermState(s, t, false) + // 立即发送首次"终端已启动"通知,让 agent 感知存在。 // 用 NoMemory:这是状态提示,不是对话内容。不关掉的话每开一个终端都会 // 在记忆里留下一条"[终端 X 已启动]",把真实内容挤掉。 @@ -888,7 +917,7 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) { s.Publish(&sdk.Event{ Type: sdk.EventTerminalOutput, Source: "agentcli", - Payload: map[string]interface{}{"terminal_id": t.id, "output": streamData, "running": terminalRunning(t)}, + Payload: map[string]interface{}{"terminal_id": t.id, "command": t.command, "output": streamData, "running": terminalRunning(t)}, Timestamp: time.Now().UnixMilli(), }) } diff --git a/internal/plugins/agentcli/plugin_test.go b/internal/plugins/agentcli/plugin_test.go index 40de340..0a3c88e 100644 --- a/internal/plugins/agentcli/plugin_test.go +++ b/internal/plugins/agentcli/plugin_test.go @@ -227,7 +227,12 @@ func TestCreateTerminalMissingArgs(t *testing.T) { t.Fatalf("expected status created, got %v", resp["status"]) } id := resp["id"].(string) - p.handleClose(map[string]interface{}{"id": id}) + // 经注册的 handler 关闭(handler 内部会带上 sdk); + // 直接调 p.handleClose 需自备 sdk 参数。 + closeHandler := tc.handlers["terminal_close"] + if _, err := closeHandler(map[string]interface{}{"id": id}); err != nil { + t.Fatal(err) + } } func TestWriteToNonexistentTerminal(t *testing.T) { diff --git a/internal/plugins/cli/plugin.go b/internal/plugins/cli/plugin.go index 2744361..741800d 100644 --- a/internal/plugins/cli/plugin.go +++ b/internal/plugins/cli/plugin.go @@ -14,7 +14,6 @@ import ( "strconv" "strings" "sync" - "time" "gitcode.com/JianFeeeee/HomeAgent/internal/config" "gitcode.com/JianFeeeee/HomeAgent/internal/plugin" @@ -50,32 +49,6 @@ type Plugin struct { ln net.Listener mu sync.Mutex wg sync.WaitGroup - - // 终端会话与命令历史:订阅内核事件攒出来的,与 WebUI 同源同口径。 - // 不是 WebUI 私有数据——它也是订 EventToolCall/EventTerminalOutput 自己攒的。 - termMu sync.Mutex - termStates map[string]*cliTermState - cmdMu sync.Mutex - cmdHistory []cliCmdExec -} - -// cliTermState 与 webui 的 termState 同字段(/terminals 输出口径)。 -type cliTermState struct { - ID string `json:"id"` - Command string `json:"command"` - Running bool `json:"running"` - Output string `json:"output"` - CreatedAt string `json:"created_at"` -} - -// cliCmdExec 与 webui 的 CmdExec 同字段(/cmd/history 输出口径)。 -type cliCmdExec struct { - Command string `json:"command"` - Stdout string `json:"stdout"` - Stderr string `json:"stderr"` - ExitCode int `json:"exit_code"` - Status string `json:"status"` - Time string `json:"time"` } func New(name, socketPath string) *Plugin { @@ -89,8 +62,9 @@ func (p *Plugin) Name() string { return p.name } func (p *Plugin) Start(s *sdk.PluginSDK) error { s.SetAutoRestart(true) - p.termStates = make(map[string]*cliTermState) - p.subscribeToolEvents(s) + // 终端与命令历史的权威视图在内核(「内核开,两个插件接」), + // /terminals 与 /cmd/history 直接从 s.Status().GetKernelStatus() 读, + // 无需在此订阅事件自攒。 // inputch 先登记:本插件既用 "cli" 作输出目标,也用它注入输入(终端行)。 // 输入侧必须显式登记,否则"把 inputch 划给驻留子"会找不到它。 @@ -350,9 +324,9 @@ func (p *Plugin) handleBuiltin(conn net.Conn, line string, s *sdk.PluginSDK) boo case "/persona": p.cmdPersona(conn, parts, s) case "/terminals": - p.cmdTerminals(conn) + p.cmdTerminals(conn, s) case "/cmd/history": - p.cmdCmdHistory(conn) + p.cmdCmdHistory(conn, s) case "/terminal": p.cmdTerminal(conn, parts, s) default: @@ -982,124 +956,41 @@ func (p *Plugin) cmdAgents(conn net.Conn, s *sdk.PluginSDK) { // ======== /terminals /cmd/history /terminal ======== -// subscribeToolEvents 订阅内核工具与终端事件,维护命令历史与终端会话视图。 +// 终端与命令历史的权威视图在内核(「内核开,两个插件接」)。 // -// 这两份数据不是 WebUI 插件私有的:WebUI 也是订阅同样的 EventToolCall / -// EventTerminalOutput 自己攒出来的(见 handler_chat.go 的 handleToolEvent、 -// handler_terminal.go 的 subscribeTerminalStream)。事件面本就是 SDK 对内部 -// 插件开放的,所以 CLI 能做到同口径,不需要新增内核接口。 -func (p *Plugin) subscribeToolEvents(s *sdk.PluginSDK) { - s.Subscribe(sdk.EventToolCall, func(ev *sdk.Event) { - payload := ev.Payload - tool, _ := payload["tool"].(string) - args, _ := payload["args"].(map[string]interface{}) - status, _ := payload["status"].(string) - switch tool { - case "cmd_run": - cmd := "" - if args != nil { - cmd, _ = args["command"].(string) - } - p.cmdMu.Lock() - p.cmdHistory = append(p.cmdHistory, cliCmdExec{ - Command: cmd, Status: status, Time: time.Now().Format(time.RFC3339), - }) - if len(p.cmdHistory) > 100 { - p.cmdHistory = p.cmdHistory[len(p.cmdHistory)-100:] - } - p.cmdMu.Unlock() - case "terminal_create": - id := "" - if args != nil { - id, _ = args["id"].(string) - } - if id == "" { - // agent 调用时不知道生成的 id,从工具结果中回填(同 WebUI) - if res, ok := payload["result"].(map[string]interface{}); ok { - id, _ = res["id"].(string) - } - } - if id == "" { - return - } - cmd := "" - if args != nil { - cmd, _ = args["command"].(string) - } - p.termMu.Lock() - if old, ok := p.termStates[id]; ok { - old.Command = cmd - old.Running = true - } else { - p.termStates[id] = &cliTermState{ - ID: id, Command: cmd, Running: true, - CreatedAt: time.Now().Format(time.RFC3339), - } - } - p.termMu.Unlock() - case "terminal_close": - id := "" - if args != nil { - id, _ = args["id"].(string) - } - if id != "" { - p.termMu.Lock() - if t, ok := p.termStates[id]; ok { - t.Running = false - } - p.termMu.Unlock() - } - } - }) - s.Subscribe(sdk.EventTerminalOutput, func(ev *sdk.Event) { - payload := ev.Payload - id, _ := payload["terminal_id"].(string) - if id == "" { - return - } - output, _ := payload["output"].(string) - running, _ := payload["running"].(bool) - p.termMu.Lock() - ts, ok := p.termStates[id] - if !ok { - ts = &cliTermState{ID: id, CreatedAt: time.Now().Format(time.RFC3339)} - p.termStates[id] = ts - } - ts.Running = running - if output != "" { - const maxTermOutput = 64 * 1024 - if len(ts.Output)+len(output) > maxTermOutput { - excess := len(ts.Output) + len(output) - maxTermOutput - if len(ts.Output) > excess { - ts.Output = ts.Output[excess:] - } else { - ts.Output = "" - } - } - ts.Output += output - } - p.termMu.Unlock() - }) -} +// 内核订自己的事件总线归并 EventToolCall(terminal_create/close、cmd_run) +// 与 EventTerminalOutput(agentcli 生命周期 + 输出),产出唯一真相挂在 +// TerminalAPI(见 internal/agent/core/terminal_registry.go)。 +// WebUI 与 CLI 都从 s.Status().GetKernelStatus() 读同一份,不再各自订阅推导—— +// 之前两份推导还各自踩过同一个坑:工具 result 是 Go map 文本,断言成 +// map[string]interface{} 永远失败 → /terminals 空空如也。 -// cmdTerminals 与 WebUI 的 GET /api/v1/terminals 同口径。 -func (p *Plugin) cmdTerminals(conn net.Conn) { - p.termMu.Lock() - list := make([]*cliTermState, 0, len(p.termStates)) - for _, t := range p.termStates { - list = append(list, t) +// cmdTerminals 读内核权威快照,与 WebUI 的 GET /api/v1/terminals 同源。 +func (p *Plugin) cmdTerminals(conn net.Conn, s *sdk.PluginSDK) { + t := s.Terminal() + if t == nil { + writeJSONContent(conn, map[string]interface{}{"terminals": []interface{}{}}) + return } - p.termMu.Unlock() - writeJSONContent(conn, map[string]interface{}{"terminals": list}) + terms := t.ListTerminals() + if terms == nil { + terms = []sdk.TerminalStatus{} + } + writeJSONContent(conn, map[string]interface{}{"terminals": terms}) } -// cmdCmdHistory 与 WebUI 的 GET /api/v1/cmd/history 同口径。 -func (p *Plugin) cmdCmdHistory(conn net.Conn) { - p.cmdMu.Lock() - out := make([]cliCmdExec, len(p.cmdHistory)) - copy(out, p.cmdHistory) - p.cmdMu.Unlock() - writeJSONContent(conn, map[string]interface{}{"history": out}) +// cmdCmdHistory 读内核权威快照,与 WebUI 的 GET /api/v1/cmd/history 同源。 +func (p *Plugin) cmdCmdHistory(conn net.Conn, s *sdk.PluginSDK) { + t := s.Terminal() + if t == nil { + writeJSONContent(conn, map[string]interface{}{"history": []interface{}{}}) + return + } + cmds := t.CmdHistory() + if cmds == nil { + cmds = []sdk.CmdExecStatus{} + } + writeJSONContent(conn, map[string]interface{}{"history": cmds}) } // cmdTerminal 通过 ToolAPI.ExecuteTool 调 agentcli 的终端工具。 diff --git a/internal/plugins/webui/handler.go b/internal/plugins/webui/handler.go index c17c44a..932cf5a 100644 --- a/internal/plugins/webui/handler.go +++ b/internal/plugins/webui/handler.go @@ -96,6 +96,7 @@ type Handler struct { settings sdk.SettingsAPI pluginMgr sdk.PluginManager status sdk.StatusAPI + term sdk.TerminalAPI llm sdk.LLMAPI sessionMu sync.Mutex @@ -125,26 +126,23 @@ type Handler struct { chatMsgMu sync.Mutex chatMsgCache map[string]*chatMsgEntry // client_msg_id -> 首次处理结果 chatMsgOrder []string // FIFO 淘汰序 - cmdMu sync.Mutex - cmdHistory []CmdExec - termMu sync.Mutex - termStates map[string]*termState } func NewHandler(s *sdk.PluginSDK) *Handler { var ( - sup sdk.SupervisorAPI - mem sdk.MemoryAPI - idx sdk.IndexerAPI - ad sdk.AdapterAPI - cfg sdk.ConfigAPI - tm sdk.TextMemoryAPI - ks sdk.KnowledgeAPI - tr sdk.TrackerAPI - se sdk.SettingsAPI - pm sdk.PluginManager - st sdk.StatusAPI - llm sdk.LLMAPI + sup sdk.SupervisorAPI + mem sdk.MemoryAPI + idx sdk.IndexerAPI + ad sdk.AdapterAPI + cfg sdk.ConfigAPI + tm sdk.TextMemoryAPI + ks sdk.KnowledgeAPI + tr sdk.TrackerAPI + se sdk.SettingsAPI + pm sdk.PluginManager + st sdk.StatusAPI + term sdk.TerminalAPI + llm sdk.LLMAPI ) if s != nil { sup, mem, idx = s.Supervisor(), s.Memory(), s.Indexer() @@ -152,6 +150,7 @@ func NewHandler(s *sdk.PluginSDK) *Handler { tm, ks, tr = s.TextMemory(), s.Knowledge(), s.Tracker() se, pm = s.Settings(), s.PluginMgr() st, llm = s.Status(), s.LLM() + term = s.Terminal() } h := &Handler{ sdk: s, @@ -167,9 +166,9 @@ func NewHandler(s *sdk.PluginSDK) *Handler { settings: se, pluginMgr: pm, status: st, + term: term, llm: llm, sessions: make(map[string]time.Time), - termStates: make(map[string]*termState), pendingIdx: -1, chatMsgCache: make(map[string]*chatMsgEntry), sseEvents: newSSEEventRing(200), @@ -184,21 +183,11 @@ func NewHandler(s *sdk.PluginSDK) *Handler { go h.chatPersistLoop() h.loadChatHistory() if s != nil { - go h.trackToolEvents() h.subscribeChatEvents() - h.subscribeTerminalStream() } return h } -func getStr(m map[string]interface{}, key string) string { - if m == nil { - return "" - } - v, _ := m[key].(string) - return v -} - func (h *Handler) getWebUIConfig() (apiKey, username, password string, ttl time.Duration) { ttl = 24 * time.Hour if h.settings == nil { diff --git a/internal/plugins/webui/handler_chat.go b/internal/plugins/webui/handler_chat.go index 2981e5c..49c5db3 100644 --- a/internal/plugins/webui/handler_chat.go +++ b/internal/plugins/webui/handler_chat.go @@ -122,15 +122,6 @@ func (h *Handler) loadChatHistory() { h.chatMu.Unlock() } -func (h *Handler) trackToolEvents() { - if h.sdk == nil { - return - } - h.sdk.Subscribe(sdk.EventToolCall, func(ev *sdk.Event) { - h.handleToolEvent(ev) - }) -} - // subscribeChatEvents 捕获所有通道(cli/qq/webui 等)的对话轮次, // 与 handleChat 的注入一起构成完整的全通道对话历史。 func (h *Handler) subscribeChatEvents() { @@ -407,80 +398,6 @@ func (h *Handler) pendingAssistantLocked() *ChatMsg { return msg } -func (h *Handler) handleToolEvent(ev *sdk.Event) { - payload := ev.Payload - tool, _ := payload["tool"].(string) - args, _ := payload["args"].(map[string]interface{}) - status, _ := payload["status"].(string) - ts := time.Now() - - switch tool { - case "cmd_run": - exec := CmdExec{ - Command: getStr(args, "command"), - Status: status, - Time: ts.Format(time.RFC3339), - } - h.cmdMu.Lock() - h.cmdHistory = append(h.cmdHistory, exec) - if len(h.cmdHistory) > maxCmdHistory { - h.cmdHistory = h.cmdHistory[len(h.cmdHistory)-maxCmdHistory:] - } - h.cmdMu.Unlock() - - case "terminal_create": - id := getStr(args, "id") - if id == "" { - // agent 调用时不知道生成的 id,从工具结果中回填 - if res, ok := payload["result"].(map[string]interface{}); ok { - id = getStr(res, "id") - } - } - if id == "" { - break - } - cmd := getStr(args, "command") - if cmd == "" { - if res, ok := payload["result"].(map[string]interface{}); ok { - cmd = getStr(res, "command") - } - } - now := time.Now() - term := &termState{ - ID: id, - Command: cmd, - Running: true, - CreatedAt: now.Format(time.RFC3339), - created: now, - } - h.termMu.Lock() - if old, ok := h.termStates[id]; ok { - old.Command = cmd - old.Running = true - old.created = now - } else { - h.termStates[id] = term - } - if len(h.termStates) > maxTerminals { - for k := range h.termStates { - delete(h.termStates, k) - break - } - } - h.termMu.Unlock() - - case "terminal_close": - id := getStr(args, "id") - if id != "" { - h.termMu.Lock() - if t, ok := h.termStates[id]; ok { - t.Running = false - } - h.termMu.Unlock() - } - } -} - // bumpSeqLocked 分配下一个聊天序号(调用方须持 chatMu)。 // 序号单调递增、随记录落盘,作为 /chat/history?after= 的增量游标。 func (h *Handler) bumpSeqLocked() int64 { diff --git a/internal/plugins/webui/handler_terminal.go b/internal/plugins/webui/handler_terminal.go index db82b0c..4c6ad1a 100644 --- a/internal/plugins/webui/handler_terminal.go +++ b/internal/plugins/webui/handler_terminal.go @@ -1,88 +1,44 @@ package webui import ( - "time" + "net/http" sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" - "net/http" ) -// 终端面:持久终端会话状态、终端接口、命令历史。 +// 终端面:终端会话列表与命令历史接口。 +// +// 「内核开,两个插件接」后,终端会话与命令历史的权威视图在内核 +// (internal/agent/core/terminal_registry.go):内核订阅自己的事件总线 +// 归并 EventToolCall(terminal_create/close、cmd_run)与 EventTerminalOutput +// (agentcli 生命周期 + 输出)。WebUI 不再自己订阅事件攒一份——直接读内核 +// 开放的 TerminalAPI(s.Terminal()),与 CLI 同源同口径。 +// +// 为什么不塞进 KernelStatus:那是全量快照,前端每 3 秒轮询 /kernel, +// 把每终端最多 64KB 的输出缓冲背进去会让轮询成本爆炸。 -type CmdExec struct { - Command string `json:"command"` - Stdout string `json:"stdout"` - Stderr string `json:"stderr"` - ExitCode int `json:"exit_code"` - Status string `json:"status"` - Time string `json:"time"` -} - -type termState struct { - ID string `json:"id"` - Command string `json:"command"` - Running bool `json:"running"` - Output string `json:"output"` - CreatedAt string `json:"created_at"` - Uptime string `json:"uptime"` - created time.Time -} - -const maxCmdHistory = 100 - -const maxTerminals = 50 - -// subscribeTerminalStream 常驻订阅终端实时画面推流(terminal_output 事件), -// 维护 termStates 的 Running 状态与全量输出缓冲,供 /api/v1/terminals 与前端轮询使用。 -func (h *Handler) subscribeTerminalStream() { - if h.sdk == nil { +// handleTerminals 返回内核权威的终端会话快照。 +func (h *Handler) handleTerminals(w http.ResponseWriter, r *http.Request) { + if h.term == nil { + writeJSON(w, http.StatusOK, map[string]interface{}{"terminals": []interface{}{}}) return } - h.sdk.Subscribe(sdk.EventTerminalOutput, func(ev *sdk.Event) { - id, _ := ev.Payload["terminal_id"].(string) - if id == "" { - return - } - output, _ := ev.Payload["output"].(string) - running, _ := ev.Payload["running"].(bool) - h.termMu.Lock() - ts, ok := h.termStates[id] - if !ok { - ts = &termState{ID: id, created: time.Now()} - h.termStates[id] = ts - } - ts.Running = running - if output != "" { - const maxTermOutput = 64 * 1024 - if len(ts.Output)+len(output) > maxTermOutput { - excess := len(ts.Output) + len(output) - maxTermOutput - if len(ts.Output) > excess { - ts.Output = ts.Output[excess:] - } else { - ts.Output = "" - } - } - ts.Output += output - } - h.termMu.Unlock() - }) -} - -func (h *Handler) handleTerminals(w http.ResponseWriter, r *http.Request) { - h.termMu.Lock() - terms := make([]*termState, 0, len(h.termStates)) - for _, ts := range h.termStates { - ts.Uptime = time.Since(ts.created).Round(time.Second).String() - terms = append(terms, ts) + terms := h.term.ListTerminals() + if terms == nil { + terms = []sdk.TerminalStatus{} } - h.termMu.Unlock() writeJSON(w, http.StatusOK, map[string]interface{}{"terminals": terms}) } +// handleCmdHistory 返回内核权威的命令执行历史快照。 func (h *Handler) handleCmdHistory(w http.ResponseWriter, r *http.Request) { - h.cmdMu.Lock() - result := make([]CmdExec, len(h.cmdHistory)) - copy(result, h.cmdHistory) - h.cmdMu.Unlock() - writeJSON(w, http.StatusOK, map[string]interface{}{"history": result}) + if h.term == nil { + writeJSON(w, http.StatusOK, map[string]interface{}{"history": []interface{}{}}) + return + } + cmds := h.term.CmdHistory() + if cmds == nil { + cmds = []sdk.CmdExecStatus{} + } + writeJSON(w, http.StatusOK, map[string]interface{}{"history": cmds}) } diff --git a/internal/sdk/plugin.go b/internal/sdk/plugin.go index 4cb5b3a..4181c15 100644 --- a/internal/sdk/plugin.go +++ b/internal/sdk/plugin.go @@ -167,6 +167,7 @@ type PluginSDK struct { config ConfigAPI tool ToolAPI indexer IndexerAPI + terminal TerminalAPI selftestMu sync.Mutex selftest *VirtualInstance @@ -314,6 +315,7 @@ type SDKConfig struct { Config ConfigAPI Tool ToolAPI Indexer IndexerAPI + Terminal TerminalAPI } func New(name string, cfg SDKConfig) *PluginSDK { @@ -353,6 +355,7 @@ func New(name string, cfg SDKConfig) *PluginSDK { config: cfg.Config, tool: cfg.Tool, indexer: cfg.Indexer, + terminal: cfg.Terminal, } } @@ -399,6 +402,7 @@ func (s *PluginSDK) Tracker() TrackerAPI { return s.tracker } func (s *PluginSDK) Config() ConfigAPI { return s.config } func (s *PluginSDK) Tool() ToolAPI { return s.tool } func (s *PluginSDK) Indexer() IndexerAPI { return s.indexer } +func (s *PluginSDK) Terminal() TerminalAPI { return s.terminal } func (s *PluginSDK) InjectInput(source, channel, eventType string, payload map[string]interface{}) { if s.iom != nil { diff --git a/internal/sdk/terminal.go b/internal/sdk/terminal.go new file mode 100644 index 0000000..9686b34 --- /dev/null +++ b/internal/sdk/terminal.go @@ -0,0 +1,39 @@ +package sdk + +// TerminalAPI 是内核开放的终端会话与命令历史接口(「内核开,两个插件接」)。 +// +// 为什么单开接口而不是塞进 KernelStatus:KernelStatus 是**全量运行态快照**, +// /api/v1/kernel 与前端总览页每 3 秒轮询一次;把每个终端最多 64KB 的输出 +// 缓冲塞进去,会让每次轮询都背一份终端全屏内容。终端输出属于**按需拉取**的 +// 明细,只该在 /terminals 被访问时取。 +// +// 权威状态由内核维护(internal/agent/core/terminal_registry.go):内核订阅 +// 自己的事件总线,归并 EventToolCall(terminal_create/close、cmd_run)与 +// EventTerminalOutput(agentcli 生命周期 + 输出),产出唯一真相。 +// WebUI 与 CLI 都从这里读,不再各自订阅推导。 +type TerminalAPI interface { + // ListTerminals 返回终端会话快照(含输出缓冲与运行状态)。 + ListTerminals() []TerminalStatus + // CmdHistory 返回命令执行历史快照(cmd_run,最近 100 条)。 + CmdHistory() []CmdExecStatus +} + +// TerminalStatus 与 WebUI termState / CLI cliTermState 同一 JSON 口径。 +type TerminalStatus struct { + ID string `json:"id"` + Command string `json:"command"` + Running bool `json:"running"` + Output string `json:"output"` + CreatedAt string `json:"created_at"` + Uptime string `json:"uptime"` +} + +// CmdExecStatus 与 WebUI CmdExec 同一 JSON 口径。 +type CmdExecStatus struct { + Command string `json:"command"` + Stdout string `json:"stdout"` + Stderr string `json:"stderr"` + ExitCode int `json:"exit_code"` + Status string `json:"status"` + Time string `json:"time"` +}