refactor(terminal): 内核开终端/命令历史权威视图,WebUI 与 CLI 都改接内核

按「内核开,两个插件接」重构终端与命令历史的数据归属。

背景:此前 WebUI 与 CLI 各订 EventToolCall/EventTerminalOutput 攅一份状态,
同一件事两份推导,还各自踩过同一个坑——工具 result 是 Go 的 map 文本
(map[cols:80 ... id:term_2 ...]),断言成 map[string]interface{} 永远失败,
terminal_create 的 id 回填不生效,/terminals 因此恒空(WebUI 也一样)。
实测确认:WebUI 自己的 /api/v1/terminals 与 /api/v1/cmd/history 同样是空的。

内核开(权威唯一真相):
- internal/agent/core/terminal_registry.go:TerminalRegistry 归并两类事件——
  EventToolCall(terminal_create/close、cmd_run,id/command 从 args 或 Go map
  文本回填)与 EventTerminalOutput(agentcli 生命周期 + 输出,含 64KB 缓冲上限、
  100 条命令历史、50 个终端上限)。
- internal/sdk/terminal.go:新增 TerminalAPI(ListTerminals/CmdHistory)与
  TerminalStatus/CmdExecStatus DTO。**不塞进 KernelStatus**:那是全量快照,
  前端每 3 秒轮询 /kernel,背上每终端最多 64KB 输出会让轮询成本爆炸;
  终端输出是按需拉取的明细,另开接口。
- Agent 订阅自己的事件总线(subscribeTerminalRegistry),且**只根 agent 建**
  (驻留子共用同一总线,每个子都建会 N+1 份重复记账)。
- SDKConfig/Registry/bootstrap 接线:pluginReg.SetTerminalAPI(agent)。

生产者补全(agentcli):终端无输出时 ticker 不发事件,内核就无从知道终端
存在。新增 emitTermState,在 handleCreate/handleClose/readLoop 退出(超时/
进程结束/读取错误/stopCh)显式上报 running 状态,并给输出事件补 command 字段。
handleClose 改为接收 *sdk.PluginSDK 以便上报。

两个插件接(消费方):
- WebUI:删掉本地 termStates/cmdHistory/subscribeTerminalStream/handleToolEvent
  及不再使用的 getStr;/terminals 与 /cmd/history 直接读 s.Terminal()。
- CLI:删掉上一轮刚加的 subscribeToolEvents 与 cliTermState/cliCmdExec;
  /terminals 与 /cmd/history 直接读 s.Terminal()。两条路(local/remote)都通。

测试:新增 terminal_registry_test.go,锁死 Go map 文本解析(旧缺陷根因)、
生命周期、CLI 直调路径(无 EventToolCall 仅凭 output 事件建条目)、历史与
终端数量上限。全量 go test ./internal/... ./cmd/... 通过。
This commit is contained in:
JianFeeeee
2026-09-17 18:56:15 +08:00
parent f671632915
commit 6b4f470b62
14 changed files with 595 additions and 330 deletions

View File

@ -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()

View File

@ -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 := ""

View File

@ -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
}

View File

@ -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)
}
}