From ec13eb391ae35d966f0419be754fdc1d52b93433 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Thu, 17 Sep 2026 20:00:47 +0800 Subject: [PATCH] =?UTF-8?q?fix(agentcli):=20=E7=BB=88=E7=AB=AF=E9=80=80?= =?UTF-8?q?=E5=87=BA=E5=89=8D=E8=A1=A5=E6=8E=A8=E6=AE=8B=E7=95=99=E8=BE=93?= =?UTF-8?q?=E5=87=BA=EF=BC=8C=E7=9F=AD=E5=91=BD=E4=BB=A4=E8=BE=93=E5=87=BA?= =?UTF-8?q?=E4=B8=8D=E5=86=8D=E4=B8=A2=E5=A4=B1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 线上验证「内核开、两个插件接」时发现的真实缺陷:`echo`、`ls` 这类在首个 200ms ticker 之前就结束的短命令,readLoop 走到 `!terminalRunning(t)` 分支 直接 return,残留在 stream 里的输出从未 flush。 症状:输出只留在 session.buf 里——agent 用 terminal_read 能看到,但 terminal_output 事件永远发不出去,于是内核权威视图(以及 WebUI/CLI 的 /terminals)的 output 恒为空。实测 term_2(echo HELLO_KERNEL_REGISTRY) 在 /terminals 里 output="" 而 agent 同期 terminal_read 拿到了正文。 修法: - 抽出 flushTermStream(s, t),ticker 与所有退出路径共用同一条推送路径 (避免以后再出现「某条退出路径忘了 flush」)。 - readLoop 顶部加 defer:defer flushTermStream 后于 defer emitTermState 声明 → LIFO 下先 flush 再报停止,保证「最后一段输出」先于 running=false 到达订阅者。 测试:plugin_flush_test.go 新增 TestReadLoopFlushesOutputOnExit——走 EventBus 捕获事件,推入输出后立即让进程退出(远早于 ticker),断言输出 已补推且末态 running=false。已验证去掉修复即 FAIL、加回即 PASS。 注:handleRead(clear=true) 会主动 Reset stream(避免与读取结果重复), 属既有设计;本修复针对的是「未被读取就退出」的路径。 --- internal/plugins/agentcli/plugin.go | 54 +++++++---- .../plugins/agentcli/plugin_flush_test.go | 95 +++++++++++++++++++ 2 files changed, 133 insertions(+), 16 deletions(-) create mode 100644 internal/plugins/agentcli/plugin_flush_test.go diff --git a/internal/plugins/agentcli/plugin.go b/internal/plugins/agentcli/plugin.go index 2b2de11..e80dcd0 100644 --- a/internal/plugins/agentcli/plugin.go +++ b/internal/plugins/agentcli/plugin.go @@ -802,6 +802,36 @@ func (p *Plugin) handleList() (interface{}, error) { }, nil } +// flushTermStream 把终端待推送的增量输出作为 terminal_output 事件发布。 +// +// 必须抽成公共函数:退出路径(进程结束/超时/读取错误/stopCh)与常规 200ms +// ticker 都要走同一条推送路径。否则短命令(echo/ls 这类在首个 ticker 之前 +// 就结束的)残留在 stream 里的输出永远发不出事件,只留在 buf 里—— +// agent 用 terminal_read 能看到,但内核权威视图的 output 恒为空。 +// 返回是否真的发布了(无残留时为 false)。 +func flushTermStream(s *sdk.PluginSDK, t *TerminalSession) bool { + if s == nil || t == nil { + return false + } + var streamData string + t.mu.Lock() + if t.stream.Len() > 0 { + streamData = t.stream.String() + t.stream.Reset() + } + t.mu.Unlock() + if streamData == "" { + return false + } + s.Publish(&sdk.Event{ + Type: sdk.EventTerminalOutput, + Source: "agentcli", + Payload: map[string]interface{}{"terminal_id": t.id, "command": t.command, "output": streamData, "running": terminalRunning(t)}, + Timestamp: time.Now().UnixMilli(), + }) + return true +} + // emitTermState 把终端存活状态作为 EventTerminalOutput 事件上报。 // // 生命周期事件(创建/关闭/退出/超时)必须显式发:终端无输出时 ticker @@ -838,8 +868,14 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) { flushTicker := time.NewTicker(200 * time.Millisecond) defer flushTicker.Stop() - // 退出时确保内核权威视图标记该终端为已停止(readLoop 的所有 return 点)。 + // 退出时先把残留输出推出去,再报停止。 + // + // defer 是 LIFO:下面这两行声明顺序决定执行顺序——先 flush 后 emit。 + // 若反了,停止事件会先于最后一段输出到达,内核会先把 running 置 false + // 再追加输出(状态看着对但顺序错);更重要的是短命令的输出 + // 只存在于 stream 里,不 flush 就彻底丢了。 defer emitTermState(s, t, false) + defer flushTermStream(s, t) // 立即发送首次"终端已启动"通知,让 agent 感知存在。 // 用 NoMemory:这是状态提示,不是对话内容。不关掉的话每开一个终端都会 @@ -906,21 +942,7 @@ func (p *Plugin) readLoop(t *TerminalSession, s *sdk.PluginSDK) { return case <-flushTicker.C: // 批量推送终端实时画面增量(独立 ticker,避免被高密度数据饿死) - var streamData string - t.mu.Lock() - if t.stream.Len() > 0 { - streamData = t.stream.String() - t.stream.Reset() - } - t.mu.Unlock() - if streamData != "" { - s.Publish(&sdk.Event{ - Type: sdk.EventTerminalOutput, - Source: "agentcli", - Payload: map[string]interface{}{"terminal_id": t.id, "command": t.command, "output": streamData, "running": terminalRunning(t)}, - Timestamp: time.Now().UnixMilli(), - }) - } + flushTermStream(s, t) case r := <-readCh: if r.err != nil { // 读取错误/EOF → 立即通知(进程可能已结束) diff --git a/internal/plugins/agentcli/plugin_flush_test.go b/internal/plugins/agentcli/plugin_flush_test.go new file mode 100644 index 0000000..42b1d4f --- /dev/null +++ b/internal/plugins/agentcli/plugin_flush_test.go @@ -0,0 +1,95 @@ +//go:build linux || windows + +package agentcli + +import ( + "sync" + "testing" + "time" + + "gitcode.com/JianFeeeee/HomeAgent/internal/events" + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// eventCapture 订阅事件总线,收下 terminal_output 事件供断言。 +type eventCapture struct { + mu sync.Mutex + outs []map[string]interface{} +} + +func (c *eventCapture) add(ev *events.Event) { + c.mu.Lock() + c.outs = append(c.outs, ev.Payload) + c.mu.Unlock() +} + +// outputFor 汇总某个终端已发布的全部 output 片段。 +func (c *eventCapture) outputFor(id string) string { + c.mu.Lock() + defer c.mu.Unlock() + var s string + for _, p := range c.outs { + if p["terminal_id"] == id { + if o, _ := p["output"].(string); o != "" { + s += o + } + } + } + return s +} + +func (c *eventCapture) statesFor(id string) []bool { + c.mu.Lock() + defer c.mu.Unlock() + var out []bool + for _, p := range c.outs { + if p["terminal_id"] == id { + r, _ := p["running"].(bool) + out = append(out, r) + } + } + return out +} + +// 短命令(在首个 200ms ticker 之前就结束)的输出必须在退出时补推, +// 否则只留在 buf 里、永远发不出 terminal_output 事件, +// 内核权威视图(以及 WebUI/CLI 的 /terminals)output 恒为空。 +func TestReadLoopFlushesOutputOnExit(t *testing.T) { + p := New("agentcli") + bus := events.NewBus() + capture := &eventCapture{} + bus.Subscribe(events.EventTerminalOutput, capture.add) + + sdkInst := sdk.New("agentcli", sdk.SDKConfig{ + RegTool: func(string, sdk.ToolDef, sdk.ToolHandler) error { return nil }, + RegStage: func(sdk.Stage, sdk.StageHandler) {}, + RegAPI: func(string) error { return nil }, + Settings: sdk.NewSettings("agentcli", nil), + EventBus: bus, + }) + sdkInst.SetIOInjector(&injectCapture{}) + + term := newMockTerm() + ts := newTestSession(term) + startReadLoop(p, sdkInst, ts) + + // 推入输出后立刻让进程退出——远早于 200ms ticker。 + term.push([]byte("HELLO_KERNEL_REGISTRY\n")) + time.Sleep(50 * time.Millisecond) + term.setRunning(false) + + select { + case <-ts.done: + case <-time.After(3 * time.Second): + t.Fatal("readLoop did not exit") + } + + if got := capture.outputFor(ts.id); got != "HELLO_KERNEL_REGISTRY\n" { + t.Fatalf("output on exit = %q, want %q", got, "HELLO_KERNEL_REGISTRY\n") + } + // 停止状态也要报到(readLoop 退出时 emitTermState(false))。 + states := capture.statesFor(ts.id) + if len(states) == 0 || states[len(states)-1] { + t.Fatalf("expected trailing running=false, got %v", states) + } +}