mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-23 10:28:06 +00:00
fix(agentcli): 终端退出前补推残留输出,短命令输出不再丢失
线上验证「内核开、两个插件接」时发现的真实缺陷:`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(避免与读取结果重复), 属既有设计;本修复针对的是「未被读取就退出」的路径。
This commit is contained in:
@ -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 → 立即通知(进程可能已结束)
|
||||
|
||||
95
internal/plugins/agentcli/plugin_flush_test.go
Normal file
95
internal/plugins/agentcli/plugin_flush_test.go
Normal file
@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user