diff --git a/.gitignore b/.gitignore index a8c7c18..e00fa9a 100644 --- a/.gitignore +++ b/.gitignore @@ -36,7 +36,7 @@ third_party/homeagent-sdk/example/ codegraph.json # 运行时产物(不提交) -adapters/ +/adapters/ knowledge/ memory/ scripts/ diff --git a/internal/lua/adapters/openai.lua b/internal/lua/adapters/openai.lua index 4ef5b6c..675b66a 100644 --- a/internal/lua/adapters/openai.lua +++ b/internal/lua/adapters/openai.lua @@ -99,7 +99,28 @@ function adapter.transform_stream_chunk(raw_chunk) unified.reasoning_content = delta.reasoning_content end if delta.tool_calls then - unified.tool_calls = delta.tool_calls + local tcs = {} + for _, tc in ipairs(delta.tool_calls) do + -- OpenAI 流式格式: {function:{name,arguments}, id, type, index} + -- homed StreamChunk.ToolCalls 期望扁平格式: {id, type, name, raw_arguments} + local fn = tc["function"] + local name = (type(fn) == "table" and fn.name) or tc.name or "" + local raw_args = "" + if type(fn) == "table" and type(fn.arguments) == "string" then + raw_args = fn.arguments + elseif type(tc.arguments) == "string" then + raw_args = tc.arguments + end + -- 不能按 name 过滤:OpenAI 流式分片中后续块 name 为空但携带 arguments + -- accumulateStream 按 index 累积并在 flushToolCall 时校验 name + table.insert(tcs, { + id = tc.id or "", + type = tc.type or "function", + name = name, + raw_arguments = raw_args + }) + end + unified.tool_calls = tcs end return json.encode(unified) end diff --git a/internal/plugins/mcp/sse.go b/internal/plugins/mcp/sse.go index 41a2a06..bdccfb4 100644 --- a/internal/plugins/mcp/sse.go +++ b/internal/plugins/mcp/sse.go @@ -5,8 +5,10 @@ import ( "encoding/json" "fmt" "io" + "net" "net/http" "strings" + "time" ) // SSETransport 通过 HTTP POST 进行 JSON-RPC 通信。 @@ -18,9 +20,19 @@ type SSETransport struct { } func NewSSETransport(url string) *SSETransport { + // 必须带超时:远程 MCP server 网络抖动/无响应时, + // 无超时的 client 会让插件加载永久阻塞(webui 等后续插件全部起不来) return &SSETransport{ - url: url, - client: &http.Client{}, + url: url, + client: &http.Client{ + Timeout: 30 * time.Second, + Transport: &http.Transport{ + DialContext: (&net.Dialer{ + Timeout: 10 * time.Second, + KeepAlive: 30 * time.Second, + }).DialContext, + }, + }, } } diff --git a/internal/plugins/mcp/stdio.go b/internal/plugins/mcp/stdio.go index d3f0f6f..6341bc4 100644 --- a/internal/plugins/mcp/stdio.go +++ b/internal/plugins/mcp/stdio.go @@ -8,6 +8,7 @@ import ( "os" "os/exec" "sync" + "time" ) // StdioTransport 通过子进程 stdin/stdout 进行 JSON-RPC 通信 @@ -105,11 +106,22 @@ func (t *StdioTransport) Send(req *rpcRequest) (*rpcResponse, error) { 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") } }