diff --git a/cmd/waiter/main.go b/cmd/waiter/main.go index 602951a..1afe38c 100644 --- a/cmd/waiter/main.go +++ b/cmd/waiter/main.go @@ -2,13 +2,13 @@ package main import ( "context" - "encoding/json" "flag" "fmt" "os" "os/signal" "path/filepath" "strings" + "sync" "syscall" "time" ) @@ -26,6 +26,48 @@ const clearLine = "\033[2K\r" var colors = true +var spinnerFrames = []string{"⠋", "⠙", "⠹", "⠸", "⠼", "⠴", "⠦", "⠧", "⠇", "⠏"} + +// isTTYFile 判断文件是否为字符终端(非终端时禁用 spinner 转圈)。 +func isTTYFile(f *os.File) bool { + fi, err := f.Stat() + if err != nil { + return false + } + return fi.Mode()&os.ModeCharDevice != 0 +} + +// startSpinner 启动 npm 风格的加载动画,返回停止函数。 +// stop() 幂等:终止动画并清除当前行。非终端环境直接空操作。 +func startSpinner(label string) func() { + if !colors || !isTTYFile(os.Stdout) { + return func() {} + } + done := make(chan struct{}) + var once sync.Once + go func() { + ticker := time.NewTicker(80 * time.Millisecond) + defer ticker.Stop() + i := 0 + for { + select { + case <-done: + return + case <-ticker.C: + fmt.Printf("%s%s %s%s\n", clearLine, colorDim, spinnerFrames[i%len(spinnerFrames)]+" "+label, colorReset) + i++ + } + } + }() + stop := func() { + once.Do(func() { + close(done) + fmt.Print(clearLine) + }) + } + return stop +} + func init() { if os.Getenv("NO_COLOR") != "" { colors = false @@ -144,12 +186,18 @@ func main() { } func oneshot(state *State, msg string) { - resp, err := state.SendChat(msg) + stop := startSpinner("thinking...") + resp, err := state.SendChatStream(msg, func(rl respLine) { + // 第一个过程帧到达即停转,后续帧直接渲染 + stop() + printServerEvent(rl) + }) + stop() if err != nil { printlnC(colorRed, fmt.Sprintf("error: %v", err)) os.Exit(1) } - fmt.Println(resp) + printlnC(colorGreen, resp) } // runCapTest 本地能力测试(无需连接服务器) @@ -203,10 +251,18 @@ func runInteractive(state *State, cfg *Config) { fmt.Println("Type /help for commands.") var readerCancel func() + var spinnerStopMu sync.Mutex + var spinnerStop = func() {} startReader := func() { ctx, cancel := context.WithCancel(context.Background()) readerCancel = cancel - go state.ReadLoop(ctx, printServerOutput) + go state.ReadLoop(ctx, func(line string) { + spinnerStopMu.Lock() + stop := spinnerStop + spinnerStopMu.Unlock() + stop() + printServerOutput(line) + }) } startReader() @@ -260,6 +316,11 @@ loop: state.Send(cmd) } + // 发送成功后启动加载动画,收到第一帧服务器输出时自动停止 + spinnerStopMu.Lock() + spinnerStop = startSpinner("thinking...") + spinnerStopMu.Unlock() + select { case <-sigCh: break loop @@ -272,22 +333,71 @@ loop: } } +// printServerOutput 渲染一行服务器输出(JSON 帧)。 func printServerOutput(content string) { + rl := parseRespLineStruct(content) if !colors { - fmt.Printf("%s%s\n", clearLine, content) + fmt.Printf("%s%s\n", clearLine, renderPlain(rl, content)) return } - var rl respLine - if err := json.Unmarshal([]byte(content), &rl); err != nil { - fmt.Printf("%s%s%s\n", clearLine, content, colorReset) + printServerEventColored(rl, content) +} + +// printServerEvent 渲染一个已解析的过程/终结事件。 +func printServerEvent(rl respLine) { + if !colors { + fmt.Printf("%s%s\n", clearLine, renderPlain(rl, "")) return } + printServerEventColored(rl, "") +} + +// renderPlain 无色模式下的纯文本渲染。 +func renderPlain(rl respLine, raw string) string { switch rl.Type { + case "reasoning": + return "[思考] " + rl.Content + case "tool_call": + return fmt.Sprintf("[工具] %s (%s) %s", rl.Tool, rl.Status, rl.Result) + case "response": + return rl.Content + case "error": + return "[错误] " + rl.Error + default: + if raw != "" { + return raw + } + return rl.Content + } +} + +// printServerEventColored 彩色模式下的帧渲染。 +func printServerEventColored(rl respLine, raw string) { + switch rl.Type { + case "reasoning": + fmt.Printf("%s%s· %s%s\n", clearLine, colorDim, rl.Content, colorReset) + case "tool_call": + mark, markColor := "⚙", colorYellow + switch rl.Status { + case "ok": + mark, markColor = "✔", colorGreen + case "denied", "interrupted", "error": + mark, markColor = "✘", colorRed + } + preview := rl.Result + if preview != "" { + preview = " " + preview + } + fmt.Printf("%s%s%s %s [%s]%s%s\n", clearLine, markColor, mark, rl.Tool, rl.Status, preview, colorReset) case "response": fmt.Printf("%s%s%s%s\n", clearLine, colorGreen, rl.Content, colorReset) case "error": fmt.Printf("%s%s%s%s\n", clearLine, colorRed, rl.Error, colorReset) default: - fmt.Printf("%s%s%s\n", clearLine, content, colorReset) + text := raw + if text == "" { + text = rl.Content + } + fmt.Printf("%s%s%s\n", clearLine, text, colorReset) } } diff --git a/cmd/waiter/state.go b/cmd/waiter/state.go index e75be56..7bf67cd 100644 --- a/cmd/waiter/state.go +++ b/cmd/waiter/state.go @@ -61,14 +61,41 @@ func (s *State) Send(line string) error { } func (s *State) SendChat(msg string) (string, error) { + return s.SendChatStream(msg, nil) +} + +// SendChatStream 发送一条对话消息并循环读取响应行直至终结帧。 +// onEvent 回调在每收到一个过程帧(reasoning/tool_call)时被调用, +// 可为 nil;返回值为最终响应内容或错误。 +func (s *State) SendChatStream(msg string, onEvent func(respLine)) (string, error) { if err := s.Send(msg); err != nil { return "", err } - line, err := s.readLine() - if err != nil { - return "", err + for { + line, err := s.readLine() + if err != nil { + return "", err + } + rl := parseRespLineStruct(line) + switch rl.Type { + case "response": + return rl.Content, nil + case "error": + if rl.Error == "" { + rl.Error = line + } + return "", fmt.Errorf("%s", rl.Error) + default: + // 过程帧:reasoning / tool_call / 旧版服务器的普通文本 + if rl.Type == "" && onEvent == nil && rl.Content == "" && rl.Error == "" { + // 非JSON旧行且无回调:直接当最终输出(向后兼容旧服务器) + return line, nil + } + if onEvent != nil { + onEvent(rl) + } + } } - return parseRespLine(line) } func (s *State) SendBuiltin(cmd string) (string, error) { @@ -96,13 +123,22 @@ type respLine struct { Type string `json:"type"` Content string `json:"content"` Error string `json:"error"` + Tool string `json:"tool"` + Status string `json:"status"` + Result string `json:"result"` +} + +// parseRespLineStruct 解析一行 JSON 响应帧,解析失败时将原文放入 Content。 +func parseRespLineStruct(line string) respLine { + var rl respLine + if err := json.Unmarshal([]byte(line), &rl); err != nil { + return respLine{Content: line} + } + return rl } func parseRespLine(line string) (string, error) { - var rl respLine - if err := json.Unmarshal([]byte(line), &rl); err != nil { - return line, nil - } + rl := parseRespLineStruct(line) switch rl.Type { case "response": return rl.Content, nil diff --git a/internal/plugins/cli/plugin.go b/internal/plugins/cli/plugin.go index f305e25..0f705bc 100644 --- a/internal/plugins/cli/plugin.go +++ b/internal/plugins/cli/plugin.go @@ -19,6 +19,11 @@ import ( // DefaultSocket 由 main.go 在 Load() 前设置,覆盖默认 socket 路径。 var DefaultSocket string +const ( + cliSource = "cli" + cliChannel = "cli" +) + func init() { plugin.RegisterPluginMeta("cli", "CLI", "CLI") plugin.RegisterFactory("cli", func(name string, config map[string]interface{}) (sdk.Plugin, error) { @@ -155,19 +160,72 @@ func (p *Plugin) handleConn(conn net.Conn, s *sdk.PluginSDK) { } } - resp := s.InjectTextSync("cli", "cli", line) - if resp != nil { - content, _ := resp.Payload["content"].(string) - writeLine(conn, map[string]interface{}{ - "type": "response", - "content": content, - }) - } else { - writeLine(conn, map[string]interface{}{ - "type": "error", - "error": "agent is not available", - }) + p.handleChat(&connWriter{conn: conn}, line, s) + } +} + +// connWriter 为单条连接提供互斥保护的 JSON 行写入。 +// 对话过程中事件订阅回调运行在事件总线的发布 goroutine 上, +// 与主循环写最终响应并发,因此写入必须串行化。 +type connWriter struct { + conn net.Conn + mu sync.Mutex +} + +func (w *connWriter) writeLine(v interface{}) { + data, err := json.Marshal(v) + if err != nil { + return + } + data = append(data, '\n') + w.mu.Lock() + w.conn.Write(data) + w.mu.Unlock() +} + +// handleChat 处理一条对话消息:订阅内核的推理/工具调用事件并实时 +// 转发给客户端(流式过程输出),InjectTextSync 返回后写出最终响应。 +// 仅插件层改动:通过 SDK 订阅事件,不触碰内核。 +func (p *Plugin) handleChat(w *connWriter, line string, s *sdk.PluginSDK) { + unsubReasoning := s.Subscribe(sdk.EventReasoning, func(evt *sdk.Event) { + if ch, _ := evt.Payload["channel"].(string); ch != cliChannel { + return } + content, _ := evt.Payload["content"].(string) + if content == "" { + return + } + w.writeLine(map[string]interface{}{"type": "reasoning", "content": content}) + }) + unsubToolCall := s.Subscribe(sdk.EventToolCall, func(evt *sdk.Event) { + if ch, _ := evt.Payload["channel"].(string); ch != cliChannel { + return + } + tool, _ := evt.Payload["tool"].(string) + status, _ := evt.Payload["status"].(string) + result, _ := evt.Payload["result"].(string) + w.writeLine(map[string]interface{}{ + "type": "tool_call", + "tool": tool, + "status": status, + "result": truncateOneLine(result, 160), + }) + }) + defer unsubReasoning() + defer unsubToolCall() + + resp := s.InjectTextSync(cliSource, cliChannel, line) + if resp != nil { + content, _ := resp.Payload["content"].(string) + w.writeLine(map[string]interface{}{ + "type": "response", + "content": content, + }) + } else { + w.writeLine(map[string]interface{}{ + "type": "error", + "error": "agent is not available", + }) } } @@ -560,6 +618,16 @@ func (p *Plugin) cmdAgents(conn net.Conn, s *sdk.PluginSDK) { // ======== helpers ======== +// truncateOneLine 将多行文本压成单行并按 rune 截断,用于事件结果预览。 +func truncateOneLine(s string, max int) string { + s = strings.Join(strings.Fields(s), " ") + r := []rune(s) + if len(r) > max { + return string(r[:max]) + "…" + } + return s +} + func writeLine(conn net.Conn, v interface{}) { data, err := json.Marshal(v) if err != nil { diff --git a/internal/plugins/webui/dashboard.html b/internal/plugins/webui/dashboard.html index e6a748c..7c671e4 100644 --- a/internal/plugins/webui/dashboard.html +++ b/internal/plugins/webui/dashboard.html @@ -3246,20 +3246,20 @@ background: } function renderReasoningCard(text, isStreaming) { - var body = renderMd(text); var preview = typeof marked !== "undefined" ? text.replace(/[\s\n]+/g, " ").slice(0, 60) : escHtml(text).replace(/<[^>]+>/g, " ").slice(0, 60); return ( - '
' + + '
' + '
' + '' + '' + (isStreaming ? __("思考中...", "Thinking...") : __("思考", "Thinking")) + '' + '
' + '
' + - '
' + body + '
' + - (isStreaming ? '
' : "") + + (isStreaming + ? '
' + escHtml(preview) + '
' + : '
' + renderMd(text) + '
') + '
' ); }