Files
HomeAgent/internal/plugins/mcp/stdio.go
JianFeeeee 51eb98e0ae fix(mcp+adapter): 流式 tool_calls 解析修复 + MCP transport 超时保护
openai adapter transform_stream_chunk:
- OpenAI 流式分片是嵌套格式 function.{name,arguments},原样透传后
  json.Unmarshal 到扁平 ToolCall{name,arguments} 时 name 恒为空,
  accumulateStream flushToolCall 因 acc.name=="" 静默丢弃整个工具调用
  (流式路径自上线起 tool_calls 全部丢失的根因)
- 现在正确解包 function.name → name,function.arguments → raw_arguments
  (保留原始 JSON 字符串分片,由 accumulateStream 按 index 拼接)
- 不按 name 过滤分片:OpenAI 流式续传块 name 为空但携带 arguments 分片

mcp stdio/sse transport:
- stdio Send() 无超时:server 进程卡死时插件加载永久阻塞
- sse http.Client 无超时:远程 server 网络抖动/无响应时永久阻塞,
  导致 webui 等后续插件全部无法启动(生产实例偶发启动卡死根因)
- stdio 加 60s 请求超时;sse client 加 30s 整体 + 10s 拨号超时
2026-08-25 18:59:00 +08:00

136 lines
2.8 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package mcp
import (
"bufio"
"encoding/json"
"fmt"
"io"
"os"
"os/exec"
"sync"
"time"
)
// StdioTransport 通过子进程 stdin/stdout 进行 JSON-RPC 通信
type StdioTransport struct {
cmd *exec.Cmd
stdin io.WriteCloser
stdout *bufio.Reader
mu sync.Mutex
pending map[int]chan *rpcResponse
done chan struct{}
}
func NewStdioTransport(command string, args []string, env []string) (*StdioTransport, error) {
cmd := exec.Command(command, args...)
if len(env) > 0 {
cmd.Env = append(os.Environ(), env...)
}
stdin, err := cmd.StdinPipe()
if err != nil {
return nil, fmt.Errorf("stdin pipe: %w", err)
}
stdout, err := cmd.StdoutPipe()
if err != nil {
return nil, fmt.Errorf("stdout pipe: %w", err)
}
// 忽略 stderrMCP 服务器可能输出日志到 stderr
cmd.Stderr = nil
if err := cmd.Start(); err != nil {
return nil, fmt.Errorf("start %s: %w", command, err)
}
t := &StdioTransport{
cmd: cmd,
stdin: stdin,
stdout: bufio.NewReader(stdout),
pending: make(map[int]chan *rpcResponse),
done: make(chan struct{}),
}
go t.readLoop()
return t, nil
}
func (t *StdioTransport) readLoop() {
dec := json.NewDecoder(t.stdout)
for {
var resp rpcResponse
if err := dec.Decode(&resp); err != nil {
close(t.done)
// 通知所有等待的请求
t.mu.Lock()
for _, ch := range t.pending {
close(ch)
}
t.pending = make(map[int]chan *rpcResponse)
t.mu.Unlock()
return
}
t.mu.Lock()
ch, ok := t.pending[resp.ID]
delete(t.pending, resp.ID)
t.mu.Unlock()
if ok {
ch <- &resp
}
}
}
func (t *StdioTransport) Send(req *rpcRequest) (*rpcResponse, error) {
data, err := json.Marshal(req)
if err != nil {
return nil, fmt.Errorf("marshal request: %w", err)
}
ch := make(chan *rpcResponse, 1)
t.mu.Lock()
t.pending[req.ID] = ch
t.mu.Unlock()
if _, err := t.stdin.Write(data); err != nil {
t.mu.Lock()
delete(t.pending, req.ID)
t.mu.Unlock()
return nil, fmt.Errorf("write stdin: %w", err)
}
if _, err := t.stdin.Write([]byte("\n")); err != nil {
t.mu.Lock()
delete(t.pending, req.ID)
t.mu.Unlock()
return nil, fmt.Errorf("write newline: %w", err)
}
// 超时保护MCP server 进程启动慢/卡死不响应时,不能让插件加载永久阻塞
// stdio server 冷启动如 npx 拉包可能数十秒,给 60s
timeout := time.After(60 * time.Second)
select {
case resp := <-ch:
return resp, nil
case <-t.done:
t.mu.Lock()
delete(t.pending, req.ID)
t.mu.Unlock()
return nil, fmt.Errorf("mcp transport closed")
case <-timeout:
t.mu.Lock()
delete(t.pending, req.ID)
t.mu.Unlock()
return nil, fmt.Errorf("mcp request timeout (60s), method may be hung")
}
}
func (t *StdioTransport) Close() error {
t.stdin.Close()
if t.cmd.Process != nil {
t.cmd.Process.Kill()
}
<-t.done
return t.cmd.Wait()
}